【发布时间】:2021-03-12 05:17:57
【问题描述】:
我正在尝试构建一个 Flink 作业,该作业将从 Kafka 源读取数据进行一系列处理,包括少量 REST 调用,然后最终陷入另一个 Kafka 主题。
我试图解决的问题是消息重试。如果 REST API 中存在暂时错误怎么办?如何像 Storm 支持的那样对这些消息进行基于指数退避的重试?
我有两种方法可以考虑
- 使用 TimerService,但如果发生故障,状态将开始不受控制地扩展。
- 将失败的消息写入不同的 Kafka 主题并延迟处理它们,但如果 Sink 本身停机几分钟就会出现问题?
有没有更好、更健壮、更简单的方法来实现这一点?
【问题讨论】:
-
对于方法 #2,您的 Kafka 集群是否经常停机,以至于如果 sink 停机然后从检查点重新启动,仅仅依靠 Flink 失败是不够的?
标签: apache-flink flink-streaming