如何在Azure Data Factory的Data Flow中使用自定义批处理活动与Python脚本
问题解答
关于Data Flow中使用自定义批处理活动的说明
Data Flow是ADF中专注于数据流ETL的组件,无法直接在Data Flow内部嵌入自定义批处理活动。自定义活动属于管道级别的执行单元,必须通过管道将自定义活动与Data Flow活动串联协作。
场景1:自定义活动处理数据 → Data Flow排序写入Sink
实现步骤:
- 自定义活动执行Python脚本:用现有脚本读取存储中的Parquet文件,完成转换后,将结果写入到临时存储路径(如ADLS Gen2/Blob的临时文件夹),输出格式选择Data Flow支持的类型(推荐Parquet,保持数据类型一致性)。示例Python代码片段:
import pandas as pd # 读取源Parquet df = pd.read_parquet("abfss://<容器名>@<存储账户名>.dfs.core.windows.net/源路径/") # 执行自定义转换逻辑 transformed_df = df[["col1", "col2"]].query("col1 > 0") # 写入临时存储 transformed_df.to_parquet("abfss://<容器名>@<存储账户名>.dfs.core.windows.net/临时路径/transformed_data.parquet") - 添加Data Flow活动:在管道中自定义活动之后,添加Data Flow活动,将Data Flow的源指向上述临时存储的转换后文件。
- Data Flow内配置排序与Sink:在Data Flow中添加
Sort转换组件,设置排序字段和规则;最后配置Sink组件,将排序后的数据写入目标存储或数据库。 - 可选:清理临时文件:管道末尾添加
Delete活动,删除临时路径下的文件,避免存储冗余。
场景2:Data Flow读源 → 自定义活动转换 → Data Flow写入Sink
实现步骤:
- 第一个Data Flow活动导出中间数据:用已创建的Data Flow读取Parquet源文件,将数据写入临时存储路径(作为中间层)。
- 自定义活动执行Python转换:读取临时存储的中间数据,用Pandas完成自定义转换逻辑,再将转换后的数据写入另一个临时存储路径。
- 第二个Data Flow活动写入Sink:添加第二个Data Flow活动,源指向转换后的临时文件,直接配置Sink组件完成最终数据写入。
Python与Data Flow结合的关键注意事项
- 依赖中间存储衔接:Data Flow与自定义活动无法直接共享内存数据,必须通过ADLS Gen2/Blob等存储作为数据中转层。
- 参数化路径提升灵活性:将存储路径设为管道参数,在自定义活动和Data Flow中引用参数,便于后续修改和多环境适配。
- 权限配置:确保ADF的托管标识(或服务主体)、Python运行环境(如自托管集成运行时、Azure Batch)拥有临时存储的读写权限。
- 临时文件生命周期管理:通过管道的Delete活动或存储的生命周期策略,自动清理临时文件,降低存储成本。
内容的提问来源于stack exchange,提问作者Nishad Nazar
相关产品推荐
相关产品推荐

