You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.01 18:01:27