Spring Integration分单事务处理MongoDB数据时的重复循环问题排查
MongoDB集成流重复处理文档问题
原事务型集成流实现
最初实现了一个带事务的MongoDb inboundChannelAdapter集成流,通过更新表达式避免重复处理:
IntegrationFlow .from(MongoDb.inboundChannelAdapter(mongoTemplate, "{'status' : 'READY'}") .collectionName("work").entityClass(Document.class) .update(Update.update("status", "P")), p -> p.poller(pm -> pm.fixedDelay(1000L).transactional()))
该轮询器逻辑为:获取符合条件的文档列表、加写锁、拆分后处理;流程正常结束时,事务会在outboundGateway将文档状态设为DONE时提交。但此设计存在缺陷:单条文档的锁冲突会导致整轮轮询失败。
并行单文档事务优化方案
为解决整轮失败问题,尝试在split()后为单文档单独开启事务并并行处理:
@Bean public IntegrationFlow mongoFlowTxnPerDoc(MongoTemplate mongoTemplate, TransactionManager tm) { return IntegrationFlow .from(MongoDb.inboundChannelAdapter(mongoTemplate, "{'status' : 'TREADY'}") .collectionName("work").entityClass(Document.class), p -> p.poller(pm -> pm.fixedRate(1000L))) .split() .channel(c -> c.executor("txProcess.input", Executors.newFixedThreadPool(3))) .get(); }
同时在outboundGateway上配置e -> e.transactional(true),通过findOneAndUpdate实现单文档事务处理:
@Bean public IntegrationFlow txProcess(MongoTemplate mongoTemplate) { return f -> f .handle(MongoDb.outboundGateway(mongoTemplate).collectionName("work").entityClass(Document.class) .collectionCallback((c, m) -> c.findOneAndUpdate(Filters.eq("uuid", ((Document) m.getPayload()).get("uuid")), Updates.set("status", "DONE")) ), e -> e.transactional(true) ) .<Document>handle((p, h) -> { System.out.println("-----PROCESSING-----"); System.out.println(p); // 模拟耗时操作 try {Thread.sleep(10000);} catch (InterruptedException e) {} return null; }); }
异常行为
优化后出现以下异常:
- 插入1-2条数据时,行为符合预期,会触发
WriteConflict错误; - 插入3条及以上数据时,轮询阶段不再出现
WriteConflict错误,初始处理流程正常,但处理完成后会持续重复轮询并处理已被标记为DONE状态的相同文档,直到应用终止; - 即使移除
e -> e.transactional(true)配置,重复处理的问题依然存在。
原因猜测
推测问题源于线程池阻塞队列,导致重复轮询的数据未及时触发锁冲突。
内容的提问来源于stack exchange,提问作者Paul
相关产品推荐
相关产品推荐

