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

Data Fusion管道失败时触发邮件通知及BigQuery记录方案咨询

Google Cloud Data Fusion 管道失败处理与校验状态记录实现思路

一、捕获管道失败事件并触发动作(发邮件/创建失败标记表)

Data Fusion 可通过全局错误处理或组件级错误分支捕获失败事件,结合对应组件实现需求:

1. 全局错误处理配置(捕获全管道失败)

  • 打开管道编辑器顶部的「设置」面板,启用全局错误处理。
  • 在错误处理流程中添加对应组件:
    • 发送邮件:添加Email组件,配置SMTP服务信息(如Cloud SMTP或企业邮箱),在邮件内容中引用内置变量填充失败详情:
      • 主题:Data Fusion 管道 ${pipelineId} 执行失败
      • 内容:失败时间:${currentTimestamp()} \n 错误信息:${errorMessage}
    • 创建BigQuery失败标记表:添加BigQuery组件,配置目标数据集与表名(勾选「自动创建表」),写入模式选「追加」,构造失败记录字段:
      • pipeline_id:值填${pipelineId}
      • failure_time:值填${currentTimestamp()}
      • status:固定值'FAILED'
      • error_details:值填${errorMessage}

2. 组件级错误分支(捕获特定步骤失败)

如果需要精准捕获某个组件(如数据源读取、校验逻辑)的失败,可给该组件单独设置「On Failure」分支:

  • 选中目标组件,在右侧面板找到「错误处理」,选择「添加分支」。
  • 在分支中添加邮件或BigQuery组件,配置逻辑同全局错误处理,还可额外添加该组件的专属标识(如failed_component: 'DataValidationStep')到BigQuery记录中。

二、结合现有数据校验,同步记录管道失败状态

针对你现有「校验一致存1、不一致存0」的逻辑,可通过分支合并实现管道失败时的状态记录:

  1. 统一状态表设计
    在BigQuery中设计一张状态记录表,字段建议包含:
  • execution_id:管道执行ID(${executionId})
  • execution_time:执行时间(${currentTimestamp()})
  • validation_status:状态值(1=校验一致,0=校验不一致,2=管道失败)
  • error_info:错误详情(可选,管道失败时填充)
  • batch_identifier:数据批次标识(如数据日期、批次ID)
  1. 分支逻辑整合
  • 现有校验流程:保持原有的校验逻辑(如用Wrangler组件对比输入输出数据,生成validation_status为1或0的记录),输出到一个「合并节点」。
  • 错误分支:在全局或组件级错误分支中,构造一条validation_status为2的记录,填充error_info与其他元数据,同样输出到「合并节点」。
  • 合并写入:将两个分支的输出合并后,用同一个BigQuery组件以「追加」模式写入状态记录表,确保所有场景(校验成功、校验失败、管道失败)的状态都被记录。
  1. 细节优化
  • 若校验步骤本身可能失败(如数据格式错误导致校验逻辑崩溃),需给校验组件单独设置「On Failure」分支,直接输出validation_status为2的记录,避免管道中断导致状态丢失。
  • 可利用Wrangler组件的if-error函数,在校验逻辑中捕获局部错误,直接生成对应状态记录,无需触发管道全局失败。

内容的提问来源于stack exchange,提问作者Henry Atehortua

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 01:37:43