【发布时间】:2014-03-01 04:14:16
【问题描述】:
我正在使用storm来处理在线问题,但我不明白为什么storm会从spout重播tuple。重试崩溃的内容可能比从根目录重播更有效,对吧? 任何人都可以帮助我吗?谢谢
【问题讨论】:
标签: java apache-storm
我正在使用storm来处理在线问题,但我不明白为什么storm会从spout重播tuple。重试崩溃的内容可能比从根目录重播更有效,对吧? 任何人都可以帮助我吗?谢谢
【问题讨论】:
标签: java apache-storm
典型的 spout 实现将仅重放 FAILED 元组。正如here 解释的那样,从 spout 发出的元组可以触发数千个其他元组,storm 基于此创建一个元组树。现在,当树中的每条消息都已被处理时,元组被称为“完全处理”。在发出 spout 时添加一个message id,用于在后期识别元组。这称为锚定,可以通过以下方式完成
_collector.emit(new Values("field1", "field2", 3) , msgId);
现在从上面发布的链接中可以看出
当一个元组的消息树在指定的超时时间内没有被完全处理时,元组被认为是失败的。可以使用 Config.TOPOLOGY_MESSAGE_TIMEOUT_SECS 配置在特定于拓扑的基础上配置此超时,默认为 30 秒。
如果元组超时,Storm 将在 spout 上调用 FAIL 方法,同样在成功的情况下,将调用 ACK 方法。
所以此时storm会告诉你哪些是它处理失败的元组,但是如果你查看源代码你会发现fail方法的实现在BaseRichSpout中是空的类,因此您需要重写 BaseRichSpout 的 fail 方法,以便在您的应用程序中具有重放功能。
【讨论】:
这种失败元组的重放应该只占整个元组流量的一小部分,因此这种简单的重放从开始策略的效率通常不是问题。
支持“replay-from-error-step”会带来很多复杂性,因为错误的位置有时很难确定,并且需要支持“replay-elsewhere”以防出现错误的集群节点发生当前(或永久)关闭。它还会减慢整个流量的执行速度,这可能无法通过错误处理获得的效率来补偿(再次假设很少触发)。
如果您认为这种从头开始重放的策略会对您的拓扑结构产生负面影响,请尝试将其分解为几个较小的部分,并由一些持久的队列系统(如 Kafka)隔开。
【讨论】: