BeamRunPythonPipelineOperator未提交Dataflow任务问题排查
问题分析与解决方案
核心问题
你的DAG配置中,pipeline_options里指定了template_location参数,这会触发Apache Beam执行生成模板文件的操作,而非直接提交Dataflow运行任务——这就是为什么只生成模板、没有启动Dataflow作业的原因。手动运行时你应该没有添加该参数,所以能正常启动任务。
解决方案
方案1:直接提交Dataflow作业(无需预生成模板)
这种方式适合一次性或按需运行的场景,直接跳过模板生成步骤,让Composer直接提交Dataflow任务到Project B。
- 修改DAG配置:移除
pipeline_options中的template_location参数,修正temp_location为目录(日志中显示它被设为了文件路径,这是错误的):
start_python_job = BeamRunPythonPipelineOperator( task_id="start-python-jobdf1", runner="DataflowRunner", py_file="/home/airflow/gcs/data/SessPubSubDataFlow.py", py_options=[], pipeline_options={ 'temp_location':"gs://abc-tempstreamsdataflow/temp/", # 改为目录路径 'project':"abc-temp" # Project B的项目ID }, py_requirements=['apache-beam[gcp]==2.37.0'], py_interpreter='python3', py_system_site_packages=False, dataflow_config={ 'location': 'us-east4', 'project_id':'abc-temp', # Project B的项目ID 'gcp_conn_id':'0-app', "wait_until_finished": True, # 改为True,方便查看作业执行状态 'job_name':'{{task.task_id}}-{{ds_nodash}}' }, )
- 配置跨项目权限:确保Project A中Composer的Worker服务账号,在Project B拥有以下权限:
roles/dataflow.developer:提交Dataflow作业的权限- PubSub相关权限(如果读取Project B的PubSub数据)
- 目标数据库的写入权限
- GCS存储桶的读写权限(访问临时文件)
方案2:先生成模板,再从模板启动作业(适合重复运行场景)
如果需要复用模板,可以拆分两个任务:先生成模板到GCS,再用模板启动Dataflow作业。
import datetime from airflow import models from airflow import DAG from airflow.utils.dates import days_ago from airflow.providers.apache.beam.operators.beam import BeamRunPythonPipelineOperator from airflow.providers.google.cloud.operators.dataflow import DataflowTemplatedJobStartOperator default_args = { "start_date": days_ago(1), 'retries':0, 'project':'abc-temp' } with models.DAG( dag_id="DataFlowPythonJob", start_date=days_ago(1), default_args=default_args, schedule_interval="@once" ) as dag: # 第一步:生成Dataflow模板 generate_template = BeamRunPythonPipelineOperator( task_id="generate-dataflow-template", runner="DataflowRunner", py_file="/home/airflow/gcs/data/SessPubSubDataFlow.py", pipeline_options={ 'template_location':"gs://abc-tempstreamsdataflow/Redislookup.txt", 'temp_location':"gs://abc-tempstreamsdataflow/temp/", 'project':"abc-temp" }, py_requirements=['apache-beam[gcp]==2.37.0'], py_interpreter='python3', py_system_site_packages=False, dataflow_config={ 'location': 'us-east4', 'project_id':'abc-temp', 'gcp_conn_id':'0-app' } ) # 第二步:从GCS模板启动Dataflow作业 start_template_job = DataflowTemplatedJobStartOperator( task_id="start-dataflow-from-template", template="gs://abc-tempstreamsdataflow/Redislookup.txt", job_name="start-python-jobdf1-{{ds_nodash}}", project_id="abc-temp", location="us-east4", gcp_conn_id="0-app", wait_until_finished=True ) # 设置任务依赖 generate_template >> start_template_job
额外注意事项
- 修正
temp_location:必须指向GCS目录,不能是文件路径,否则会导致Dataflow作业执行失败 - 解决Python版本警告:日志提示Python版本兼容问题,确保Composer环境的Python版本与Beam 2.37.0兼容(推荐3.7-3.9版本)
- 检查
gcp_conn_id:确保该连接对应的服务账号拥有Project B的所有必要权限
内容的提问来源于stack exchange,提问作者Nagesh B Viswanadham
相关产品推荐
相关产品推荐

