@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
相关产品推荐
相关产品推荐

