【发布时间】:2017-03-04 03:33:06
【问题描述】:
我正在使用 Kafka,我们有一个用例来构建一个容错系统,在这个系统中甚至不会遗漏任何一条消息。所以问题来了: 如果由于任何原因(ZooKeeper 宕机、Kafka 代理宕机等)导致发布到 Kafka 失败,我们如何能够稳健地处理这些消息并在事情再次备份时重播它们。正如我所说,我们甚至无法承受单个消息失败。 另一个用例是,我们还需要在任何给定时间点知道有多少消息由于任何原因未能发布到 Kafka,例如计数器功能,现在这些消息需要重新发布。
其中一个解决方案是将这些消息推送到某个数据库(例如 Cassandra,其中写入速度非常快,但我们还需要计数器功能,我猜 Cassandra 计数器功能不是那么好,我们不想使用它。)它可以处理这种负载,还为我们提供了非常准确的计数器设施。
这个问题更多是从架构的角度来看,然后是使用哪种技术来实现。
PS:我们处理一些像 3000TPS 的地方。因此,当系统开始失败时,这些失败的消息会在很短的时间内快速增长。我们正在使用基于 java 的框架。
感谢您的帮助!
【问题讨论】:
-
嗨@Nishant,你找到“解决方案”了吗?愿意与社区分享吗?提前致谢。
-
也许您需要一个仅附加数据库,例如 timescaledb 或 influxdb。对于那些每秒 3k 个事件来说,这没什么大不了的。
-
我对这个话题了解不多,但使用拉动方法而不是推动方法似乎更容易做到这一点。因此,您可以将 Web 服务添加到发送方,并且可以从接收方轮询 Web 服务。所以接收者将负责获取消息,而不是发送者或中间的其他组件将负责将其传递给所有接收者,维护接收者列表,重试等......但我想这并不总是一个选项,因为它不够快,或者可能是我不知道的其他原因。
标签: java cassandra redis apache-kafka kafka-consumer-api