Spring Batch:Slave步骤DB Writer提交后推送数据至队列方案咨询
最优解决方案:Spring Batch Slave步骤事务提交后发送队列数据
核心思路
要确保队列消息发送在Slave步骤的Chunk事务提交之后,必须依托事务同步机制实现"事务提交成功后再触发消息发送"的逻辑,同时要能直接获取当前Chunk的处理数据。
方案一:TransactionSynchronizationAdapter(推荐)
通过自定义ItemWriter包装类,在原始DB写入完成后注册事务同步器,仅当DB事务提交成功时才执行队列发送操作,完美满足时序要求。
实现步骤
- 编写事务同步型ItemWriter包装类:
public class TransactionalQueueItemWriter<T> implements ItemWriter<T> { private final ItemWriter<T> delegateWriter; private final QueueSender queueSender; public TransactionalQueueItemWriter(ItemWriter<T> delegateWriter, QueueSender queueSender) { this.delegateWriter = delegateWriter; this.queueSender = queueSender; } @Override public void write(List<? extends T> items) throws Exception { // 先执行原始DB写入逻辑 delegateWriter.write(items); // 注册事务同步器,仅在事务提交成功后执行队列发送 TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronizationAdapter() { @Override public void afterCommit() { // 提取数据所需属性,发送至嵌入式队列 List<PartialData> partialDataList = items.stream() .map(item -> new PartialData(item.getBizId(), item.getOperateType())) // 替换为实际需要的属性 .collect(Collectors.toList()); queueSender.send(partialDataList); } }); } }
- 修改Slave Step配置,替换原Writer为包装后的实例:
@Bean public Step slaveStep(ItemReader<BaseDTO> dispatchReader, CompositeItemProcessor<BaseDTO, BaseDTO> dispatchProcessor, CompositeItemWriter<BaseDTO> dispatchWriter, QueueSender queueSender) throws Exception { // 用事务同步Writer包装原始DB Writer ItemWriter<BaseDTO> transactionalWriter = new TransactionalQueueItemWriter<>(dispatchWriter, queueSender); return new StepBuilder("slaveStep", jobRepository) .<BaseDTO, BaseDTO>chunk(chunkSize, dispatchTransactionManager) .reader(dispatchReader) .processor(dispatchProcessor) .writer(transactionalWriter) .faultTolerant() .skipLimit(skipErrorCount) .skip(Exception.class) .noRetry(Exception.class) .noRollback(Exception.class) .processorNonTransactional() .listener(itemSkipListener()) .listener(stepExecutionListener()) .build(); }
方案优势
- 严格保证时序:仅当DB写入事务提交成功后才发送队列消息,完全避免"DB回滚但消息已发送"的数据不一致问题
- 数据获取直接:可直接拿到当前Chunk的所有处理数据,无需额外存储或查询
- 多线程安全:每个Slave线程的事务同步器独立运行,无并发冲突风险
方案二:StepExecutionListener.afterStep(补充方案)
若业务允许批量发送而非按Chunk实时发送,可在Step执行完成后,从StepExecution中获取成功处理的数据标识,批量查询后发送队列。
注意事项
- 需要提前在ItemWriter或自定义Listener中记录成功写入的数据ID
- 适合对实时性要求不高、允许批量处理的场景
关于ChunkListener的说明
ChunkListener的afterChunk方法确实无法直接获取Chunk的处理数据(仅能拿到ChunkContext),即便在多线程场景下无并发问题,也因数据获取限制无法满足你的需求,不推荐使用。
内容的提问来源于stack exchange,提问作者Rakesh
相关产品推荐
相关产品推荐

