Clean DAG

This commit is contained in:
Giambattista Bloisi 2024-10-28 21:10:38 +01:00
parent 6c25db9ac2
commit 7af06fbda5
1 changed files with 3 additions and 3 deletions

View File

@ -24,7 +24,7 @@ default_args = {
default_args=default_args,
params={
"S3_CONN_ID": Param("s3_conn", type='string', description="Airflow connection of S3 endpoint"),
"POSTGRES_CONN_ID": Param("postgres_conn", type='string', description="Airflow connection of S3 endpoint"),
"POSTGRES_CONN_ID": Param("postgres_conn", type='string', description=""),
"INPUT_PATH": Param("s3a://graph/tmp/prod_provision/graph/02_graph_grouped", type='string', description=""),
"OUTPUT_PATH": Param("s3a://graph/tmp/prod_provision/graph/03_graph_cleaned", type='string', description=""),
@ -43,7 +43,7 @@ default_args = {
def clean_graph_dag():
getdatasourcefromcountry = SparkKubernetesOperator(
task_id='getdatasourcefromcountry',
task_display_name="Propagate Relations",
task_display_name="Get datasource from Country",
namespace='dnet-spark-jobs',
template_spec=SparkConfigurator(
name="getdatasourcefromcountry-{{ ds }}-{{ task_instance.try_number }}",
@ -101,7 +101,7 @@ def clean_graph_dag():
kubernetes_conn_id="kubernetes_default"
))
chain([getdatasourcefromcountry, masterduplicateaction], clean_tasks)
chain(getdatasourcefromcountry, masterduplicateaction, clean_tasks)
clean_graph_dag()