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

@KafkaListener中MongoDB事务异常:事务已关闭却重复处理消息

问题:Kafka监听器中MongoDB事务批量插入后事务中止导致重复消费

在@KafkaListener方法内使用MongoDB事务执行无序批量插入(允许部分失败,但需追踪成功插入的数据)时,遇到以下问题:

  • 日志显示事务已关闭,但Kafka监听器持续重处理该批次(抛出ListenerExecutionFailedException)
  • 怀疑是MongoDB驱动在捕获数据完整性违规时自动中止事务,或是MongoDB事务与Kafka事务(未启用)存在冲突

想确认是否有其他开发者遇到过类似问题,以及是否在事务/Kafka监听器配置上存在遗漏。

错误日志

org.springframework.kafka.listener.ListenerExecutionFailedException: Listener method 'public void com.example.MyListener.listen(java.util.List<java.lang.String>)' threw exception
...
Caused by: org.springframework.transaction.TransactionSystemException: Could not commit Mongo transaction for session [ClientSessionImpl@660ef6de id = {"id": {"$binary": {"base64": "lRRKccYZTAasvsU3oVUCCg==", "subType": "04"}}}, causallyConsistent = true, txActive = false, txNumber = 11, closed = false, clusterTime = {"clusterTime": {"$timestamp": {"t": 1743775965, "i": 11}}, "signature": {"hash": {"$binary": {"base64": "AAAAAAAAAAAAAAAAAAAAAAAAAAA==", "subType": "00"}}, "keyId": 0}}].
...
Caused by: com.mongodb.MongoCommandException: Command failed with error 251 (NoSuchTransaction): 'Transaction with { txnNumber: 11 } has been aborted.' on server localhost:27017. The full response is {"errorLabels": ["TransientTransactionError"], "ok": 0.0, "errmsg": "Transaction with { txnNumber: 11 } has been aborted.", "code": 251, "codeName": "NoSuchTransaction", "$clusterTime": {"clusterTime": {"$timestamp": {"t": 1743775965, "i": 11}}, "signature": {"hash": {"$binary": {"base64": "AAAAAAAAAAAAAAAAAAAAAAAAAAA==", "subType": "00"}}, "keyId": 0}}, "operationTime": {"$timestamp": {"t": 1743775965, "i": 11}}}
    at com.mongodb.internal.connection.ProtocolHelper.getCommandFailureException(ProtocolHelper.java:210) ~[mongodb-driver-core-5.2.1.jar:na]
    ...

相关代码

MyRepository.java

@Component
@AllArgsConstructor
class MyRepository {

  private final MongoOperations mongoOperations;

  public void transactionalBulkSave(List<MyDocument> documents) {
    BulkOperations bulkOperations = mongoOperations.bulkOps(BulkMode.UNORDERED, MyDocument.class);
    bulkOperations.insert(documents);
    bulkOperations.execute();
  }
}

MyService.java

@Service
@AllArgsConstructor
@Slf4j
class MyService {

  private final MyRepository myRepository;

  @Transactional(transactionManager = "mongoTransactionManager")
  public void transactionalBulkSave(List<MyDocument> documents) {
    try {
      myRepository.transactionalBulkSave(documents);
    } catch (BulkOperationException e) {
      log.error(e.getMessage());
      Set<String> collect = e
          .getResult()
          .getInserts()
          .stream()
          .map(el -> el
              .getId()
              .asObjectId()
              .getValue()
              .toString())
          .collect(Collectors.toSet());
      List<MyDocument> valid = documents
          .stream()
          .filter(doc -> collect.contains(doc
              .id()
              .toString()))
          .toList();
      log.info("saved {}", valid);
    }
  }
}

MyListener.java

@Component
@AllArgsConstructor
@Slf4j
class MyListener {

  private final MyService myService;

  @KafkaListener(
      topics = "test",
      batch = "true"
  )
  public void listen(List<String> strings) {
    log.info("Received: {}", strings);
    List<MyDocument> list = strings
        .stream()
        .map(e -> new MyDocument(ObjectId.get(), e))
        .toList();
    myService.transactionalBulkSave(list);
  }
}

MyConfiguration.java

@Configuration(proxyBeanMethods = false)
class MyConfiguration {

  @Bean
  MongoTransactionManager mongoTransactionManager(MongoDatabaseFactory databaseFactory) {
    return new MongoTransactionManager(databaseFactory);
  }
}

