Airflow 2下DataflowPythonOperator报404及区域配置冲突如何解决?
问题根因
Airflow 2.x 对 Dataflow 相关 Operator 的参数解析逻辑和 1.10.x 版本存在两处关键差异,是遇到问题的核心原因:
- 1.10.x 版本中 DAG default_args 里的
dataflow_default_options会自动和任务级的options参数合并,2.x 版本中任务级options会直接覆盖默认配置,不会合并,所以自定义的job_args里没有携带 region 配置时,Dataflow 作业会默认使用 us-central1。 - 2.x 版本中 Dataflow Operator 轮询作业状态时,使用的区域是从单独的
location参数读取,不再从options里的 region 字段读取,默认值为 us-central1。就算手动在options里加了 region 配置,轮询逻辑还是会去默认的 us-central1 找作业,就会抛出404错误。
解决方案
需要同时完成两个配置修改即可解决问题:
- 合并默认的
dataflow_default_options和自定义的job_args,避免默认配置被覆盖 - 给 DataFlowPythonOperator 显式传入
location参数,指定作业所在区域
修改后的代码示例
DEFAULT_DAG_ARGS = { 'start_date': YESTERDAY, 'email': models.Variable.get('email'), 'email_on_failure': True, 'email_on_retry': False, 'retries': 0, 'project_id': models.Variable.get('gcp_project'), 'dataflow_default_options': { 'region': 'europe-west1', 'project': models.Variable.get('gcp_project'), 'temp_location': models.Variable.get('gcp_temp_location'), 'runner': 'DataflowRunner', 'zone': 'europe-west1-d' } } with models.DAG(dag_id='GcsToBigQueryTriggered', description='A DAG triggered by an external Cloud Function', schedule_interval=None, default_args=DEFAULT_DAG_ARGS, max_active_runs=1) as dag: # 合并默认配置和自定义作业参数 job_args = { **DEFAULT_DAG_ARGS['dataflow_default_options'], 'input': 'gs://{{ dag_run.conf["bucket"] }}/{{ dag_run.conf["name"] }}', 'output': models.Variable.get('bq_output_table'), 'fields': models.Variable.get('input_field_names'), 'load_dt': DS_TAG } # 显式传入location参数指定轮询区域 dataflow_task = dataflow_operator.DataFlowPythonOperator( task_id="data-ingest-gcs-process-bq", py_file=DATAFLOW_FILE, options=job_args, location="europe-west1" )
可选优化
如果有多个 Dataflow 任务,不想每个任务都单独写 location 参数,可以直接把 location 加到 DAG 的 default_args 里,所有 Dataflow 任务会自动读取该配置:
DEFAULT_DAG_ARGS = { # 原有其他配置不变 'location': 'europe-west1' }
内容的提问来源于stack exchange,提问作者ankie
相关产品推荐
相关产品推荐

