如何在Azure Data Factory中处理Parquet文件并输出Parquet文件
最优实现方案(ADF原生映射数据流方案,无额外服务依赖)
该方案完全满足优先使用ADF原生能力的要求,无需额外申请服务密钥、无需升级云服务层级,适配所有核心需求。
第一步:新建独立参数化Parquet数据集(解决原有数据集参数报错问题)
- 不要复用原有同步管道的通用数据集,单独新建Azure Data Lake Storage Gen2类型的Parquet加工数据集
- 数据集新增两个参数:
cw_folderPath、cw_fileName,将两个参数的默认值均设置为空字符串——此前报错No value provided for Parameter 'cw_fileName'的核心原因是原有同步数据集将该参数设为必填项,且未配置默认值,数据流调用时未传参就会触发报错 - 数据集路径配置:容器选择存储账户对应容器,目录路径填动态内容
@dataset().cw_folderPath,文件名填动态内容@dataset().cw_fileName,开启「允许通配符路径」选项
第二步:配置数据流源,实现通配符批量读取历史文件
- 新建通用映射数据流,拖入源转换,关联上述新建的加工数据集
- 源设置中「源类型」选择通配符路径,通配符路径填入对应实体的路径规则,例如
Raw/CRM/*/*/*/campaign.parquet,无需手动传入cw_fileName参数,通配符模式会直接覆盖数据集的文件名配置,参数使用默认空值即可避免报错 - 开启源侧「允许架构漂移」选项,兼容不同日期抽取的Parquet文件存在的微小字段差异,避免任务运行失败
第三步:按实体类型配置定制化加工逻辑
- 通用字段处理:拖入「选择」转换,通过规则映射配置需要保留的字段,直接在别名列填写重命名后的目标字段名,无需额外编写转换代码
- 仅保留最新版本数据的实体:
- 拖入「聚合」转换,按实体业务主键(如leadid、contactid)分组,聚合列配置
max(RowStartDate),别名为LatestRowStartDate - 拖入「存在」转换,将原数据流与聚合结果做内连接,匹配条件为
业务主键相等 AND RowStartDate == LatestRowStartDate,过滤后即可得到每个主键对应的最新版本数据
- 拖入「聚合」转换,按实体业务主键(如leadid、contactid)分组,聚合列配置
- 需要构建缓慢变化维(SCD Type2)的实体:
- 拖入「窗口」转换,按业务主键分区,按
RowStartDate升序排序,窗口函数选择lead(RowStartDate,1),生成中间字段NextRowStartDate - 拖入「派生列」转换,新增
RowEndDate字段,转换逻辑为iif(isNull(NextRowStartDate), toTimestamp('9999-12-31 23:59:59'), NextRowStartDate) - 最后通过「选择」转换移除中间生成的
NextRowStartDate字段即可
- 拖入「窗口」转换,按业务主键分区,按
第四步:配置接收器输出单个Parquet文件
- 拖入「接收器」转换,关联指向目标暂存目录的Parquet数据集,路径配置为对应实体的暂存路径,例如
Staged/CRM/campaign/ - 接收器设置中「文件名选项」选择输出到单个文件,自定义文件名填写对应实体的Parquet文件名,例如
campaign.parquet - 进入优化面板,将接收器分区设置为「单分区」,确保最终只生成1个完整Parquet文件,不会产生分片小文件
第五步:批量调度配置
- 新建ADF管道,添加「ForEach」循环活动,循环项配置为12个实体的名称数组:
['campaign','lead','contact',<其余9个实体名>] - ForEach活动内部拖入「执行数据流」活动,关联上述通用加工数据流,通过管道参数将当前循环的实体名传入数据流,动态替换通配符路径中的实体名、输出路径中的实体名、输出文件名,一次配置即可完成12个实体的加工调度,无需重复搭建数据流
原有尝试路径的阻塞问题修复方法
如果需要复用此前验证过的路径,可按以下方法解决对应阻塞:
- 数据流复用旧数据集报
No value provided for Parameter 'cw_fileName'修复:
选中数据流的源转换后,画布下方配置栏最底部有默认收起的数据集参数折叠区,展开后即可找到cw_fileName参数的配置入口,传入对应通配符文件名即可;也可以直接编辑原有数据集,将cw_fileName参数的默认值设置为空字符串,即可跳过参数校验。 - 脚本活动连接Synapse缺少Service Principal Key修复:
无需使用Service Principal认证,ADF连接Synapse可直接使用系统分配托管身份认证:在Synapse工作区访问控制页,给ADF对应的托管身份分配Synapse SQL Administrator角色,同时在存储账户的访问控制页给该身份分配Storage Blob Data Contributor角色,即可完成连接配置,无需持有SP密钥。需要输出结果到存储时,直接使用CETAS语法即可将查询结果直接写入指定存储路径。 - Databricks Standard Tier无法配置密钥作用域挂载存储修复:
无需在Databricks侧配置密钥作用域,在ADF的Databricks Notebook活动配置中,将存储账户访问密钥通过ADF加密变量传入Notebook入参,Notebook内部直接调用入参传入的密钥执行dbutils.fs.mount即可完成挂载,Standard Tier完全支持该方式。 - Synapse加工无法直接输出到存储暂存修复:
使用Synapse无服务器SQL池的CETAS(CREATE EXTERNAL TABLE AS SELECT)能力,创建外部数据源时指定目标存储容器路径和Parquet格式,执行查询后结果会直接写入指定的存储目录,不会落到Synapse内置存储中,完全满足暂存输出需求。
内容的提问来源于stack exchange,提问作者DaarioNaharis
相关产品推荐
相关产品推荐

