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

开启自动提交的Kafka客户端关闭时提交未消费的最新生产消息偏移量是否为预期行为?

Issue: Auto-committing offsets includes unconsumed self-produced messages in Akka Kafka

TLDR: Is it expected behavior for Kafka clients with auto-commit enabled (in apps that consume and produce to the same topic) to mark offsets of unconsumed self-produced messages as committed?

Problem Description

I built a simple Scala application that uses Akka Actors to consume messages from a Kafka topic. When exceptions occur during message processing, the application re-produces the message to the same topic. The consumer runs in cycles: it starts every minute, runs for 20 seconds, then stops.

Relevant Code Snippets

TestActor.scala (Message Processing Logic)

override protected def processMessage(messages: Seq[ConsumerRecord[String, String]]): Future[Done] = {
  Future.sequence(messages.map(message => {
    logger.info(s"--CONSUMED: offset: ${message.offset()} message: ${message.value()}")
    // In production, this runs processing logic; on exception, sends to the same topic
    sendToExceptionTopic(Instant.now().toEpochMilli)
    Thread.sleep(1000)
    Future(Done)
  })).transformWith(_ => Future(Done))
}

Starter.scala (Scheduling Logic)

def init(): Unit = {
  exceptionManagerActor ! InitExceptionActors
  system.scheduler.schedule(2.second, 60.seconds) {
    logger.info("started consuming messages")
    exceptionManagerActor ! ConsumeExceptions
  }
}

ExceptionManagerActor.scala (Actor Lifecycle Scheduling)

private def startScheduledActor(actorRef: ActorRef): Unit = {
  actorRef ! Start
  context.system.scheduler.scheduleOnce(20.seconds) {
    logger.info("stopping consuming messages")
    actorRef ! Stop
  }
}

BaseActorWithAutoCommit.scala (Stream Start/Stop Logic)

override def receive: Receive = {
  case Start =>
    consumerBase = consumer
      .groupedWithin(20, 2000.millisecond)
      .mapAsyncUnordered(10)(processMessage)
      .toMat(Sink.seq)(DrainingControl.apply)
      .run()
  case Stop =>
    consumerBase.drainAndShutdown().transformWith {
      case Success(value) =>
        logger.info("actor stopped")
        Future(value)
      case Failure(ex) =>
        logger.error("error: ", ex)
        Future.failed(ex)
    }
  //Await.result(consumerBase.drainAndShutdown(), 1.minute)
}

Example Logs

14:28:48.868 INFO - started consuming messages
14:28:50.945 INFO - --CONSUMED: offset: 97 message: 1
14:28:51.028 INFO - ----PRODUCED: offset: 98 message: 1643542130945
...
14:29:08.886 INFO - stopping consuming messages
14:29:08.891 INFO - --CONSUMED: offset: 106 message: 1643542147106
14:29:08.895 INFO - ----PRODUCED: offset: 107 message: 1643542148891 <------ This message is lost
14:29:39.946 INFO - actor stopped
14:29:39.956 INFO - Message [akka.kafka.internal.KafkaConsumerActor$Internal$StopFromStage] from Actor[akka://test-consumer/system/Materializers/StreamSupervisor-2/$$a#1541548736] to Actor[akka://test-consumer/system/kafka-consumer-1#914599016] was not delivered. [1] dead letters encountered.
14:29:48.866 INFO - started consuming messages <----- Expected to consume offset 107 here, but it's skipped
14:30:08.871 INFO - stopping consuming messages
14:30:38.896 INFO - actor stopped

Dependencies

lazy val versions = new {
  val akka = "2.6.13"
  val akkaHttp = "10.1.9"
  val alpAkka = "2.0.7"
  val logback = "1.2.3"
  val apacheCommons = "1.7"
  val json4s = "3.6.7"
}
libraryDependencies ++= {
  Seq(
    "com.typesafe.akka" %% "akka-slf4j" % versions.akka,
    "com.typesafe.akka" %% "akka-stream-kafka" % versions.alpAkka,
    "com.typesafe.akka" %% "akka-http" % versions.akkaHttp,
    "com.typesafe.akka" %% "akka-protobuf" % versions.akka,
    "com.typesafe.akka" %% "akka-stream" % versions.akka,
    "ch.qos.logback" % "logback-classic" % versions.logback,
    "org.json4s" %% "json4s-jackson" % versions.json4s,
    "org.apache.commons" % "commons-text" % versions.apacheCommons,
  )
}

Root Cause Analysis

This issue stems from how Kafka's auto-commit mechanism interacts with Akka Kafka's stream shutdown logic:

  1. When your app produces a message to the same topic, the topic's highest offset is immediately updated on the broker.
  2. When calling drainAndShutdown(), Akka Kafka's auto-commit submits the current highest offset known to the consumer—this includes the unconsumed message you just produced, because the consumer's position is synced to the topic's latest offset (even if the message hasn't been pulled or processed).
  3. The dead letter log suggests the consumer actor terminated before properly handling offset commit logic, exacerbating the problem.

Solutions

Auto-commit is risky for scenarios where you consume and produce to the same topic. Manual commits let you precisely control which offsets are marked as consumed:

  • Update consumer config: enable.auto.commit = false
  • Use CommittableSource to handle offset commits only after successful message processing:
val committableSource = Consumer.committableSource(consumerSettings, Subscriptions.topics(topic))
committableSource
  .groupedWithin(20, 2000.millisecond)
  .mapAsyncUnordered(10) { msgs =>
    val processTasks = msgs.map { msg =>
      logger.info(s"--CONSUMED: offset: ${msg.record.offset()} message: ${msg.record.value()}")
      sendToExceptionTopic(Instant.now().toEpochMilli)
      Thread.sleep(1000)
      Future.successful(msg.committableOffset)
    }
    Future.sequence(processTasks).map { offsets =>
      offsets.foreach(_.commitScaladsl())
      Done
    }
  }
  .toMat(Sink.seq)(DrainingControl.apply)
  .run()

2. Separate Consumption and Production Topics

If your workflow allows, use a dedicated retry topic for failed messages instead of reusing the same topic. This eliminates offset interference entirely, makes monitoring easier, and clarifies business logic.

3. Adjust Shutdown Logic to Commit Only Processed Offsets

If you must keep auto-commit:

  • Pause production before triggering shutdown to avoid new messages being added mid-shutdown.
  • Ensure drainAndShutdown() completes fully before allowing the consumer to terminate, so only offsets for processed messages are committed.

Summary

This is not expected business behavior, but it's a side effect of Kafka's auto-commit mechanism in this specific consume-produce-to-same-topic scenario. Manual offset commits or topic separation will resolve the message loss issue.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.28 11:17:34