Kafka消费者组重连后丢失已提交偏移量,重置为最早偏移量问题
Kafka消费者组重连后丢失已提交偏移量问题
问题描述
应用因网络问题与5个Kafka Broker中的3个断开连接,出现如下错误:
Error sending fetch request (sessionId=2108009367, epoch=INITIAL) to node 4: org.apache.kafka.common.errors.DisconnectException. Error sending fetch request (sessionId=1169924518, epoch=INITIAL) to node 1: org.apache.kafka.common.errors.DisconnectException. Error sending fetch request (sessionId=1156529071, epoch=INITIAL) to node 5: org.apache.kafka.common.errors.DisconnectException.
随后大量重复出现组协调器不可用日志:
Group coordinator <host>:9093 (id: 2147483642 rack: null) is unavailable or invalid, will attempt rediscovery Discovered group coordinator <host>:9093 (id: 2147483642 rack: null)
多次重试后消费者组成功重新加入,但丢失了已提交的偏移量,被重置为分区最早偏移量:
Revoking previously assigned partitions [<topic>-0, <topic>-1, <topic>-2] (Re-)joining group Successfully joined group with generation 1 Setting newly assigned partitions: <topic>-0, <topic>-1, <topic>-2 Found no committed offset for partition <topic>-0 Found no committed offset for partition <topic>-1 Found no committed offset for partition <topic>-2 Resetting offset for partition <topic>-0 to offset 3812363 (first within retention period).
该问题导致消费者重新处理保留期内大量数据,业务影响严重。
补充信息
消费者代码
def receive: Receive = { case job: Job => Consumer.committableSource(consumerSettings, Subscriptions.topics(topic)) .collect { case msg@JsonTx(tx) => (msg, tx) } .mapAsync(streamParallelism) { case (msg, tx) => currencyNormalizer.normalize(tx).map { msg -> _ } } .mapAsync(streamParallelism) { case (kafkaMsg, message) => //build tx and persist to DB } .map { case (kafkaMsg, message) => //send to other topic ProducerMessage.Message(new ProducerRecord[String, String](configService.topic, message.toString()), kafkaMsg.committableOffset) } .via(Producer.flow(producerSettings)) .map(_.message.passThrough) .batch(max = 100, first => CommittableOffsetBatch.empty.updated(first)) { (batch, elem) => batch.updated(elem) } .mapAsync(streamParallelism) { msg => commitOffsets(msg) } .withAttributes(ActorAttributes.supervisionStrategy(decider)) .runWith(Sink.ignore) () }
偏移量提交代码
private def commitOffsets(x: ConsumerMessage.CommittableOffsetBatch): Future[akka.Done] = { retry .Backoff(max = maxRetries, delay = offsetRetryDelay, base = offsetBackoffBase) .apply { Try(x.commitScaladsl()) match { case Success(value) => value case Failure(e) => logger.warn("Failed to commit offsets. Retrying...", e) throw e } }(retry.Success[akka.Done] {_ => true}, ec) }
依赖信息
"com.typesafe.akka" %% "akka-stream-kafka" % Versions.akkaKafka
排查与解决方向
组协调器与偏移量存储状态验证
- 检查组协调器Broker在断开期间是否发生重启、副本故障,导致消费者组元数据不可用。
- 查看
__consumer_offsets主题状态:该主题存储消费者偏移量,若其副本分布在断开的Broker上,会导致偏移量读取失败。需确认该主题的分区副本分配、ISR状态是否正常。
偏移量提交可靠性检查
- 验证
commitOffsets重试逻辑是否覆盖所有异常:当前代码仅捕获Try包裹的异常,需确认x.commitScaladsl()是否会抛出未被捕获的非检查型异常,导致提交失败且未重试。 - 增加提交成功日志:在偏移量提交成功后记录具体的偏移量值,确认提交操作是否真的写入
__consumer_offsets。
- 验证
消费者配置优化
- 调整
auto.offset.reset配置:临时改为none避免自动重置为最早偏移量,同时在代码中增加偏移量不存在时的自定义处理(比如从业务数据库读取已处理的最大偏移量)。 - 缩短
metadata.max.age.ms:加快元数据刷新速度,减少组协调器重复发现的次数。 - 确认
enable.auto.commit为false:使用手动提交时必须关闭自动提交,避免偏移量提交冲突。
- 调整
依赖版本排查
- 检查
akka-stream-kafka版本是否存在已知的偏移量提交或组协调器发现BUG,尝试升级到最新稳定版本。
- 检查
如需进一步定位,可提供以下信息:
- Kafka集群版本
- 消费者
consumerSettings的完整配置参数 __consumer_offsets主题的分区数、副本配置及状态- 断开期间组协调器节点的Broker日志
内容的提问来源于stack exchange,提问作者Zixel
相关产品推荐
相关产品推荐

