ADF管道:数据流前后行数校验及过滤行数统计需求
ADF管道验证数据流行数匹配(含过滤行数统计)方案
一、统计过滤行数的两种方式
方式1:通过数据流转换计算
- 在过滤转换前添加聚合转换,用
Count函数统计输入总行数,输出列命名为total_rows。 - 保留原过滤转换处理数据。
- 在过滤转换后添加第二个聚合转换,统计过滤后的数据行数,输出列命名为
written_rows。 - 过滤行数可通过
total_rows - written_rows得到,可将这三个值写入临时存储(如Azure SQL表、Blob文件)供管道调用。
方式2:直接读取数据流运行指标
ADF数据流运行后,过滤转换的运行日志会记录被丢弃的行数,可通过管道表达式直接获取:
@activity('你的数据流活动名称').output.runStatus.metrics.你的过滤转换名称.dropped
(替换表达式中的活动名称和转换名称为实际值)
二、管道中添加验证逻辑
执行数据流活动后,定义变量存储关键数值:
- 读取行数:
@int(activity('你的数据流活动名称').output.runStatus.metrics.源转换名称.rowsRead) - 写入行数:
@int(activity('你的数据流活动名称').output.runStatus.metrics.接收器转换名称.rowsWritten) - 过滤行数:
@int(activity('你的数据流活动名称').output.runStatus.metrics.你的过滤转换名称.dropped)(方式2)或@sub(variables('total_rows'), variables('written_rows'))(方式1)
- 读取行数:
添加If Condition活动,设置验证条件:
@equals(variables('读取行数'), add(variables('写入行数'), variables('过滤行数')))
- 在If Condition的
False分支中添加Fail活动,设置失败提示信息,例如:
"行数验证失败:读取行数{variables('读取行数')},写入行数{variables('写入行数')},过滤行数{variables('过滤行数')},不满足读取行数=写入行数+过滤行数"
当条件不满足时,管道会触发失败。
内容的提问来源于stack exchange,提问作者MicB
相关产品推荐
相关产品推荐

