【问题标题】:Apache Storm Join Pattern - At least onceApache Storm Join Pattern - 至少一次
【发布时间】:2015-11-02 12:27:57
【问题描述】:

我正在 Storm 中实现一个 Bolt,它接收来自 RabbitMQ spout (https://github.com/ppat/storm-rabbitmq) 的消息。

我必须在 Storm 中处理的每个事件都作为来自 Rabbit 的两条消息到达,因此我在 bolt 上有一个 fieldsGrouping 以便这两条消息在同一个 bolt 中到达。

我的第一种方法是:

  1. 接收第一个元组并将消息保存在内存中
  2. 确认第一个元组
  3. 当第二个元组到达时,从内存中获取第一个元组并从 spout 发出一个锚定到第二个的新元组。

这行得通,但如果工人死亡,我可能会丢失消息,因为我会在获得第二个元组并处理之前确认第一个元组。

我把它改成了:

  1. 接收第一个元组并将其保存在内存中
  2. 当第二个元组到达时,从内存中获取第一个元组,发出一个锚定到两个输入元组的新元组并确认两个输入元组。

内存中的缓存是一个有时间期限的 Guava 缓存,当一个元组由于超时而被驱逐时,我将在拓扑中失败()它,以便稍后对其进行重新处理。

这似乎可行,但是当我进行一些测试时,我遇到了系统停止从 Rabbit 队列获取消息的情况。

队列上的预取设置为 5,并在 7 处使用 setMaxSpoutPending 喷出。在 Rabbit 界面中,我看到 5 条 Unacked 消息。

在风暴日志中,我看到相同的元组一遍又一遍地从缓存中驱逐。

我知道问题在于 spout 只会获取 5 条消息,这些消息都是一对的第一部分。我可以增加预取,但这不能保证这不会在生产中发生。

所以我的问题是:如何在 Storm 中处理这些问题时实现连接?

【问题讨论】:

    标签: rabbitmq apache-storm


    【解决方案1】:

    Storm 没有为此提供好的解决方案...您需要的是一个 可靠的 存储来缓冲第一个元组(即,有状态的运算符)。因此,您可以立即确认第一个元组并在失败后恢复状态。

    1. 据我所知,Trident 支持一些状态处理。但我从来没有用过。
    2. 作为第二种选择,您可以使用分布式键值存储(如 Casandra)作为缓冲区。当然,这将是一个手写的解决方案,即您需要自己编写所有 Casandra 交互的代码。
    3. 最后但同样重要的是,您可以切换到支持 Apache Flink 等有状态操作符的流处理系统。 (免责声明:我是 Flink 的提交者)

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2015-12-25
      • 1970-01-01
      • 1970-01-01
      • 2021-09-11
      • 1970-01-01
      • 2022-01-13
      相关资源
      最近更新 更多