【问题标题】:How can I update a configuration in a Flink transformation?如何更新 Flink 转换中的配置?
【发布时间】:2022-10-14 07:49:04
【问题描述】:

给定一个 Flink 流作业,它将 map() 操作应用于流。

这个map() 操作从一些属性中读取其配置,并相应地映射数据。例如,配置指定读取属性“输入”,并使用不同的属性名称“输出”将其写入流。这已经很好了。

现在配置发生了变化,例如转换是为输出使用不同的属性名称。

因此,我正在寻找一种让所有 Flink 任务在运行时重新读取新配置的方法。

有没有可能

  • 暂停KafkaSource
  • 等到管道排空(冲洗)
  • 触发集群中的所有任务重新读取一个配置文件(协调)
  • 恢复KafkaSource

以编程方式在 Flink 中无需重新部署?

万一这很重要

  • 我目前正在使用 Flink 1.14,但我们必须尽快迁移到 1.15。
  • 作业使用检查点。
  • 作业使用 Flink 提供的KafkaSourceJdbcSinkKafkaSink
  • 还有用于 JDBC 和 InfluxDB 的其他自定义接收器

【问题讨论】:

    标签: java configuration apache-flink flink-streaming distributed-system


    【解决方案1】:

    通常,这是通过读取Stream 中的配置更改,然后使用connect 操作来完成的。这样,您可以使用 map1 函数处理您的数据流的映射,然后如果检测到配置中的任何更改,它可以在 map2 中处理并存储在状态中,您可以使 map1 函数依赖于最后收到的配置改变。

    不确定这是否可以解决您的问题,但似乎应该可以正常工作。

    【讨论】:

    • 在某些情况下,广播配置流可能更合适。
    【解决方案2】:

    通常您会广播您的配置流更改,这意味着它们将被发送到执行属性转换的操作员的每个实例。有关将规则应用于形状流的示例,请参阅https://nightlies.apache.org/flink/flink-docs-master/docs/dev/datastream/fault-tolerance/broadcast_state/,这似乎模仿了您的一些要求。

    【讨论】:

      猜你喜欢
      • 2014-05-11
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-03-22
      • 1970-01-01
      • 1970-01-01
      • 2020-07-19
      相关资源
      最近更新 更多