求助:Azure Data Factory数据流动态列名设置及ADLS文件去重
Azure Data Factory 数据流实现动态列名+多文件去重写入ADLS
一、动态列名在窗口/聚合活动中的配置
先在数据流中定义数组类型参数,比如dedupKeyColumns(存储去重的键列)、aggregateColumns(存储需要聚合的列),然后通过表达式动态生成列配置:
- 聚合活动:分组依据选择
map(dedupKeyColumns, #item),聚合列使用map(aggregateColumns, @(name = #item, expression = concat('first(', #item, ')'))),自动生成所有目标列的聚合规则,无需手动逐个添加。 - 窗口活动:分区列设置为
map(dedupKeyColumns, #item),排序列按需配置(比如按时间列降序以保留最新记录),新增列用rowNumber()生成行号,列名可设为rn。
二、读取ADLS多文件
在数据流的源数据集(ADLS Gen2)中:
- 勾选允许读取多个文件,使用通配符(如
*.parquet)或参数化路径(比如用管道参数传递文件前缀/文件夹路径)读取批量文件。 - 开启允许架构漂移并勾选推断漂移列的类型,确保动态列能被正确识别。
三、去重逻辑(无需Azure Databricks)
两种实现方式按需选择:
- 窗口+筛选去重:在窗口活动生成行号后,添加筛选活动,保留
rn == 1的记录,这种方式会保留原始数据结构,仅去除重复键的冗余记录。 - 聚合去重:通过聚合活动按去重键分组,对其他列取
first()/max()等聚合值,适合需要合并重复记录的场景。
四、完整数据流流程
- 源活动:配置ADLS多文件读取,开启架构漂移。
- 参数配置:定义数组参数存储动态列名,支持从管道传递变量或配置值。
- 去重处理:用窗口/聚合活动结合动态列参数实现去重。
- 筛选(可选):若用窗口活动,筛选行号为1的记录。
- Sink活动:配置ADLS目标数据集,开启架构漂移,选择覆盖/追加模式写入去重后的数据。
关键注意事项
- 源和Sink必须开启架构漂移,否则动态列无法被处理和写入。
- 若列名包含特殊字符,用
quote(#item)包裹列名,比如concat('max(', quote(#item), ')'),避免表达式语法错误。 - 测试时先传入固定数组参数验证逻辑,再切换为管道动态传递列名。
内容的提问来源于stack exchange,提问作者Azure_legend
相关产品推荐
相关产品推荐

