Spring KafkaListener断连停消费:CommitFailedException恢复方案
问题描述
我有一个有意不清理消息的Kafka Topic,期望消费者即使离线数天/周/月后仍能持续消费。运行时触发CommitFailedException,之后@KafkaListener不再消费topic1的新消息,即便该Topic持续有消息写入。尝试重启Spring Boot应用、执行偏移量重置命令:
kafka-consumer-groups.sh --bootstrap-server localhost:9092 --group aggregator-txnstokafka-group-ethereum --topic topic1 --reset-offsets --to-earliest --execute
但消费仍未恢复。请问如何让@KafkaListener恢复消息消费?
异常信息
2023-10-05 11:14:31,755 ERROR org.springframework.kafka.core.DefaultKafkaProducerFactory [aggregator-txnstokafka-listener-ethereum-0-C-1] - commitTransaction failed: CloseSafeProducer [delegate=org.apache.kafka.clients.producer.KafkaProducer@7de7fc0] org.apache.kafka.common.KafkaException: Cannot execute transactional method because we are in an error state at org.apache.kafka.clients.producer.internals.TransactionManager.maybeFailWithError(TransactionManager.java:1125) at org.apache.kafka.clients.producer.internals.TransactionManager.lambda$beginCommit$2(TransactionManager.java:373) at org.apache.kafka.clients.producer.internals.TransactionManager.handleCachedTransactionRequestResult(TransactionManager.java:1231) at org.apache.kafka.clients.producer.internals.TransactionManager.beginCommit(TransactionManager.java:372) at org.apache.kafka.clients.producer.KafkaProducer.commitTransaction(KafkaProducer.java:755) at org.springframework.kafka.core.DefaultKafkaProducerFactory$CloseSafeProducer.commitTransaction(DefaultKafkaProducerFactory.java:1156) at org.springframework.kafka.core.KafkaResourceHolder.commit(KafkaResourceHolder.java:58) at org.springframework.kafka.transaction.KafkaTransactionManager.doCommit(KafkaTransactionManager.java:186) at org.springframework.transaction.support.AbstractPlatformTransactionManager.processCommit(AbstractPlatformTransactionManager.java:743) at org.springframework.transaction.support.AbstractPlatformTransactionManager.commit(AbstractPlatformTransactionManager.java:711) at org.springframework.transaction.support.TransactionTemplate.execute(TransactionTemplate.java:152) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeInTransaction(KafkaMessageListenerContainer.java:2387) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeRecordListenerInTx(KafkaMessageListenerContainer.java:2353) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeRecordListener(KafkaMessageListenerContainer.java:2329) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeListener(KafkaMessageListenerContainer.java:2003) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeIfHaveRecords(KafkaMessageListenerContainer.java:1373) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.pollAndInvoke(KafkaMessageListenerContainer.java:1364) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.run(KafkaMessageListenerContainer.java:1255) at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:539) at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264) at java.base/java.lang.Thread.run(Thread.java:833) Caused by: org.apache.kafka.clients.consumer.CommitFailedException: Transaction offset Commit failed due to consumer group metadata mismatch: The coordinator is not aware of this member. at org.apache.kafka.clients.producer.internals.TransactionManager$TxnOffsetCommitHandler.handleResponse(TransactionManager.java:1771) at org.apache.kafka.clients.producer.internals.TransactionManager$TxnRequestHandler.onComplete(TransactionManager.java:1322) at org.apache.kafka.clients.ClientResponse.onComplete(ClientResponse.java:109) at org.apache.kafka.clients.NetworkClient.completeResponses(NetworkClient.java:583) at org.apache.kafka.clients.NetworkClient.poll(NetworkClient.java:575) at org.apache.kafka.clients.producer.internals.Sender.maybeSendAndPollTransactionalRequest(Sender.java:418) at org.apache.kafka.clients.producer.internals.Sender.runOnce(Sender.java:316) at org.apache.kafka.clients.producer.internals.Sender.run(Sender.java:243) ... 1 more 2023-10-05 11:14:31,760 ERROR org.springframework.kafka.listener.KafkaMessageListenerContainer [aggregator-txnstokafka-listener-ethereum-0-C-1] - Transaction rolled back org.apache.kafka.common.KafkaException: Cannot execute transactional method because we are in an error state at org.apache.kafka.clients.producer.internals.TransactionManager.maybeFailWithError(TransactionManager.java:1125) at org.apache.kafka.clients.producer.internals.TransactionManager.lambda$beginCommit$2(TransactionManager.java:373) at org.apache.kafka.clients.producer.internals.TransactionManager.handleCachedTransactionRequestResult(TransactionManager.java:1231) at org.apache.kafka.clients.producer.internals.TransactionManager.beginCommit(TransactionManager.java:372) at org.apache.kafka.clients.producer.KafkaProducer.commitTransaction(KafkaProducer.java:755) at org.springframework.kafka.core.DefaultKafkaProducerFactory$CloseSafeProducer.commitTransaction(DefaultKafkaProducerFactory.java:1156) at org.springframework.kafka.core.KafkaResourceHolder.commit(KafkaResourceHolder.java:58) at org.springframework.kafka.transaction.KafkaTransactionManager.doCommit(KafkaTransactionManager.java:186) at org.springframework.transaction.support.AbstractPlatformTransactionManager.processCommit(AbstractPlatformTransactionManager.java:743) at org.springframework.transaction.support.AbstractPlatformTransactionManager.commit(AbstractPlatformTransactionManager.java:711) at org.springframework.transaction.support.TransactionTemplate.execute(TransactionTemplate.java:152) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeInTransaction(KafkaMessageListenerContainer.java:2387) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeRecordListenerInTx(KafkaMessageListenerContainer.java:2353) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeRecordListener(KafkaMessageListenerContainer.java:2329) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeListener(KafkaMessageListenerContainer.java:2003) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeIfHaveRecords(KafkaMessageListenerContainer.java:1373) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.pollAndInvoke(KafkaMessageListenerContainer.java:1364) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.run(KafkaMessageListenerContainer.java:1255) at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:539) at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264) at java.base/java.lang.Thread.run(Thread.java:833) Caused by: org.apache.kafka.clients.consumer.CommitFailedException: Transaction offset Commit failed due to consumer group metadata mismatch: The coordinator is not aware of this member. at org.apache.kafka.clients.producer.internals.TransactionManager$TxnOffsetCommitHandler.handleResponse(TransactionManager.java:1771) at org.apache.kafka.clients.producer.internals.TransactionManager$TxnRequestHandler.onComplete(TransactionManager.java:1322) at org.apache.kafka.clients.ClientResponse.onComplete(ClientResponse.java:109) at org.apache.kafka.clients.NetworkClient.completeResponses(NetworkClient.java:583) at org.apache.kafka.clients.NetworkClient.poll(NetworkClient.java:575) at org.apache.kafka.clients.producer.internals.Sender.maybeSendAndPollTransactionalRequest(Sender.java:418) at org.apache.kafka.clients.producer.internals.Sender.runOnce(Sender.java:316) at org.apache.kafka.clients.producer.internals.Sender.run(Sender.java:243) ... 1 more
服务器配置
offsets.retention.minutes=5256000 group.max.session.timeout.ms=2147483647
客户端配置
session.timeout.ms=2147483646 max.poll.interval.ms=900000
相关代码片段
(从topic1拉取消息发送至topic2)
@Autowired @Qualifier("evmBlockToKakfaTemplate") KafkaTemplate<Number, Object> destTemplate; @KafkaListener(id = "aggregator-txnstokafka-listener-ethereum", groupId = "aggregator-txnstokafka-group-ethereum", topics = "topic1", containerFactory = "evmTxnsToKafkaContainerFactory") public void onMessage(ConsumerRecord<Long, String> record, Acknowledgement acknowledgement) { // Do some stuff destTemplate.send("topic2", key, value); } @Bean public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, String>> evmTxnsToKafkaContainerFactory( KafkaTransactionManager<Number, Object> evmTxnAggregatorTransactionManager, ConsumerFactory<? super String, ? super String> evmTxnToKakfaConsumerFactory) { ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.getContainerProperties().setTransactionManager(evmTxnAggregatorTransactionManager); factory.setConsumerFactory(evmTxnToKakfaConsumerFactory); factory.getContainerProperties().setAckMode(AckMode.MANUAL_IMMEDIATE); factory.setCommonErrorHandler(new TransactionAggregatorErrorHandler()); factory.setConcurrency(1); factory.setBatchListener(false); return factory; } @Bean public KafkaTemplate<Number, Object> evmTxnToKakfaTemplate() { return new KafkaTemplate<Number, Object>(evmTxnToKakfaProducerFactory()); } @Bean public DefaultKafkaProducerFactory<Number, Object> evmTxnToKakfaProducerFactory() { Map<String, Object> config = new HashMap<>(); config.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, destKafkaBroker); config.put(ProducerConfig.CLIENT_ID_CONFIG, "aggregator-txnstokafka-" + chain); config.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, LongSerializer.class.getName()); config.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class.getName()); config.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true); config.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "aggregator-txnstokafka-" + chain + "-"); return new DefaultKafkaProducerFactory<>(config); } @Bean public ConsumerFactory<? super String, ? super BlockchainDataAggregator> evmTxnToKakfaConsumerFactory() { Map<String, Object> config = new HashMap<>(); config.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, srcKafkaBroker); config.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, LongDeserializer.class.getName()); config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class.getName()); config.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false"); config.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, IsolationLevel.READ_COMMITTED.toString().toLowerCase(Locale.ROOT)); config.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); config.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, "2147483646"); // Integer max value. About 24.8 days //config.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, ...); config.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, "900000"); // 15 minutes return new DefaultKafkaConsumerFactory<>(config); } @Bean KafkaTransactionManager<Number, Object> evmTxnAggregatorTransactionManager( DefaultKafkaProducerFactory<Number, Object> evmTxnToKakfaProducerFactory) { KafkaTransactionManager<Number, Object> kafkaTransactionManager = new KafkaTransactionManager<>( evmTxnToKakfaProducerFactory); kafkaTransactionManager .setTransactionSynchronization(AbstractPlatformTransactionManager.SYNCHRONIZATION_ALWAYS); return kafkaTransactionManager; }
使用环境
- Ubuntu上的Kafka 2.12-3.2.0
- Java 18
- Spring Boot Starter 2.7.4(集成Spring Kafka)
- topic1为1分区3副本(有意配置)
解决方案
问题根源分析
- 事务状态异常:生产者事务管理器因之前的提交失败被标记为错误状态,后续所有事务操作直接拒绝执行。
- 消费者组元数据不匹配:超大的
session.timeout.ms与max.poll.interval.ms不匹配,消费者离线超过15分钟后被协调器踢出组,但事务提交偏移量时仍使用旧的组元数据,触发不匹配错误。 - 偏移量重置无效:事务模式下,偏移量通过生产者事务提交,直接重置消费者组偏移量无法覆盖事务内的状态记录。
分步恢复步骤
1. 终止异常事务
事务ID在Kafka集群中留存的异常状态会阻塞后续操作,需先终止:
# 替换为实际的transactional.id(去掉末尾多余的-) kafka-transactions.sh --bootstrap-server localhost:9092 --transactional-id aggregator-txnstokafka-ethereum --abort
2. 清理消费者组元数据并重置偏移量
确保应用处于停止状态,执行以下命令:
# 重置偏移量到最早位置 kafka-consumer-groups.sh --bootstrap-server localhost:9092 --group aggregator-txnstokafka-group-ethereum --topic topic1 --reset-offsets --to-earliest --execute # 删除消费者组,强制协调器清除旧元数据 kafka-consumer-groups.sh --bootstrap-server localhost:9092 --group aggregator-txnstokafka-group-ethereum --delete
3. 修复代码配置缺陷
(1)修正事务ID格式
移除事务ID末尾多余的-,避免状态识别异常:
config.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "aggregator-txnstokafka-" + chain);
(2)同步超时参数配置
设置heartbeat.interval.ms为session.timeout.ms的1/3,避免协调器误判消费者离线:
config.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, "715827882"); // 约8.2天
(3)增强错误处理器的事务恢复能力
在TransactionAggregatorErrorHandler中添加事务重置逻辑:
@Override public void handleOtherException(Exception thrownException, Consumer<?, ?> consumer, MessageListenerContainer container) { // 重置生产者事务状态 DefaultKafkaProducerFactory<?, ?> producerFactory = (DefaultKafkaProducerFactory<?, ?>) ((KafkaTransactionManager<?, ?>) container.getContainerProperties().getTransactionManager()).getProducerFactory(); producerFactory.reset(); // 关闭消费者触发重新加入组 consumer.close(Duration.ofSeconds(10)); }
4. 重启应用验证
启动Spring Boot应用,观察日志中是否出现消费者组重新加入的记录,同时检查topic1的消息是否开始被消费、topic2是否收到转发消息。
长期优化建议
- 避免设置过大的
session.timeout.ms,依赖offsets.retention.minutes保留偏移量即可,同时将max.poll.interval.ms调整为合理值(如2小时),处理慢消息时手动调用pause()/resume()。 - 事务模式下,确保错误处理器能及时重置生产者状态,避免长期阻塞。
- 监控消费者组状态与事务状态,提前发现离线或异常情况。
内容的提问来源于stack exchange,提问作者Michael C
相关产品推荐
相关产品推荐

