【问题标题】:Recommended approach to fatal technical errors during Kafka processing in a partition分区中 Kafka 处理期间致命技术错误的推荐方法
【发布时间】:2020-12-15 04:14:50
【问题描述】:
我对在 Kafka 流式传输期间处理致命技术错误的推荐方法有疑问。
场景:
- 业务事务的所有消息都在一个分区中
- 重复处理消息会导致应用出现异常(与 Kafka 无关)
- 我们不能跳过消息,因为顺序很重要(一个交易的所有消息需要一起处理)
- Kafka 自动将分区分配给消费者(无需手动分配)
鉴于这些限制,
- 如果我停止消费者,那么带有问题消息的分区只会被分配给不同的消费者,并且会重复同样的问题。
- 如果我停止整个消费者组,我将延迟处理所有分区,而如果它们仍在处理,我可以处理其他没有问题的事务。
对于这种情况,推荐的方法是什么?
另外,是否可以在没有应用程序同步机制的情况下以某种方式关闭整个消费者组(对于多节点消费者组)?
【问题讨论】:
标签:
apache-kafka
kafka-consumer-api
【解决方案1】:
Kafka 流只有一次语义,即一条消息只会被处理一次,即使在失败的情况下也是如此。您只需要在您的流配置中设置processing.guarantee=exactly_once。更多信息请参考following article。
关于消息的顺序,按照产生的顺序被消费。由于事务中所有消息的键都相同,因此它们位于同一个分区中(至少在默认分区逻辑下)。
您还可以利用 中间主题 使用 KStream.through() 方法存储您的中间状态。
除此之外,您还有状态存储(本地(仅分配给该实例的分区)和全局(所有分区))来存储您的状态。
如果任何一个消费者死亡并被重新部署到其他地方(在另一个没有存储数据的节点中),那么状态存储将从 Kafka 更改日志中构建。