如何将Apache Beam DataFlow Pipeline同步运行?代码改写求助
Apache Beam DataFlow 管道执行顺序调整方案
需求说明
原管道为异步并行执行,需调整为以下执行顺序:
- 先执行
last_sync_data任务,读取并提取最后同步时间 - 待
last_sync_data完成后,并行执行wisgen_data和recruiter_data的读取与写入任务 - 待上述所有任务完成后,执行末尾的时间戳写入任务,更新同步时间
实现原理
在Apache Beam中,任务的执行顺序由数据依赖关系控制:如果一个PTransform依赖于另一个PCollection的输出,Beam会确保上游任务完成后才启动下游任务。我们可以通过构建数据依赖来强制任务执行顺序:
- 将
last_sync_data的输出作为wisgen_data和recruiter_data的启动触发信号 - 将
wisgen_data和recruiter_data的写入任务输出合并,作为时间戳写入任务的依赖
改写后的完整代码
import datetime import apache_beam as beam from apache_beam.io.jdbc import ReadFromJdbc from apache_beam.io.gcp.bigquery import ReadFromBigQuery, WriteToBigQuery, BigQueryDisposition # 先初始化Pipeline对象(原代码顺序错误,需调整) p = beam.Pipeline(options=options) # Step 1: 执行last_sync_data任务,读取并提取最后同步时间 last_sync_data = ( p | "Read last sync data" >> ReadFromBigQuery( query='SELECT MAX(time) as last_sync FROM `wisdomcircle-350611.custom_test_data.last_sync`' ) | "Extract last sync time" >> beam.Map(lambda elem: elem['last_sync']) ) # Step 2: 依赖last_sync_data完成后,并行执行wisgen_data任务 wisgen_data = ( last_sync_data | "Trigger wisgen job" >> beam.Map(lambda _: None) # 用last_sync_data触发任务启动 | "wisgen job" >> ReadFromJdbc( jdbc_url=jdbc_url, username=username, password=password, driver_class_name='org.postgresql.Driver', query="""SELECT users.id AS user_id, CONCAT(users.first_name,' ', users.last_name) AS full_name FROM users""", ) | "Convert TableRow to dict(wisgen data)" >> beam.Map(lambda row: row._asdict()) | "Write to BigQuery in wisgen data table" >> WriteToBigQuery( wisgen_table, write_disposition=BigQueryDisposition.WRITE_APPEND, create_disposition=BigQueryDisposition.CREATE_IF_NEEDED, schema='user_id:INTEGER,full_name:STRING' ) ) # Step 2: 依赖last_sync_data完成后,并行执行recruiter_data任务 recruiter_data = ( last_sync_data | "Trigger recruiter job" >> beam.Map(lambda _: None) # 用last_sync_data触发任务启动 | "recruiter job" >> ReadFromJdbc( jdbc_url=jdbc_url, username=username, password=password, driver_class_name='org.postgresql.Driver', query="""SELECT users.id AS user_id, '' AS full_name FROM users""", ) | "Convert TableRow to dict(recruiter data)" >> beam.Map(lambda row: row._asdict()) | "Write to BigQuery in recruiter data table" >> WriteToBigQuery( recruiter_table, write_disposition=BigQueryDisposition.WRITE_APPEND, create_disposition=BigQueryDisposition.CREATE_IF_NEEDED, schema='user_id:INTEGER,full_name:STRING' ) ) # Step 3: 等待wisgen和recruiter任务完成后,执行时间戳写入 _ = ( [wisgen_data, recruiter_data] | "Combine completion signals" >> beam.Flatten() # 合并两个任务的输出信号 | "Create current timestamp" >> beam.Map(lambda _: {'time': datetime.datetime.utcnow()}) | "Write to last_sync" >> WriteToBigQuery( last_sync_tab, schema='time:TIMESTAMP', write_disposition=BigQueryDisposition.WRITE_TRUNCATE, create_disposition=BigQueryDisposition.CREATE_IF_NEEDED, ) ) # 启动管道并等待所有任务完成 result = p.run() result.wait_until_finish()
关键修改点说明
- 修正Pipeline初始化顺序:原代码中
last_sync_data定义在p = beam.Pipeline()之前,属于语法错误,必须先创建Pipeline对象才能定义后续任务。 - 构建任务依赖链:
- 通过
last_sync_data | "Trigger xxx job" >> beam.Map(lambda _: None),让wisgen和recruiter任务依赖last_sync_data的完成,确保它们在last_sync_data之后启动。 - 通过
beam.Flatten合并wisgen和recruiter任务的输出,让时间戳写入任务依赖这两个任务的完成,确保最后执行。
- 通过
- 修复原代码细节问题:
- 将BigQuery表名修正为标准反引号格式:
project.dataset.table - 为Recruiter的SQL查询补充
AS full_name,确保输出字段与schema匹配 - 添加管道启动和等待完成的代码,确保任务按顺序执行完毕
- 将BigQuery表名修正为标准反引号格式:
内容的提问来源于stack exchange,提问作者Shubham Belwal
相关产品推荐
相关产品推荐

