【发布时间】:2015-11-02 12:27:57
【问题描述】:
我正在 Storm 中实现一个 Bolt,它接收来自 RabbitMQ spout (https://github.com/ppat/storm-rabbitmq) 的消息。
我必须在 Storm 中处理的每个事件都作为来自 Rabbit 的两条消息到达,因此我在 bolt 上有一个 fieldsGrouping 以便这两条消息在同一个 bolt 中到达。
我的第一种方法是:
- 接收第一个元组并将消息保存在内存中
- 确认第一个元组
- 当第二个元组到达时,从内存中获取第一个元组并从 spout 发出一个锚定到第二个的新元组。
这行得通,但如果工人死亡,我可能会丢失消息,因为我会在获得第二个元组并处理之前确认第一个元组。
我把它改成了:
- 接收第一个元组并将其保存在内存中
- 当第二个元组到达时,从内存中获取第一个元组,发出一个锚定到两个输入元组的新元组并确认两个输入元组。
内存中的缓存是一个有时间期限的 Guava 缓存,当一个元组由于超时而被驱逐时,我将在拓扑中失败()它,以便稍后对其进行重新处理。
这似乎可行,但是当我进行一些测试时,我遇到了系统停止从 Rabbit 队列获取消息的情况。
队列上的预取设置为 5,并在 7 处使用 setMaxSpoutPending 喷出。在 Rabbit 界面中,我看到 5 条 Unacked 消息。
在风暴日志中,我看到相同的元组一遍又一遍地从缓存中驱逐。
我知道问题在于 spout 只会获取 5 条消息,这些消息都是一对的第一部分。我可以增加预取,但这不能保证这不会在生产中发生。
所以我的问题是:如何在 Storm 中处理这些问题时实现连接?
【问题讨论】:
标签: rabbitmq apache-storm