Azure Data Factory自定义日志:数据流规则校验日志表构建求助
Azure Data Factory 数据流日志记录实现方案
核心思路
通过管道层面捕获数据流各接收器的写入行数,将规则与行数关联后,循环调用存储过程写入日志表。
步骤1:数据流接收器命名规范
给每个对应业务规则的坏记录接收器设置明确关联规则的名称,比如sink_rule_手机号格式错误、sink_rule_必填字段为空,确保后续能直接通过名称对应到具体业务规则。
步骤2:管道中捕获接收器写入行数
在执行数据流的活动之后,通过ADF动态表达式获取每个接收器的写入行数:
- 单个接收器行数表达式示例:
@activity('你的数据流活动名称').output.sinks.sink_rule_手机号格式错误.rowsWritten - 整理规则与行数为数组变量:
- 新建一个数组类型变量(比如
ruleLogList) - 使用
Set Variable活动,将规则信息和对应行数组装成数组,示例表达式:@createArray( createObject('ruleDesc', '手机号格式不符合要求', 'failedCount', activity('你的数据流活动名称').output.sinks.sink_rule_手机号格式错误.rowsWritten), createObject('ruleDesc', '核心必填字段为空', 'failedCount', activity('你的数据流活动名称').output.sinks.sink_rule_必填字段为空.rowsWritten) )
- 新建一个数组类型变量(比如
步骤3:循环调用存储过程写入日志
- 添加ForEach活动,设置
Items为刚才的数组变量@variables('ruleLogList') - 在ForEach内部添加存储过程活动,配置连接到日志表所在的数据源(比如Snowflake)
- 给存储过程传递参数:
- 规则描述参数值:
@item().ruleDesc - 失败记录数参数值:
@item().failedCount - 可额外添加执行时间参数:
@utcNow()(根据日志表结构调整)
- 规则描述参数值:
存储过程参考逻辑(Snowflake示例)
假设日志表字段为RULE_DESCRIPTION、FAILED_RECORDS、EXECUTION_TIME,存储过程可写为:
CREATE OR REPLACE PROCEDURE INSERT_VALIDATION_LOG( p_rule_desc VARCHAR, p_failed_count INT, p_exec_time TIMESTAMP ) RETURNS VARCHAR LANGUAGE SQL AS $$ BEGIN INSERT INTO VALIDATION_LOG(RULE_DESCRIPTION, FAILED_RECORDS, EXECUTION_TIME) VALUES(p_rule_desc, p_failed_count, p_exec_time); RETURN 'Success'; END; $$;
内容的提问来源于stack exchange,提问作者Smitha
相关产品推荐
相关产品推荐

