无法在Google Cloud Composer中通过DataflowCreatePythonJobOperator启动任务
Airflow DAG使用DataflowCreatePythonJobOperator无任务生成问题
使用以下Airflow DAG代码调用DataflowCreatePythonJobOperator时,代码目录等配置均无异常且未出现报错,但Composer UI中始终无法生成任务:
import airflow from airflow import DAG from airflow.providers.google.cloud.operators.dataflow import DataflowCreatePythonJobOperator from datetime import datetime, timedelta default_args = { 'start_date': airflow.utils.dates.days_ago(0), 'retries': 0, 'retry_delay': timedelta(minutes=1) } dag = DAG( 'custom_python_operator', default_args=default_args, description='Just runs templates', schedule_interval='0 * * * *', max_active_runs=1, catchup=False, dagrun_timeout=timedelta(minutes=10), ) dataflow_job = DataflowCreatePythonJobOperator( task_id='run_dataflow_pipeline', py_file='./template_runner_scripts/MultiSoapAPI_repeat_bkp.py', # GCS path job_name='direct-script-1617' , project_id='ai-data-ingestion-staging', location='us-central1', options={ 'ip_configuration': 'WORKER_IP_PUBLIC' }, py_requirements=['apache-beam[gcp]==2.21.0'], py_interpreter='python3', py_system_site_packages=False ) dataflow_job
问题原因
代码中的DataflowCreatePythonJobOperator实例未关联到创建的DAG对象,Airflow无法识别该任务所属的DAG,因此不会在UI中展示任务。
解决方案
在初始化DataflowCreatePythonJobOperator时添加dag=dag参数,将任务与DAG关联:
修改后的任务代码片段:
dataflow_job = DataflowCreatePythonJobOperator( task_id='run_dataflow_pipeline', py_file='./template_runner_scripts/MultiSoapAPI_repeat_bkp.py', # GCS path job_name='direct-script-1617' , project_id='ai-data-ingestion-staging', location='us-central1', options={ 'ip_configuration': 'WORKER_IP_PUBLIC' }, py_requirements=['apache-beam[gcp]==2.21.0'], py_interpreter='python3', py_system_site_packages=False, dag=dag # 添加该行,关联到目标DAG )
也可通过任务依赖方式关联(适用于多任务场景):
dataflow_job >> None # 显式将任务加入DAG的依赖链
修改后重新部署DAG,Composer UI即可正常显示任务。
内容的提问来源于stack exchange,提问作者Vamshi
相关产品推荐
相关产品推荐

