【发布时间】: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