updated client config dag

This commit is contained in:
Umang Mistry
2024-03-28 13:15:58 -05:00
parent 41ad10fbd5
commit 8122345dd2
+21 -6
View File
@@ -5,10 +5,25 @@ 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):
SNOWFLAKE_CONN_ID = "doczy_dev_snowflake"
DAG_ID = "load_client_config"
DATABASE="DOCZY_DEV"
# bucket = "airflow-data-ingestion"
TAGS=["dev","config_interface","dataload"]
default_params = {"client_list_file_name": ""}
def process_csv_in_s3(params):
# Instantiate S3Hook
file_name = params['client_list_file_name']
s3_hook = S3Hook(aws_conn_id="aws_default")
bucket_name = "doczy-dev-infra-raw-data-ingestion"
key = 'client_names_openair/' + file_name
# Get the file from S3
file_content = s3_hook.read_key(key, bucket_name)
@@ -40,15 +55,15 @@ def generate_s3_path(client_name):
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/'
base_s3_path = 's3://'
return base_s3_path + final_name + "/"
default_args = {
'owner': 'airflow',
'depends_on_past': False,
'start_date': datetime(2024, 3, 26),
'start_date': datetime(2024, 3, 28),
'retries': 1,
'retry_delay': timedelta(minutes=5),
'retry_delay': timedelta(minutes=5)
}
dag = DAG(
@@ -56,13 +71,13 @@ dag = DAG(
default_args=default_args,
description='Process a CSV in S3 and replace it',
schedule_interval='@daily',
params = default_params
)
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,
dag=dag
)
process_csv_task