【问题标题】:Correctly sending Flink state to Kafka正确地将 Flink 状态发送到 Kafka
【发布时间】:2022-02-10 07:08:57
【问题描述】:

我正在构建一个 Kafka -> Flink -> Kafka 管道,它适用于描述的“会话”数据。我的输入 Kafka 主题具有以下格式的数据,并构成 session_key 的一个会话:

start_event(session_key, some_other_data...)
entry_event(session_key, some_other_data...)
entry_event(session_key, some_other_data...)
...
entry_event(session_key, some_other_data...)
end_event(session_key, some_other_data...)

像这样的每个会话大约有 100 个事件长,很快就会出现(每 1-2 秒),所有事件共享相同的 session_key,我正在将会话转换为一系列 20 个左右的事件进入输出主题。要构建这些事件,我需要了解整个会话,因此我需要等待end_event 到达才能运行处理并将输出事件推送到输出主题。

实现相当简单 - 通过session_key 键入密钥,将start_event 存储到ValueState,将条目存储到ListState,然后当end_event 到达时,对所有事件运行处理逻辑并将结果推送到输出 Kafka 主题。

我的问题是关于检查点和可能的故障 - 假设检查点在 end_event 退出 Kafka 之后开始。偏移量已提交给 Kafka,检查点屏障到达我的处理操作员,该操作员在它之前失败(Kafka 现在已关闭)。

我应该如何正确地从中恢复?如果 Kafka 偏移量已经提交,并且没有 end_event 将永远不会因为那个 session_key 而离开 Kafka,那么我以后如何触发处理操作符来处理我保存的状态?或者在这种情况下不会提交 Kafka 偏移量,end_event 会再次通过 Flink?

【问题讨论】:

    标签: apache-kafka apache-flink flink-streaming checkpointing


    【解决方案1】:

    Flink 提交 Kafka 偏移量只是为了监控,而不依赖它们来实现容错(顺便说一句,它会在检查点完成时这样做):

    请注意,Kafka 源不依赖已提交的偏移量来实现容错。提交offset只是为了暴露consumer和consumer group的进度以便监控。

    Consumer Offset Committing

    主题偏移量在检查点期间保存为 Kafka 源状态的一部分。在所描述的场景中,整个检查点将失败,Flink 将从前一个检查点中保存的偏移量开始消费主题。不会丢失任何消息,但有些消息可能会重复(假设 AT_LEAST_ONCE 检查点模式)。

    所以是的,end_event 将再次通过 Flink。

    【讨论】:

      【解决方案2】:

      我认为kafka offset在这种情况下不会被提交,offset是在checkpoint的notify阶段提交的

      只有在所有算子的检查点都成功时才会触发通知阶段。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2019-02-04
        • 2019-07-03
        • 2011-05-24
        • 2021-03-22
        • 2020-12-17
        • 1970-01-01
        • 1970-01-01
        • 2018-11-04
        相关资源
        最近更新 更多