【问题标题】:Flink message retries like StormFlink 消息重试,如 Storm
【发布时间】:2021-03-12 05:17:57
【问题描述】:

我正在尝试构建一个 Flink 作业,该作业将从 Kafka 源读取数据进行一系列处理,包括少量 REST 调用,然后最终陷入另一个 Kafka 主题。

我试图解决的问题是消息重试。如果 REST API 中存在暂时错误怎么办?如何像 Storm 支持的那样对这些消息进行基于指数退避的重试?

我有两种方法可以考虑

  1. 使用 TimerService,但如果发生故障,状态将开始不受控制地扩展。
  2. 将失败的消息写入不同的 Kafka 主题并延迟处理它们,但如果 Sink 本身停机几分钟就会出现问题?

有没有更好、更健壮、更简单的方法来实现这一点?

【问题讨论】:

  • 对于方法 #2,您的 Kafka 集群是否经常停机,以至于如果 sink 停机然后从检查点重新启动,仅仅依靠 Flink 失败是不够的?

标签: apache-flink flink-streaming


【解决方案1】:

我会使用 Flink 的 AsyncFunction 进行 REST 调用。如果需要,它将对源进行反压,而不是使用超过配置数量的状态。如需重试,请参阅AsyncFunction retries

【讨论】:

  • 谢谢,但是我必须放弃一致性,对吗?出于实际原因,不能将 REST API 设为链中的第一个调用。
  • AsyncFunction 更适合进行查找而不是更新外部系统。它提供至少一次保证。
  • 根据文档,AsyncIO 只能用于操作链的头部。我理解这意味着它必须是源流上的第一个运算符,对吗?
  • 该限制不再成立;见stackoverflow.com/a/66705952/2000823。但是不,这绝不意味着它必须是工作中的第一个操作员,而是它必须是其任务中的第一个操作员。
猜你喜欢
  • 1970-01-01
  • 2017-01-20
  • 2016-06-29
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-02-23
  • 1970-01-01
  • 2015-09-20
相关资源
最近更新 更多