Spring Batch中如何实现多Writer并行写入多个MongoDB集合
实现Spring Batch多Mongo集合并行写入的方案
你可以选择两种方案实现并行写入,根据你的业务场景选择即可:
方案1:改造CompositeWriter内部并行写不同集合(推荐,改动最小)
这个方案是在单Chunk的写入阶段,将不同集合的批量写入操作并行执行,改动最小,也符合你现有的代码逻辑。
步骤1:配置异步线程池Bean
首先在配置类中添加线程池配置:
@Bean public TaskExecutor writerTaskExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); // 核心线程数和你要写入的集合数量匹配即可 executor.setCorePoolSize(4); executor.setMaxPoolSize(8); executor.setQueueCapacity(100); executor.setThreadNamePrefix("mongo-writer-"); // 拒绝策略直接抛异常,保证Chunk失败回滚 executor.setRejectedExecutionHandler(new ThreadPoolExecutor.AbortPolicy()); executor.initialize(); return executor; }
步骤2:优化并改造CompositeMongoItemWriter的write方法
先按集合对要写入的数据分组,再并行执行不同集合的批量写入,避免原来单条写入的性能损耗:
public class CompositeMongoItemWriter extends MongoItemWriter<CompositeWriterData> { @Autowired MongoItemWriter<ProfileRecommendationInfo> recommendationsDataWriter; @Autowired MongoItemWriter<ProfileLifebandInfo> lifeBandWriter; @Autowired private MongoTemplate secondaryMongoTemplate; @Autowired MongoItemWriter<Profile> profileWriter; // 注入刚才配置的线程池 @Autowired private TaskExecutor writerTaskExecutor; @Override public void write(List<? extends CompositeWriterData> items) throws Exception { if( items== null || items.isEmpty()) { return; } // 先按集合分组,攒批量数据 List<ProfileRecommendationInfo> recList = new ArrayList<>(); List<ProfileLifebandInfo> lifebandList = new ArrayList<>(); // 存储要更新的Profile映射,key是profileId,value是对应的recommendationId Map<String, String> profileUpdateMap = new HashMap<>(); for(CompositeWriterData compositeWriterData : items) { for( Map.Entry<String, Object> collection : compositeWriterData.getCollectionsPOJODataMap().entrySet() ) { String collectionName = collection.getKey(); if(CommonConstants.PROFILE_RECOMMENDATION_INFO.equalsIgnoreCase(collectionName)) { ProfileRecommendationInfo rec = (ProfileRecommendationInfo) collection.getValue(); recList.add(rec); profileUpdateMap.put(rec.getProfileId(), rec.getDataId()); }else if(CommonConstants.PROFILE_LIFEBAND_INFO.equalsIgnoreCase(collectionName)) { lifebandList.add((ProfileLifebandInfo) collection.getValue()); } } } // 提交并行写入任务 List<Callable<Void>> writeTasks = new ArrayList<>(); if(!recList.isEmpty()) { writeTasks.add(() -> { recommendationsDataWriter.write(recList); return null; }); } if(!lifebandList.isEmpty()) { writeTasks.add(() -> { lifeBandWriter.write(lifebandList); return null; }); } // 等待所有写入任务完成,任意任务失败都会抛出异常触发Chunk回滚 CompletableFuture.allOf( writeTasks.stream() .map(task -> CompletableFuture.runAsync(() -> { try { task.call(); } catch (Exception e) { throw new RuntimeException("写入集合失败", e); } }, writerTaskExecutor)) .toArray(CompletableFuture[]::new) ).join(); // 所有集合写入完成后,再批量更新Profile if(!profileUpdateMap.isEmpty()) { List<Profile> profileList = new ArrayList<>(); for (Map.Entry<String, String> entry : profileUpdateMap.entrySet()) { Profile profile = secondaryMongoTemplate.findById(entry.getKey(), Profile.class); profile.setProfileRecommendationInfoId(entry.getValue()); profileList.add(profile); } profileWriter.write(profileList); } } // 其余原有方法保持不变即可 }
方案2:启用Spring Batch多线程Step(更高并行度,适合数据量大的场景)
如果需要更高的并行度,可以直接将Step改造为多线程执行,多个Chunk会并行被读取、处理、写入,也是官方推荐的并行处理方案。
只需要修改你的Step配置,注入线程池即可:
@Bean public Step singleETLStep(MongoItemWriter<CompositeWriterData> writer, HttpServletRequest request, TaskExecutor writerTaskExecutor) { return stepBuilderFactory.get("singleETLStep") .<UserInfo, CompositeWriterData>chunk(etlConfiguration.getBatchChunkSize()) .reader(reader(("#{jobParameters[profileId]}"))).faultTolerant().skipPolicy(readerSkipper()) .processor(processor(request,"#{jobParameters[executeProcessing]}")) .listener(processorListener()).faultTolerant().skipPolicy(writerSkipper()) .writer(writer) // 新增多线程配置 .taskExecutor(writerTaskExecutor) // 并发线程数可根据实际资源调整 .throttleLimit(10) .build(); }
注意:MongoItemReader本身是线程安全的,无需额外改造,如果你的自定义Processor是无状态的,就可以直接使用该方案,配合方案1的并行写入可以达到最高性能。
注意事项
- 要保证你的MongoDB服务端能支撑对应的并发写入压力,调整线程池参数时要结合Mongo的性能情况
- 如果写入操作之间有强顺序依赖,不要强行并行,比如你现有的Profile更新逻辑必须等ProfileRecommendationInfo写入完成后再执行,不能并行
内容的提问来源于stack exchange,提问作者Sanjay Naik
相关产品推荐
相关产品推荐