测试代码

@Test
void shouldInsertOthersButNotDuplicates() {
  kafkaTemplate.send("test", UUID
      .randomUUID()
      .toString(), "1");

  await()
      .atMost(1000, TimeUnit.SECONDS)
      .untilAsserted(() -> assertThat(mongoOperations.findAll(MyDocument.class)).hasSize(1));

  kafkaTemplate.send("test", UUID
      .randomUUID()
      .toString(), "1");
  kafkaTemplate.send("test", UUID
      .randomUUID()
      .toString(), "2");

  await()
      .atMost(10, TimeUnit.SECONDS)
      .untilAsserted(() -> assertThat(mongoOperations.findAll(MyDocument.class)).hasSize(2));
}

MyDocument.java

@Document
record MyDocument(
    @Id
    ObjectId id,
    @Indexed(unique = true)
    String value
) {
}

问题根源与解决方案

1. 核心问题:事务已中止但Spring仍尝试提交

当MongoDB执行无序批量插入时,若出现唯一键冲突(数据完整性违规),MongoDB驱动会自动中止当前事务——事务内的批量操作失败会触发事务中止规则。你的代码捕获BulkOperationException后未处理事务状态,Spring事务管理器在方法结束时仍尝试提交已中止的事务,导致NoSuchTransaction错误,最终Kafka监听器因异常触发无限重试。

另外需注意:MongoDB事务内的UNORDERED批量操作不支持部分提交,只要有一个操作失败,整个事务会被标记为中止,所有已执行的操作都会回滚。

2. 解决方案

方案一:禁用事务,手动追踪批量结果

既然需求是允许部分失败并追踪成功数据,事务的原子性反而会成为阻碍。直接移除@Transactional注解,让无序批量操作的部分成功结果直接保留:

修改MyService方法:

public void transactionalBulkSave(List<MyDocument> documents) {
  try {
    myRepository.transactionalBulkSave(documents);
  } catch (BulkOperationException e) {
    log.error(e.getMessage());
    Set<String> successIds = e.getResult().getInserts()
        .stream()
        .map(el -> el.getId().asObjectId().getValue().toString())
        .collect(Collectors.toSet());
    List<MyDocument> savedDocs = documents.stream()
        .filter(doc -> successIds.contains(doc.id().toString()))
        .toList();
    log.info("已成功插入: {}", savedDocs);
  }
}

方案二:保留事务,拆分批量为单个插入

若必须使用事务(比如有其他原子操作需求),需将批量操作拆分为单个插入,逐个处理重复异常:

@Transactional(transactionManager = "mongoTransactionManager")
public void transactionalBulkSave(List<MyDocument> documents) {
  List<MyDocument> success = new ArrayList<>();
  for (MyDocument doc : documents) {
    try {
      mongoOperations.insert(doc);
      success.add(doc);
    } catch (DuplicateKeyException e) {
      log.warn("文档 {} 已存在,跳过", doc.id());
    }
  }
  log.info("已成功插入: {}", success);
}

这种方式下,单个插入的异常不会导致整个事务中止,仅跳过重复数据,事务可正常提交,避免Kafka重试。

方案三:配置Kafka监听器重试策略

若必须保留当前事务+批量操作模式,需调整Kafka监听器的重试策略,避免无限重试:

  • 配置retryTemplate设置重试次数上限
  • 排除NoSuchTransaction这类不可重试的异常

示例配置:

@KafkaListener(
    topics = "test",
    batch = "true",
    retryTemplate = "@myRetryTemplate"
)

对应的RetryTemplate Bean:

@Bean
public RetryTemplate myRetryTemplate() {
  RetryTemplate template = new RetryTemplate();
  SimpleRetryPolicy policy = new SimpleRetryPolicy();
  policy.setMaxAttempts(3);
  policy.setRetryableExceptions(Map.of(
      BulkOperationException.class, true,
      NoSuchTransactionException.class, false
  ));
  template.setRetryPolicy(policy);
  return template;
}

3. 额外注意事项

  • MongoDB事务仅支持副本集或分片集群,单实例MongoDB无法使用事务,此时事务配置会直接导致异常
  • 避免混用MongoDB事务与Kafka事务,确保Kafka的enable.auto.commit为true(默认),且未配置kafkaTransactionManager

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 13:27:09