- 示例工程
- 大数据
【免费下载链接】flink-learning
flink learning blog. http://www.54tianzhisheng.cn/ 含 Flink 入门、概念、原理、实战、性能调优、源码解析等内容。涉及 Flink Connector、Metrics、Library、DataStream API、Table API & SQL 等内容的学习案例,还有 Flink 落地应用的大型项目案例(PVUV、日志存储、百亿数据实时去重、监控告警)分享。欢迎大家支持我的专栏《大数据实时计算引擎 Flink 实战与性能优化》
配置管理是 Flink 作业从"能跑"走向"可运维"的关键一环。本文以 flink-learning 仓库中 flink-in-action-10.2.md 为骨架,系统讲解如何使用 Flink 的ParameterTool类统一读取运行参数、系统属性与环境变量、properties 配置文件,并结合仓库内 ExecutionEnvUtil.java 的多层配置合并方案,给出让同一个 Job 无需修改代码即可在不同环境(开发、测试、预发、生产)运行的实战范式。读完本文,你将掌握Configuration + withParameters与ParameterTool两条配置路线的差异、全部读取方式、全局参数注册、mergeWith覆盖优先级设计,以及广播变量动态更新配置的进阶方向。
一、为什么要做参数化配置:写死配置的代价
在一个真实 Flink 项目中,需要打交道的配置远不止一处:算子的并行度、Kafka 数据源地址(broker 地址、topic 名、group.id)、Checkpoint 是否开启、状态后端存储路径、数据库地址、用户名和密码……这些配置如果全部在代码中写死,就会遇到一个非常现实的问题:
你的作业能否不修改任何配置,就直接在开发、测试、预发、生产等不同环境运行?
答案通常是否定的——每个环境的配置值都不一样。如果配置是写死的,那么每换一个环境运行测试作业,都要经历"修改代码 → 编译打包 → 提交运行"的重复劳动,大量时间被消耗在毫无技术含量的重复操作上。参数化配置解决的就是这个问题:把"变的值"从代码中剥离出去,让代码本身与环境无关。
在 Flink 中有几种管理配置的方式,下面分别说明,并重点讲解最实用的ParameterTool。
二、方式一:Configuration + withParameters(批程序专用,局限明显)
Flink 提供了withParameters方法,可以把Configuration中的参数传递给函数。要使用它,需要实现Rich 函数(如RichFlatMapFunction、RichMapFunction),而不是普通函数,因为 Rich 函数才有open方法,可以在open中通过Configuration读取传入的参数值。完整示例:
ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment(); // Configuration 类来存储参数 Configuration configuration = new Configuration(); configuration.setString("name", "zhisheng"); env.fromElements(WORDS) .flatMap(new RichFlatMapFunction<String, Tuple2<String, Integer>>() { String name; @Override public void open(Configuration parameters) throws Exception { // 读取配置 name = parameters.getString("name", ""); } @Override public void flatMap(String value, Collector<Tuple2<String, Integer>> out) throws Exception { String[] splits = value.toLowerCase().split("\\W+"); for (String split : splits) { if (split.length() > 0) { out.collect(new Tuple2<>(split + name, 1)); } } } }).withParameters(configuration) // 将参数传递给函数 .print();这段代码有两个非常明显的局限性,需要特别注意:
withParameters只在批程序中支持,流程序中是没有该方法的(DataStreamAPI 不提供此能力);withParameters要在每个算子后面单独调用,并不是一次设置所有算子都能获取到。如果所有算子都需要同一份配置,就要在每个算子后重复设置,非常繁琐。
因此在实际项目中,这种方式的实用性很低,它更适合作为理解 Flink 参数传递机制的入门示例。
三、方式二:ParameterTool 统一管理配置
ParameterTool(org.apache.flink.api.java.utils.ParameterTool)是 Flink 官方提供的配置工具类,它的核心价值在于:用一套统一的 API 读取来自不同来源的配置,并且它本身是可序列化的,可以随函数传递、注册为全局参数。
ParameterTool支持三种数据来源:运行参数(arguments)、系统属性(system properties)、配置文件(properties file)。下面逐一讲解。
3.1 读取运行参数(fromArgs)
Flink UI 上支持为每个 Job 单独传入 arguments(参数),格式要求如下两种写法均可:
--brokers 127.0.0.1:9200 --username admin --password 123456或者单横线形式:
-brokers 127.0.0.1:9200 -username admin -password 123456在 Flink 程序中,通过ParameterTool.fromArgs(args)一次性获取所有参数,再通过parameterTool.get("username")按 key 取值:
ParameterTool parameterTool = ParameterTool.fromArgs(args); String username = parameterTool.get("username");fromArgs的解析规则是:--key value/-key value,同时也支持--key=value的写法。值得注意的是,fromArgs会把形如--brokers的裸参数(没有 value 跟随)解析为布尔值true,使用时需留意。
这个能力的实际意义在于:你可以把一份配置放在一个第三方接口中,通过参数传入该接口地址,Job 启动后请求该接口拉取更多配置,从而把配置彻底从代码与打包产物中解放出来。
3.2 读取系统属性(fromSystemProperties)
ParameterTool还支持通过ParameterTool.fromSystemProperties()读取 JVM 系统属性。仓库中的示例 ParameterToolGetSystemMain.java 展示了它的用法:
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.getConfig().setGlobalJobParameters(ParameterTool.fromSystemProperties()); env.addSource(new RichSourceFunction<String>() { @Override public void run(SourceContext<String> sourceContext) throws Exception { while (true) { ParameterTool parameterTool = (ParameterTool) getRuntimeContext().getExecutionConfig().getGlobalJobParameters(); sourceContext.collect(System.currentTimeMillis() + parameterTool.get("os.name") + parameterTool.get("user.home")); } } @Override public void cancel() { } }).print(); env.execute("ParameterTool Get config from SystemProperties");这里读取的os.name、user.home都是 JVM 内置系统属性。系统属性通常用于传递 JVM 级别的信息(如-Dkey=value启动参数注入的属性),作为配置来源之一。
3.3 读取配置文件(fromPropertiesFile)
除了上述两种,ParameterTool还支持ParameterTool.fromPropertiesFile("/application.properties")读取 properties 配置文件。它接受文件路径,也接受InputStream(因此可以用class.getResourceAsStream("/application.properties")读取 classpath 下的资源文件)。
有了它,你可以把所有要配置的地方(并行度、Kafka、MySQL 等配置)都写成可配置项,key 和 value 统一写在配置文件中,最后通过ParameterTool读取并分发。仓库示例 ParameterToolGetPropertiesMain.java 展示了从 classpath 资源读取配置的写法:
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.getConfig().setGlobalJobParameters( ParameterTool.fromPropertiesFile( ParameterToolGetPropertiesMain.class.getResourceAsStream("/application.properties"))); env.addSource(new RichSourceFunction<String>() { @Override public void run(SourceContext<String> sourceContext) throws Exception { while (true) { ParameterTool parameterTool = (ParameterTool) getRuntimeContext().getExecutionConfig().getGlobalJobParameters(); sourceContext.collect(System.currentTimeMillis() + parameterTool.get("metrics.topic")); } } @Override public void cancel() { } }).print(); env.execute("ParameterTool Get config from SystemProperties");注意fromPropertiesFile会抛出IOException,需要在main方法中声明throws Exception或自行捕获处理。
3.4 ParameterTool 获取值:类型安全的方法族
ParameterTool提供了一系列便捷方法获取不同类型的值,常用方法如下:
| 方法 | 说明 | 典型用途 |
|---|---|---|
get(String key) | 获取字符串值 | 通用配置、用户名 |
get(String key, String defaultValue) | 获取字符串值,缺失时返回默认值 | 可缺省配置 |
getRequired(String key) | 获取必需值,缺失抛异常 | 必填配置(如 topic 名) |
getInt(String key, int defaultValue) | 获取 int 值 | 算子并行度 |
getLong(String key, long defaultValue) | 获取 long 值 | Checkpoint 间隔 |
getBoolean(String key, boolean defaultValue) | 获取 boolean 值 | Checkpoint 开关 |
getProperties() | 返回底层Properties对象 | 构造 Kafka Consumer 配置 |
has(String key) | 判断 key 是否存在 | 分支逻辑 |
getNumberOfParameters() | 返回参数数量 | 调试 |
你可以在应用程序的main()方法中直接使用这些方法返回值,例如设置算子的并行度:
ParameterTool parameters = ParameterTool.fromArgs(args); int parallelism = parameters.get("mapParallelism", 2); DataStream<Tuple2<String, Integer>> counts = data.flatMap(new Tokenizer()).setParallelism(parallelism);3.5 将 ParameterTool 作为参数传递给自定义函数
因为ParameterTool是可序列化的,所以可以把它当作构造参数直接传给自定义函数:
ParameterTool parameters = ParameterTool.fromArgs(args); DataStream<Tuple2<String, Integer>> counts = data.flatMap(new Tokenizer(parameters));然后在函数内部使用ParameterTool获取参数值。这意味着你在作业的任何地方都可以获取到参数,而不像withParameters那样需要在每个算子后面重复设置。这是ParameterTool相比Configuration路线的核心优势之一。
3.6 注册全局参数:setGlobalJobParameters
在ExecutionConfig中,可以将ParameterTool注册为全作业级别的参数,这样它就能被 JobManager 的 Web 端以及用户自定义函数以配置值的形式访问:
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.getConfig().setGlobalJobParameters(ParameterTool.fromArgs(args));注册之后,用户自定义的 Rich 函数中可以通过getRuntimeContext()取回该参数对象:
env.addSource(new RichSourceFunction<String>() { @Override public void run(SourceContext<String> sourceContext) throws Exception { while (true) { ParameterTool parameterTool = (ParameterTool) getRuntimeContext().getExecutionConfig().getGlobalJobParameters(); sourceContext.collect(System.currentTimeMillis() + parameterTool.get("os.name") + parameterTool.get("user.home")); } } @Override public void cancel() { } })这一模式在仓库中被广泛使用:例如 ParameterToolGetArgsMain.java 演示了用--name zhisheng或-name zhisheng传参并全局注册;而 KafkaConfigUtil.java 中则通过(ParameterTool) env.getConfig().getGlobalJobParameters()取出全局参数来构建 Kafka Source,并用parameter.getRequired(PropertiesConstants.METRICS_TOPIC)读取必填 topic:
public static DataStreamSource<MetricEvent> buildSource(StreamExecutionEnvironment env) throws IllegalAccessException { ParameterTool parameter = (ParameterTool) env.getConfig().getGlobalJobParameters(); String topic = parameter.getRequired(PropertiesConstants.METRICS_TOPIC); Long time = parameter.getLong(PropertiesConstants.CONSUMER_FROM_TIME, 0L); return buildSource(env, topic, time); }从源码结构可以看到,全局参数注册 +getRuntimeContext().getExecutionConfig().getGlobalJobParameters()是跨算子共享配置的标准通道,这也是ParameterTool比withParameters更推荐的根本原因——注册一次,处处可读。
四、进阶实战:多层配置合并与 K8s 环境变量优先策略
在真实生产环境中,往往不是单一数据源,而是配置文件兜底、运行参数覆盖、环境变量优先的多层配置。仓库中的 ExecutionEnvUtil.java 给出了一个非常经典的实现:
public static ParameterTool createParameterTool(final String[] args) throws Exception { return ParameterTool .fromPropertiesFile(ExecutionEnvUtil.class.getResourceAsStream(PropertiesConstants.PROPERTIES_FILE_NAME)) .mergeWith(ParameterTool.fromArgs(args)) .mergeWith(ParameterTool.fromSystemProperties()); }其作者在书中给出过更完整的版本——将 K8s 环境变量也纳入合并链路,且优先以环境变量为准:
public static ParameterTool createParameterTool(final String[] args) throws Exception { return ParameterTool .fromPropertiesFile(ExecutionEnv.class.getResourceAsStream("/application.properties")) .mergeWith(ParameterTool.fromArgs(args)) .mergeWith(ParameterTool.fromSystemProperties()) .mergeWith(ParameterTool.fromMap(getenv()));// mergeWith 会使用最新的配置 } // 获取 Job 设置的环境变量 private static Map<String, String> getenv() { Map<String, String> map = new HashMap<>(); for (Map.Entry<String, String> entry : System.getenv().entrySet()) { map.put(entry.getKey().toLowerCase().replace('_', '.'), entry.getValue()); } return map; }这个方案的核心语义有两点,值得深入理解:
mergeWith的覆盖规则:mergeWith会使用最新的配置覆盖之前的同名 key。在上述代码中,合并顺序是"配置文件 → 运行参数 → 系统属性 → 环境变量",因此最终的优先级是:环境变量 > 系统属性 > 运行参数 > 配置文件。- 环境变量 key 规范化:K8s 环境变量名不允许出现点号(
.),通常用下划线(_)代替,因此getenv()中把环境变量的 key 统一转成小写并把_替换为.,从而与 properties 配置文件的 key 命名风格(如kafka.brokers)对齐。
结合仓库里的 PropertiesConstants.java,可以看到这套常量类把配置文件中的 key 集中管理(如kafka.brokers、kafka.group.id、stream.parallelism、stream.checkpoint.enable、stream.checkpoint.interval、metrics.topic、mysql.host等),再配合 ExecutionEnvUtil.prepare() 读取这些 key 完成环境初始化:
public static StreamExecutionEnvironment prepare(ParameterTool parameterTool) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(parameterTool.getInt(PropertiesConstants.STREAM_PARALLELISM, 5)); env.getConfig().setRestartStrategy(RestartStrategies.fixedDelayRestart(4, 60000)); if (parameterTool.getBoolean(PropertiesConstants.STREAM_CHECKPOINT_ENABLE, true)) { env.enableCheckpointing(parameterTool.getLong(PropertiesConstants.STREAM_CHECKPOINT_INTERVAL, 10000)); } env.getConfig().setGlobalJobParameters(parameterTool); env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime); return env; }可以看到,并行度、重启策略、Checkpoint 开关与间隔、全局参数注册全部由ParameterTool驱动,且每个读取都带有默认值,保证配置缺失时作业仍能按兜底值启动。类似的,KafkaConfigUtil.buildKafkaProps() 直接用parameterTool.getProperties()拿到底层Properties,再put覆盖 Kafka 必填项(bootstrap.servers、group.id、反序列化器等),把 Kafka 连接配置也完全参数化。
这套方案在生产中的运行闭环
以运行在 K8s 上的 Flink Job 为例,完整闭环是:
- 为 Flink Job 在 K8s 上设置大量环境变量(对应各类配置项);
- Job 启动时,
createParameterTool按"配置文件 → 运行参数 → 系统属性 → 环境变量"的顺序合并,环境变量优先级最高; - 要修改配置,只需修改 K8s 上的环境变量,重启 Job 即可,全程无需重新编译打包。
当然,这个方案也有代价:重启一个作业的代价很大——重启后必须保证状态恢复到重启前的状态,尽管 Flink 的 Checkpoint 和 Savepoint 已经很强大,但对于复杂作业而言,能少一次重启就少一次。因此更理想的方向是动态获取配置:配置变更后作业能自行感知。Flink 中有的配置(如并行度)不能动态设置,但业务类配置可以,此时就要借助广播变量(Broadcast State)——本书 3.4 节 已介绍过广播变量基础用法,而 11.4 节 则通过实际案例演示如何用广播变量动态更新配置数据。
五、ParameterTool 源码分析视角:设计要点
虽然书中 10.2.3 节的源码分析部分在开源版本中未展开,但从使用方式与 API 行为可以推断ParameterTool的几个关键设计:
- 统一抽象:
fromArgs、fromSystemProperties、fromPropertiesFile、fromMap四条静态工厂方法,分别负责把不同来源的数据解析并规整到内部统一的键值存储中,上层使用方完全感知不到来源差异; - 覆盖式合并:
mergeWith(ParameterTool other)以"后者覆盖前者"的方式合并两份参数,这是实现"环境变量优先"多级配置的基础; - 序列化能力:
ParameterTool implements Serializable,因此可以作为函数构造参数随算子分发到各 TaskManager,也可以通过setGlobalJobParameters随ExecutionConfig全局广播; - 类型安全取值:
getInt/getLong/getBoolean/getRequired等方法的实现会对值做类型解析与默认值兜底,缺 key 时返回默认值、getRequired则直接抛异常,将"配置缺失"问题暴露在作业启动阶段而非运行期。
这些设计叠加起来,构成了它成为 Flink 社区配置管理事实标准的原因。建议读者在阅读源码时重点对比fromArgs对--key value、-key value与--key=value三种写法的解析差异,以及mergeWith的覆盖实现。
六、自定义配置参数类:更进一步
在实际工程中,还可以在ParameterTool之上再封装一层自定义配置参数类,把"取 key"的过程收敛成类型安全的强类型字段访问。比如:
public class JobConfig { private final ParameterTool parameterTool; private JobConfig(ParameterTool parameterTool) { this.parameterTool = parameterTool; } public static JobConfig from(ParameterTool parameterTool) { return new JobConfig(parameterTool); } public String getKafkaBrokers() { return parameterTool.get(PropertiesConstants.KAFKA_BROKERS, PropertiesConstants.DEFAULT_KAFKA_BROKERS); } public String getMetricsTopic() { return parameterTool.getRequired(PropertiesConstants.METRICS_TOPIC); } public int getParallelism() { return parameterTool.getInt(PropertiesConstants.STREAM_PARALLELISM, 5); } public boolean isCheckpointEnabled() { return parameterTool.getBoolean(PropertiesConstants.STREAM_CHECKPOINT_ENABLE, true); } }这样做的好处是:业务代码中不再散落字符串 key,所有配置项集中在单一入口(配合PropertiesConstants常量类),key 拼写错误、类型错误都在编译期暴露,同时保留了ParameterTool的序列化与全局注册能力。仓库中PropertiesConstants+ExecutionEnvUtil的组合已经体现了这一思想:key 收口到常量类,读取逻辑收口到工具类。
七、小结与反思
本章(第 10 章)实际讲了两个实践内容:一是作业的重启策略(从真实线上故障分析出发,讲解如何配置重启策略及重启策略的种类,见 flink-in-action-10.1.md),二是本文所讲的ParameterTool管理配置。两者都是贴近真实生产场景、具有一定工程价值的实践。
回到配置管理本身,最终要记住的决策要点:
| 方案 | 适用场景 | 局限 |
|---|---|---|
Configuration+withParameters | 批程序、单算子传参 | 流程序不支持;需逐算子设置,繁琐 |
ParameterTool.fromArgs | Job 启动参数 | 需在提交时显式传参 |
ParameterTool.fromSystemProperties | JVM 系统属性 | 信息量有限 |
ParameterTool.fromPropertiesFile | 固定配置文件 | 环境间切换需更换文件 |
mergeWith多层合并 | K8s/多环境部署 | 改配置仍需重启 Job |
| 广播变量动态配置 | 业务配置热更新 | 仅限业务类配置,如 3.4 节、11.4 节 |
ParameterTool是这套体系中最核心的组件:一条 API 覆盖运行参数、系统属性、配置文件、环境变量四类来源,支持全局注册与序列化传递,配合mergeWith实现"配置文件兜底 + 环境变量优先"的层级覆盖。将本文的方法应用到你的项目中,即可实现"同一个编译产物,零代码改动,无缝运行在开发、测试、预发、生产"。
如果你想进一步深入,推荐直接阅读仓库源码:ExecutionEnvUtil.java、KafkaConfigUtil.java、PropertiesConstants.java,以及三个最小可运行的示例:ParameterToolGetArgsMain.java、ParameterToolGetPropertiesMain.java、ParameterToolGetSystemMain.java。
- 示例工程
- 大数据
【免费下载链接】flink-learning
flink learning blog. http://www.54tianzhisheng.cn/ 含 Flink 入门、概念、原理、实战、性能调优、源码解析等内容。涉及 Flink Connector、Metrics、Library、DataStream API、Table API & SQL 等内容的学习案例,还有 Flink 落地应用的大型项目案例(PVUV、日志存储、百亿数据实时去重、监控告警)分享。欢迎大家支持我的专栏《大数据实时计算引擎 Flink 实战与性能优化》
相关推荐
Pixi 多环境(Multi-Environment)实战指南:在同一个 Workspace 中管理开发、测试与生产环境
Pixi 多环境(Multi Environment)实战指南:在同一个 Workspace 中管理开发、测试与生产环境 导读 Pixi 是基于 Conda 生
开发工具CLI包管理器任务调度《Flink 实战与性能优化》——Flink Job 反压(BackPressure)机制详解与定位调优实战
《Flink 实战与性能优化》——Flink Job 反压(BackPressure)机制详解与定位调优实战 本文是《Flink 实战与性能优化》第九章「Fli
示例工程大数据终极指南:Docker Compose多环境配置管理——一套配置轻松适配开发、测试、生产环境
终极指南:Docker Compose多环境配置管理——一套配置轻松适配开发、测试、生产环境 Docker Compose作为Docker官方提供的多容器应用编
云原生容器编排DevOpsCLI
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考