在Flink中如何为Python脚本传递参数或读取配置文件?
解答
传递自定义参数给Python脚本
完全可以,直接在提交命令的末尾追加参数即可,示例命令:
bin/flink run --python my_driver.py --input /data/raw --output /data/result
在my_driver.py里,可以用Python原生的sys.argv直接获取参数:
import sys if __name__ == "__main__": # sys.argv[0]是脚本自身路径,从索引1开始为传入的自定义参数 input_path = sys.argv[1] output_path = sys.argv[2] # 后续业务逻辑中使用参数
如果参数较多或需要键值对形式,推荐用argparse模块做更规范的参数解析:
import argparse if __name__ == "__main__": parser = argparse.ArgumentParser(description="Flink Python 任务参数") parser.add_argument("--input", required=True, help="输入数据路径") parser.add_argument("--output", required=True, help="输出结果路径") parser.add_argument("--parallelism", type=int, default=2, help="任务并行度") args = parser.parse_args() # 使用解析后的参数 print(f"输入路径: {args.input}")
在Python脚本中读取配置文件
当然支持,Python常见的配置文件格式(ini、json、yaml等)都可以直接读取,以下是几种常见方式:
1. 读取INI配置文件
假设存在config.ini:
[data] input_path = /data/raw output_path = /data/result [job] parallelism = 4
在脚本中用configparser读取:
import configparser if __name__ == "__main__": config = configparser.ConfigParser() config.read("config.ini") input_path = config.get("data", "input_path") parallelism = config.getint("job", "parallelism")
集群环境注意事项:如果是提交到Flink集群,需要确保配置文件能被所有TaskManager节点访问。可以用--pyFiles参数将配置文件作为依赖上传,Flink会自动分发到各个节点:
bin/flink run --python my_driver.py --pyFiles config.ini
2. 读取JSON配置文件
如果用JSON格式的config.json:
{ "data": { "input_path": "/data/raw", "output_path": "/data/result" }, "job": { "parallelism": 4 } }
用json模块读取:
import json if __name__ == "__main__": with open("config.json", "r") as f: config = json.load(f) input_path = config["data"]["input_path"]
3. 读取YAML配置文件
若使用YAML格式,需先安装pyyaml依赖,然后在脚本中读取:
import yaml if __name__ == "__main__": with open("config.yaml", "r") as f: config = yaml.safe_load(f) input_path = config["data"]["input_path"]
提交任务时需将依赖包和配置文件一起上传:
bin/flink run --python my_driver.py --pyFiles config.yaml,requirements.txt
内容的提问来源于stack exchange,提问作者Rinze
相关产品推荐
相关产品推荐

