testing client_name_dag

This commit is contained in:
Umang Mistry
2024-03-28 12:15:15 -05:00
parent 559b87f5c8
commit 41ad10fbd5
+68
View File
@@ -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