You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.28 23:27:17