【发布时间】:2019-07-02 12:18:09
【问题描述】:
昨天我从日志中发现,在 Kafka 组协调器启动组重新平衡后,kafka 正在重新消费一些消息。这些消息已在两天前被消费(从日志中确认)。
日志中报告了另外两个重新平衡,但它们不再重新使用消息。那么为什么第一次重新平衡会导致重复消费消息呢?有什么问题?
我正在使用 golang kafka 客户端。这是代码
config := sarama.NewConfig()
config.Version = version
config.Consumer.Offsets.Initial = sarama.OffsetOldest
我们在声明消息之前处理消息,所以我们似乎正在使用 kafka 的“至少发送一次”策略。我们在一台机器上有三个代理,而在另一台机器上只有一个消费者线程(goroutine)。
对这个现象有什么解释吗? 我认为这些消息一定已经提交了,因为它们是在两天前被消费的,或者为什么 kafka 会在不提交的情况下保持偏移量超过两天?
消费代码示例:
func (consumer *Consumer) ConsumeClaim(session
sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error {
for message := range claim.Messages() {
realHanlder(message) // consumed data here
session.MarkMessage(message, "") // mark offset
}
return nil
}
添加:
重新平衡在应用重新启动后发生。还有另外两次重启并没有导致重新消费
-
kafka 的配置
log.retention.check.interval.ms=300000
log.retention.hours=168
zookeeper.connection.timeout.ms=6000
group.initial.rebalance.delay.ms=0
delete.topic.enable = true
auto.create.topics.enable=false
【问题讨论】:
-
当您使用最旧的偏移量时,您将从最旧的偏移量中获取您尚未提交的消息。你能分享一下你代码的消费阶段吗?
-
您的服务器保留政策是什么?再平衡过程中你的群体识别是否发生了变化?
标签: go apache-kafka kafka-consumer-api sarama