GCSFileTransformOperator的transform_script可执行脚本结构咨询
GCSFileTransformOperator转换脚本构建问题解答
问题描述
在Airflow中执行任务时,需要使用GCSFileTransformOperator预处理大型CSV文件,查阅文档后仍不清楚transform_script对应的可执行转换脚本该如何构建,询问以下脚本结构是否正确,以及Airflow是否会通过命令行传递参数调用该脚本:
# Import the required modules import preprocessing modules import sys # Define the function that passes source_file and destination_file params def preprocess_file(source_file, destination_file): # (1) code that processes the source_file # (2) code then writes to destination_file # Extract source_file and destination_file from the list of command-line arguments source_file = sys.argv[1] destination_file = sys.argv[2] preprocess_file(source_file, destination_file)
解答
脚本结构正确性
你提供的脚本结构完全符合要求,是正确的。
执行逻辑说明
Airflow确实会通过命令行调用该可执行脚本并传递参数,具体流程为:
- Operator先将GCS上的源文件下载到Airflow Worker节点的本地临时目录
- 调用你的转换脚本时,会把本地临时目录中的源文件路径作为第一个参数、本地临时目标文件路径作为第二个参数传入
- 脚本处理完成后,Operator会自动将本地的处理结果文件上传回GCS指定的目标路径
额外注意事项
- 脚本开头需添加Shebang声明(如
#!/usr/bin/env python3),并赋予可执行权限(执行chmod +x your_script.py),确保Airflow能直接调用 - 处理大型CSV时,建议使用流式读取/处理方式(比如
pandas.read_csv的chunksize参数),避免内存溢出 - 脚本依赖的所有第三方库,需要提前在Airflow Worker的运行环境中安装好
内容的提问来源于stack exchange,提问作者Christine
相关产品推荐
相关产品推荐

