Confluent Cloud中Kafka消费者遇异常后停止消费问题排查
问题场景
AWS ECS中部署的Spring Kafka消费者出现以下异常,导致消费者停止消费:
org.springframework.transaction.CannotCreateTransactionException: Could not create Kafka transaction; nested exception is org.apache.kafka.common.KafkaException: TransactionalId item-availability-view-tx:41e90b10-0690-487b-8a26-9a4dcb295fbditem_availability_view_service_menu_item_changes.MENU_ITEMS_CHANGE_SYNC_FOR_86N.2: Invalid transition attempted from state IN_TRANSACTION to state IN_TRANSACTION at org.springframework.kafka.transaction.KafkaTransactionManager.doBegin(KafkaTransactionManager.java:162) ~[spring-kafka-2.5.0.RELEASE.jar:2.5.0.RELEASE] at org.springframework.transaction.support.AbstractPlatformTransactionManager.startTransaction(AbstractPlatformTransactionManager.java:400) ~[spring-tx-5.2.6.RELEASE.jar:5.2.6.RELEASE] at org.springframework.transaction.support.AbstractPlatformTransactionManager.getTransaction(AbstractPlatformTransactionManager.java:373) ~[spring-tx-5.2.6.RELEASE.jar:5.2.6.RELEASE] at org.springframework.transaction.support.TransactionTemplate.execute(TransactionTemplate.java:137) ~[spring-tx-5.2.6.RELEASE.jar:5.2.6.RELEASE] at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeRecordListenerInTx(KafkaMessageListenerContainer.java:1621) ~[spring-kafka-2.5.0.RELEASE.jar:2.5.0.RELEASE] at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeRecordListener(KafkaMessageListenerContainer.java:1596) ~[spring-kafka-2.5.0.RELEASE.jar:2.5.0.RELEASE] at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeListener(KafkaMessageListenerContainer.java:1330) ~[spring-kafka-2.5.0.RELEASE.jar:2.5.0.RELEASE] at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.pollAndInvoke(KafkaMessageListenerContainer.java:1062) ~[spring-kafka-2.5.0.RELEASE.jar:2.5.0.RELEASE] at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.run(KafkaMessageListenerContainer.java:970) ~[spring-kafka-2.5.0.RELEASE.jar:2.5.0.RELEASE] at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Unknown Source) ~[na:na] at java.base/java.util.concurrent.FutureTask.run(Unknown Source) ~[na:na] at java.base/java.lang.Thread.run(Unknown Source) ~[na:na] Caused by: org.apache.kafka.common.KafkaException: TransactionalId item-availability-view-tx:41e90b10-0690-487b-8a26-9a4dcb295fbditem_availability_view_service_menu_item_changes.MENU_ITEMS_CHANGE_SYNC_FOR_86N.2: Invalid transition attempted from state IN_TRANSACTION to state IN_TRANSACTION at org.apache.kafka.clients.producer.internals.TransactionManager.transitionTo(TransactionManager.java:1053) ~[kafka-clients-2.5.0.jar:na] at org.apache.kafka.clients.producer.internals.TransactionManager.transitionTo(TransactionManager.java:1046) ~[kafka-clients-2.5.0.jar:na] at org.apache.kafka.clients.producer.internals.TransactionManager.beginTransaction(TransactionManager.java:338) ~[kafka-clients-2.5.0.jar:na] at org.apache.kafka.clients.producer.KafkaProducer.beginTransaction(KafkaProducer.java:612) ~[kafka-clients-2.5.0.jar:na] at org.springframework.kafka.core.DefaultKafkaProducerFactory$CloseSafeProducer.beginTransaction(DefaultKafkaProducerFactory.java:795) ~[spring-kafka-2.5.0.RELEASE.jar:2.5.0.RELEASE] at brave.kafka.clients.TracingProducer.beginTransaction(TracingProducer.java:69) ~[brave-instrumentation-kafka-clients-5.12.7.jar:na] at org.springframework.kafka.core.ProducerFactoryUtils.getTransactionalResourceHolder(ProducerFactoryUtils.java:99) ~[spring-kafka-2.5.0.RELEASE.jar:2.5.0.RELEASE] at org.springframework.kafka.transaction.KafkaTransactionManager.doBegin(KafkaTransactionManager.java:146) ~[spring-kafka-2.5.0.RELEASE.jar:2.5.0.RELEASE]
相关配置与代码如下:
消费者定义
public interface UpdateSyncConsumer { @KafkaListener( id = "update_sync_consumer", topics = "${app.kafka.topics.UPDATES_TOPIC}", containerFactory = "updatesListenerFactory") @Transactional void handleUpdate( ConsumerRecord<String, UpdateAvro> consumerRecord, Acknowledgment ack); }
消费者实现
@Component @Log4j2 @RequiredArgsConstructor public class DefaultUpdateSyncConsumer implements UpdateSyncConsumer { private final UpdateAvroMapper avroMapper; private final UpdateJobScheduler updateJobScheduler; @Override public void handleUpdate( final ConsumerRecord<String, UpdateAvro> consumerRecord, final Acknowledgment ack) { try { log.info( "Scheduling Update Sync: payload: {} count:{}", consumerRecord.value(), consumerRecord.value().getItemIds().size()); updateJobScheduler.scheduleUpdateSyncAttempt(avroMapper.map(consumerRecord.value())); log.info("Scheduled Update Sync: DONE"); } catch (Exception e) { log.error("ERROR: Scheduling Update Sync: {}", Throwables.getStackTraceAsString(e)); } finally { ack.acknowledge(); log.info("Acknowledged: {}", consumerRecord.key()); } } }
消费者配置
@Log4j2 @EnableKafka @Configuration @RequiredArgsConstructor public class KafkaConsumerConfig extends KafkaBasicConfig { private final KafkaTemplate<String, Object> kafkaTemplate; @Bean public ConsumerFactory<String, Object> consumerFactory() { Map<String, Object> basicConfig = getBasicConfig(); basicConfig.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest"); basicConfig.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); basicConfig.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed"); basicConfig .put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class); basicConfig.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class); basicConfig .put(ErrorHandlingDeserializer.KEY_DESERIALIZER_CLASS, StringDeserializer.class.getName()); basicConfig.put(ErrorHandlingDeserializer.VALUE_DESERIALIZER_CLASS, KafkaAvroDeserializer.class.getName()); basicConfig.put(KafkaAvroDeserializerConfig.SPECIFIC_AVRO_READER_CONFIG, true); return new DefaultKafkaConsumerFactory<>(basicConfig); } @Bean(name = "updatesListenerFactory") public ConcurrentKafkaListenerContainerFactory<String, Object> kafkaListenerContainerFactory( KafkaTransactionManager<String, Object> ktm) { ConcurrentKafkaListenerContainerFactory<String, Object> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConcurrency(6); factory.setConsumerFactory(consumerFactory()); factory.getContainerProperties().setTransactionManager(ktm); factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE); factory.setAfterRollbackProcessor(afterRollbackProcessor()); return factory; } private DefaultAfterRollbackProcessor<String, Object> afterRollbackProcessor() { DefaultAfterRollbackProcessor<String, Object> processor = new DefaultAfterRollbackProcessor<>((rec, exception) -> { log.error( "!!!!!!!!!!! Following Record keeps failing despite multiple attempts, therefore would be skipped !! : " + rec); log.error(exception); }, new FixedBackOff(1000L, 2L)); processor.setCommitRecovered(true); processor.setKafkaTemplate(kafkaTemplate); processor.setLogLevel(KafkaException.Level.ERROR); processor.addNotRetryableException(SerializationException.class); return processor; } }
异常原因分析
核心异常是Invalid transition attempted from state IN_TRANSACTION to state IN_TRANSACTION,说明Kafka Producer的事务管理器在事务未正常结束的情况下,被要求再次开启新事务,导致状态转换冲突。具体原因如下:
手动ACK与事务机制冲突
配置中已为容器指定KafkaTransactionManager,同时方法标注@Transactional,此时Spring Kafka会将消息消费和偏移量提交绑定到同一事务:事务提交时自动提交偏移量,事务回滚时偏移量也会回滚。但代码中在finally块手动调用ack.acknowledge(),会绕过事务机制直接提交偏移量,干扰事务生命周期,导致Producer事务管理器状态无法重置为IDLE,下一次处理消息时触发状态转换异常。AckMode配置与事务不兼容
容器配置了AckMode.MANUAL_IMMEDIATE,但使用Kafka事务时,正确的AckMode应为AckMode.RECORD或AckMode.BATCH(由事务自动管理偏移量),手动ACK模式与事务机制无法协同,破坏事务原子性。并发数设置不合理
Topic仅3个分区,但容器并发数设为6,导致3个线程无法分配分区处于空闲状态,虽不是直接异常原因,但会增加线程资源浪费,加剧事务状态管理复杂度。
解决方案
移除手动ACK调用
删除finally块中的ack.acknowledge()代码,由事务管理器自动处理偏移量的提交与回滚。调整AckMode配置
将容器AckMode改为ContainerProperties.AckMode.RECORD,或直接移除AckMode配置(Spring Kafka在使用事务时会默认采用兼容模式)。修正并发数设置
将容器并发数调整为不超过Topic分区数(即3),避免空闲线程,保证每个线程都能分配到分区。验证事务管理器配置
确保KafkaTransactionManager对应的ProducerFactory已正确配置transactionIdPrefix,这是开启Kafka事务的必要条件。
内容的提问来源于stack exchange,提问作者Romeo Sierra

