如何在Dataflow任务成功后触发指定Cloud Function?
Dataflow任务成功后触发Cloud Function的实现方案
以下是修正后的完整代码,既修复了原代码中的语法问题,又实现了仅当Dataflow任务执行成功时,向指定Pub/Sub主题发送消息触发Cloud Function的需求:
from apache_beam import Pipeline, PipelineOptions from apache_beam.io import ReadFromBigQuery, WriteToBigQuery from apache_beam.io.jdbc import ReadFromJdbc from google.cloud import pubsub_v1 import datetime def run_pipeline(jdbc_url, username, password, wisgen_table): # 先查询BigQuery获取最后同步时间(避免侧输入的复杂处理) last_sync_query = """ SELECT IFNULL(MAX(time), '1970-01-01 00:00:00') as last_sync FROM `wisdomcircle-350611.custom_test_data.last_sync` """ # 直接执行查询获取单值 last_sync_result = ReadFromBigQuery(query=last_sync_query).run().get() last_sync_time = last_sync_result[0]['last_sync'] options = PipelineOptions() p = Pipeline(options=options) # 读取Wisgen数据(修正SQL语法错误) wisgen_data = ( p | "Wisgen数据读取" >> ReadFromJdbc( jdbc_url=jdbc_url, username=username, password=password, driver_class_name='org.postgresql.Driver', query=f""" SELECT users.id AS user_id, users.full_name FROM users WHERE to_char(users.updated_at, 'YYYY-MM-DD HH24:MI:SS') >= '{last_sync_time}' """, table_name="users" ) ) # 写入BigQuery(修复schema引号闭合问题) ( wisgen_data | "转换为字典格式" >> lambda row: row._asdict() | "写入BigQuery" >> WriteToBigQuery( wisgen_table, write_disposition=WriteToBigQuery.WRITE_APPEND, create_disposition=WriteToBigQuery.CREATE_IF_NEEDED, schema='user_id:INTEGER,full_name:STRING' ) ) # 运行流水线并获取状态 result = p.run() result.wait_until_finish() # 如果任务执行成功,发送消息到Pub/Sub主题触发Cloud Function if result.state == 'DONE': publisher = pubsub_v1.PublisherClient() topic_path = "projects/wisdomcircle-350611/topics/uat-timestamp-job-trigger" # 构造消息内容(可根据需求调整) message = f"Dataflow任务执行成功,完成时间:{datetime.datetime.now().isoformat()}" publisher.publish(topic_path, message.encode("utf-8")) print("已发送触发消息到Cloud Function") if __name__ == '__main__': run_pipeline( jdbc_url='xxxx', username='xxxx', password='xxxx', wisgen_table='wisdomcircle-350611:uat_custom_dataset.wisgen_data_table', )
关键修改说明
- 修正SQL语法:移除了原JDBC查询中位置错误的
where,补充了查询字段users.full_name,避免语法报错 - 简化同步时间获取:直接在流水线外查询BigQuery获取最后同步时间,更简便地构造动态SQL
- 修复BigQuery写入配置:补全了schema字符串的闭合引号,避免配置错误
- 任务成功触发逻辑:通过
result.wait_until_finish()等待任务完成,判断状态为DONE时,使用Pub/Sub客户端发送消息到指定主题,触发订阅该主题的Cloud Function - 依赖说明:需要提前安装Pub/Sub客户端库,执行命令:
pip install google-cloud-pubsub
内容的提问来源于stack exchange,提问作者Shubham Belwal
相关产品推荐
相关产品推荐

