【发布时间】:2015-12-30 22:49:57
【问题描述】:
我们将 Storm 与 Kafka Spout 一起使用。当消息失败时,我们希望重播它们,但在某些情况下,错误的数据或代码错误会导致消息总是失败 Bolt,因此我们将进入无限重播循环。显然,当我们发现错误时,我们正在修复它们,但希望我们的拓扑通常具有容错性。重放 N 次以上后,我们如何 ack() 一个元组?
查看 Kafka Spout 的代码,我发现它旨在使用指数退避计时器和 comments on the PR 状态重试:
“spout 不会终止重试周期(我认为它不应该这样做,因为它无法报告有关发生中止请求的失败的上下文),它只处理延迟重试。一个螺栓仍然期望拓扑最终会调用 ack() 而不是 fail() 来停止循环。”
我看到了建议编写自定义 spout 的 StackOverflow 响应,但如果有推荐的方法在 Bolt 中执行此操作,我宁愿不要被困在维护 Kafka Spout 内部的自定义补丁。
在 Bolt 中执行此操作的正确方法是什么?我在元组中没有看到任何状态会显示它被重放了多少次。
【问题讨论】:
-
如果您在螺栓中有一些错误检查,您可以根据您的业务逻辑得出特定元组“坏”的结论,您可以“确认”而不是失败......所以它会不能重播.....