【发布时间】:2022-10-14 07:49:04
【问题描述】:
给定一个 Flink 流作业,它将 map() 操作应用于流。
这个map() 操作从一些属性中读取其配置,并相应地映射数据。例如,配置指定读取属性“输入”,并使用不同的属性名称“输出”将其写入流。这已经很好了。
现在配置发生了变化,例如转换是为输出使用不同的属性名称。
因此,我正在寻找一种让所有 Flink 任务在运行时重新读取新配置的方法。
有没有可能
- 暂停
KafkaSource - 等到管道排空(冲洗)
- 触发集群中的所有任务重新读取一个配置文件(协调)
- 恢复
KafkaSource
以编程方式在 Flink 中无需重新部署?
万一这很重要
- 我目前正在使用 Flink 1.14,但我们必须尽快迁移到 1.15。
- 作业使用检查点。
- 作业使用 Flink 提供的
KafkaSource、JdbcSink、KafkaSink。 - 还有用于 JDBC 和 InfluxDB 的其他自定义接收器
【问题讨论】:
标签: java configuration apache-flink flink-streaming distributed-system