From 41ad10fbd574583a02707d8b2588e01dd98cf287 Mon Sep 17 00:00:00 2001 From: Umang Mistry Date: Thu, 28 Mar 2024 12:15:15 -0500 Subject: [PATCH] testing client_name_dag --- airflow/dags/client_name_dag.py | 68 +++++++++++++++++++++++++++++++++ 1 file changed, 68 insertions(+) create mode 100644 airflow/dags/client_name_dag.py diff --git a/airflow/dags/client_name_dag.py b/airflow/dags/client_name_dag.py new file mode 100644 index 0000000..a9eccb5 --- /dev/null +++ b/airflow/dags/client_name_dag.py @@ -0,0 +1,68 @@ +import re +import pandas as pd +from airflow import DAG +from airflow.operators.python_operator import PythonOperator +from airflow.providers.amazon.aws.hooks.s3 import S3Hook +from datetime import datetime, timedelta + +def process_csv_in_s3(bucket_name, key): + # Instantiate S3Hook + s3_hook = S3Hook(aws_conn_id="aws_default") + + # Get the file from S3 + file_content = s3_hook.read_key(key, bucket_name) + + # Convert string to DataFrame + from io import StringIO + df = pd.read_csv(StringIO(file_content)) + + # Apply the generate_s3_path function + df['s3_path'] = df['client_name'].apply(generate_s3_path) + + # Convert DataFrame to CSV string + csv_buffer = StringIO() + df.to_csv(csv_buffer, index=False) + csv_content = csv_buffer.getvalue() + + # Replace the file in S3 with the processed content + s3_hook.load_string( + string_data=csv_content, + key=key, + bucket_name=bucket_name, + replace=True + ) + +def generate_s3_path(client_name): + # Your function as provided + client_name = client_name.rstrip('.') + pattern = r'[^0-9a-zA-Z!_.()*\'-]' + multi_underscore_pattern = r'_{2,}' + intermediate_name = re.sub(pattern, '_', client_name) + final_name = re.sub(multi_underscore_pattern, '_', intermediate_name) + final_name = final_name.lower() + base_s3_path = 's3://your-bucket-name/' + return base_s3_path + final_name + "/" + +default_args = { + 'owner': 'airflow', + 'depends_on_past': False, + 'start_date': datetime(2024, 3, 26), + 'retries': 1, + 'retry_delay': timedelta(minutes=5), +} + +dag = DAG( + 's3_csv_processing', + default_args=default_args, + description='Process a CSV in S3 and replace it', + schedule_interval='@daily', +) + +process_csv_task = PythonOperator( + task_id='process_csv', + python_callable=process_csv_in_s3, + op_kwargs={'bucket_name': 'your-bucket-name', 'key': 'path/to/your/file.csv'}, + dag=dag, +) + +process_csv_task