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

Confluent Cloud中Kafka消费者遇异常后停止消费问题排查

问题分析: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的事务管理器在事务未正常结束的情况下,被要求再次开启新事务,导致状态转换冲突。具体原因如下:

  1. 手动ACK与事务机制冲突
    配置中已为容器指定KafkaTransactionManager,同时方法标注@Transactional,此时Spring Kafka会将消息消费和偏移量提交绑定到同一事务:事务提交时自动提交偏移量,事务回滚时偏移量也会回滚。但代码中在finally块手动调用ack.acknowledge(),会绕过事务机制直接提交偏移量,干扰事务生命周期,导致Producer事务管理器状态无法重置为IDLE,下一次处理消息时触发状态转换异常。

  2. AckMode配置与事务不兼容
    容器配置了AckMode.MANUAL_IMMEDIATE,但使用Kafka事务时,正确的AckMode应为AckMode.RECORD或AckMode.BATCH(由事务自动管理偏移量),手动ACK模式与事务机制无法协同,破坏事务原子性。

  3. 并发数设置不合理
    Topic仅3个分区,但容器并发数设为6,导致3个线程无法分配分区处于空闲状态,虽不是直接异常原因,但会增加线程资源浪费,加剧事务状态管理复杂度。

解决方案

  1. 移除手动ACK调用
    删除finally块中的ack.acknowledge()代码,由事务管理器自动处理偏移量的提交与回滚。

  2. 调整AckMode配置
    将容器AckMode改为ContainerProperties.AckMode.RECORD,或直接移除AckMode配置(Spring Kafka在使用事务时会默认采用兼容模式)。

  3. 修正并发数设置
    将容器并发数调整为不超过Topic分区数(即3),避免空闲线程,保证每个线程都能分配到分区。

  4. 验证事务管理器配置
    确保KafkaTransactionManager对应的ProducerFactory已正确配置transactionIdPrefix,这是开启Kafka事务的必要条件。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 09:55:24