【问题标题】:Apache Flink Stateful Functions forwarding the same message to N functionsApache Flink Stateful Functions 将相同的消息转发给 N 个函数
【发布时间】:2020-05-09 13:38:56
【问题描述】:

我正在尝试将传入消息发送到多个有状态函数,但我无法完全理解如何去做。为了便于理解,假设我的一个有状态函数获取一些整数并将它们发送到几个远程函数。这些函数将此整数添加到它们的状态值中,并将其保存为新状态。

当这两个远程功能之一失败时,另一个应该继续以相同的方式工作。 当失败的函数恢复时,它应该处理它在失败期间无法处理的消息。

我想像下面这样一个接一个地发送它们,但我认为它不会起作用

context.send(RemoteFuncType1,someID,someInteger);
context.send(RemoteFuncType2,someID,someInteger);
...
  • 如何以容错方式做到这一点?
  • 如果可能的话,它在后台如何工作?

【问题讨论】:

    标签: apache-flink fault-tolerance flink-statefun


    【解决方案1】:

    您建议的方式是正确的方式!

    StateFun 会以一致的方式将消息传递给远程函数。如果其中一个函数的停机时间很短,StateFun 将重试发送消息,直到:

    • 它会成功交付它(带回退)
    • 将达到重试的最大超时。当达到超时时,整个 StateFun 作业将倒退到 以前一致的检查点。

    由于 StateFun 正在管理消息传递和函数的状态(包括远程),它将确保将一致的状态和消息传递给每个函数。 在您的示例中:一旦恢复,第二个远程函数将接收 someInteger 之前具有的任何状态。

    为了更深入地了解检查点在 Flink 中的工作原理以及它如何启用恰好一次处理,我建议以下内容:

    https://ci.apache.org/projects/flink/flink-docs-stable/internals/stream_checkpointing.html

    【讨论】:

    • 对于顺序问题,恢复的遥控功能是否按整数顺序排列?
    • 是的,消息将按照先进先出的顺序传递。
    猜你喜欢
    • 1970-01-01
    • 2022-06-12
    • 2021-09-30
    • 2020-08-14
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-02-04
    相关资源
    最近更新 更多