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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 13:54:05