开启自动提交的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:
- When your app produces a message to the same topic, the topic's highest offset is immediately updated on the broker.
- 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). - The dead letter log suggests the consumer actor terminated before properly handling offset commit logic, exacerbating the problem.
Solutions
1. Disable Auto-Commit, Use Manual Offset Commit (Recommended)
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
CommittableSourceto 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

