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

