You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

排查与解决方向

  1. 组协调器与偏移量存储状态验证

    • 检查组协调器Broker在断开期间是否发生重启、副本故障,导致消费者组元数据不可用。
    • 查看__consumer_offsets主题状态:该主题存储消费者偏移量,若其副本分布在断开的Broker上,会导致偏移量读取失败。需确认该主题的分区副本分配、ISR状态是否正常。
  2. 偏移量提交可靠性检查

    • 验证commitOffsets重试逻辑是否覆盖所有异常:当前代码仅捕获Try包裹的异常,需确认x.commitScaladsl()是否会抛出未被捕获的非检查型异常,导致提交失败且未重试。
    • 增加提交成功日志:在偏移量提交成功后记录具体的偏移量值,确认提交操作是否真的写入__consumer_offsets。
  3. 消费者配置优化

    • 调整auto.offset.reset配置:临时改为none避免自动重置为最早偏移量,同时在代码中增加偏移量不存在时的自定义处理(比如从业务数据库读取已处理的最大偏移量)。
    • 缩短metadata.max.age.ms:加快元数据刷新速度,减少组协调器重复发现的次数。
    • 确认enable.auto.commit为false:使用手动提交时必须关闭自动提交,避免偏移量提交冲突。
  4. 依赖版本排查

    • 检查akka-stream-kafka版本是否存在已知的偏移量提交或组协调器发现BUG,尝试升级到最新稳定版本。

如需进一步定位,可提供以下信息:

  • Kafka集群版本
  • 消费者consumerSettings的完整配置参数
  • __consumer_offsets主题的分区数、副本配置及状态
  • 断开期间组协调器节点的Broker日志

内容的提问来源于stack exchange,提问作者Zixel

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.04 17:47:03