BeamRunPythonPipelineOperator无法跟踪Dataflow任务状态及无日志输出问题
问题:Airflow中无法跟踪Dataflow任务状态,无法按状态触发后续任务
我在Cloud Composer(版本composer-2.9.8-airflow-2.9.3)环境中,通过BeamRunPythonPipelineOperator()提交Dataflow流水线任务。任务已成功提交到Dataflow,但Airflow侧任务持续运行却没有输出任何Dataflow任务状态更新日志,最终Airflow任务以INFO - Process exited with return code: 0状态退出。需要实现在Airflow中跟踪Dataflow任务状态的功能,以便根据任务状态(如JOB_STATE_DONE)触发后续关联任务。
现有算子配置
start_dataflow_job = BeamRunPythonPipelineOperator( task_id="start_dataflow_job", runner="DataflowRunner", py_file=GCS_FILE_LOCATION, pipeline_options={ "tempLocation": GCS_BUCKET, "stagingLocation": GCS_BUCKET, "output_project": PROJECT, "service_account_email": GCP_CUSTOM_SERVICE_ACCOUNT, "requirements_file": "gs://GCS_CODE_BUCKET/requirements.txt", "max_num_workers": "2", "region": "us-east1", "experiments": [ "streaming_boot_disk_size_gb=100", "workerLogLevelOverrides=com.google.cloud.dataflow#DEBUG", "dataflow_service_options=enable_prime" ], }, py_options=[], py_requirements=["apache-beam[gcp]~=2.60.0"], py_interpreter="python3", py_system_site_packages=False, dataflow_config=DataflowConfiguration( job_name="{{task.task_id}}", project_id=PROJECT, location="us-east1", wait_until_finished=False, gcp_conn_id="google_cloud_default", ), do_xcom_push=True, )
日志信息
[2024-11-04, 04:31:42 UTC] {beam.py:151} INFO - Start waiting for Apache Beam process to complete. [2024-11-04, 04:41:09 UTC] {beam.py:172} INFO - Process exited with return code: 0 [2024-11-04, 04:41:11 UTC] {taskinstance.py:441} ▼ Post task execution logs [2024-11-04, 04:41:11 UTC] {taskinstance.py:1206} INFO - Marking task as SUCCESS. dag_id=test_data_pipeline, task_id=start_dataflow_job, run_id=manual__2024-11-04T04:24:59.026429+00:00, execution_date=20241104T042459, start_date=20241104T042502, end_date=20241104T044111 [2024-11-04, 04:41:12 UTC] {local_task_job_runner.py:243} INFO - Task exited with return code 0 [2024-11-04, 04:41:12 UTC] {taskinstance.py:3506} INFO - 1 downstream tasks scheduled from follow-on schedule check [2024-11-04, 04:41:12 UTC] {local_task_job_runner.py:222} ▲▲▲ Log group end
解决方案
1. 直接等待Dataflow任务完成(同步模式)
当前配置中dataflow_config的wait_until_finished=False是核心问题,该参数设置为False时,Airflow仅负责提交任务,不会等待Dataflow任务结束,也不会跟踪状态。将其改为True后,Airflow会持续监听Dataflow任务状态并输出日志,直到任务进入最终状态(成功/失败),后续任务会自动触发:
修改后的dataflow_config部分:
dataflow_config=DataflowConfiguration( job_name="{{task.task_id}}", project_id=PROJECT, location="us-east1", wait_until_finished=True, # 调整为True gcp_conn_id="google_cloud_default", )
2. 异步提交+状态传感器监听(非阻塞模式)
如果不想让Airflow任务长时间阻塞等待Dataflow完成,可使用DataflowJobStatusSensor单独监听任务状态:
- 利用已设置的
do_xcom_push=True,BeamRunPythonPipelineOperator会将Dataflow的job_id推送到XCom中。 - 添加状态传感器任务,依赖于提交任务:
from airflow.providers.google.cloud.sensors.dataflow import DataflowJobStatusSensor wait_for_dataflow_job = DataflowJobStatusSensor( task_id="wait_for_dataflow_job", job_id="{{ ti.xcom_pull(task_ids='start_dataflow_job')['dataflow_job_id'] }}", project_id=PROJECT, location="us-east1", expected_statuses=["JOB_STATE_DONE"], gcp_conn_id="google_cloud_default", poke_interval=60, # 每隔60秒检查一次任务状态 ) # 配置任务依赖 start_dataflow_job >> wait_for_dataflow_job >> downstream_task # downstream_task为你的后续任务
3. 优化日志输出
为确保Airflow能输出Dataflow状态更新日志,可添加logging_level="INFO"参数强制输出详细日志:
start_dataflow_job = BeamRunPythonPipelineOperator( # 保留原有参数 logging_level="INFO", )
内容的提问来源于stack exchange,提问作者mgore
相关产品推荐
相关产品推荐

