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

基于Spring Batch实现MongoDB订单数据拆分导出至100MB以内CSV

Spring Batch 解决方案:MongoDB采购记录CSV导出(多线程+同订单同文件+大小控制)

核心思路

针对需求,核心解决点在于均衡负载的分区策略、同订单记录的连续读取与写入、动态文件大小控制,利用Spring Batch的分区机制+自定义Writer实现。

1. 优化分区策略解决负载不均

放弃固定数量分区,改用基于订单数据量的动态分区:

  • 先通过MongoDB聚合统计每个order_id对应的采购项数量,计算总记录数后,将order_id划分成N个组,每个组的累计记录数大致相等(比如目标每个分区处理10万条左右)
  • 每个分区负责处理一组order_id的所有采购记录,确保每个线程的处理量均衡,同时保证同订单的记录在同一个分区内

2. 保证同订单记录连续读取

配置MongoItemReader时,按order_id排序读取,确保同订单的所有采购项连续返回,为后续写入同一文件提供基础:

  • 给MongoDB的order_id字段添加索引,提升排序和查询性能

3. 自定义CSV Writer实现文件大小控制

实现自定义ItemWriter,逻辑如下:

  • 维护当前写入的文件句柄、当前文件大小、当前处理的order_id
  • 每次写入前,若当前记录的order_id与上一个不同:
    • 检查当前文件大小是否接近100MB(比如设置95MB阈值,预留空间)
    • 若达到阈值,关闭当前文件,创建新文件
  • 将当前订单的所有采购项写入同一文件,直到订单切换

代码实现示例

第一步:MongoDB索引配置

在MongoDB中创建索引,提升查询性能:

db.purchase_records.createIndex({order_id: 1})

第二步:Spring Batch 核心配置

1. 分区器(Partitioner)实现

@Component
public class PurchaseOrderPartitioner implements Partitioner {

    @Autowired
    private MongoTemplate mongoTemplate;

    private static final long TARGET_RECORDS_PER_PARTITION = 100000; // 每个分区目标处理10万条

    @Override
    public Map<String, ExecutionContext> partition(int gridSize) {
        // 统计每个order_id的记录数
        Aggregation aggregation = Aggregation.newAggregation(
                Aggregation.group("order_id").count().as("count"),
                Aggregation.sort(Sort.Direction.ASC, "_id")
        );
        List<OrderCount> orderCounts = mongoTemplate.aggregate(aggregation, "purchase_records", OrderCount.class).getMappedResults();

        Map<String, ExecutionContext> partitions = new HashMap<>();
        int partitionIndex = 0;
        long currentPartitionTotal = 0;
        List<String> currentOrderIds = new ArrayList<>();

        for (OrderCount orderCount : orderCounts) {
            currentOrderIds.add(orderCount.getOrderId());
            currentPartitionTotal += orderCount.getCount();

            // 当当前分区累计记录数接近目标值,或者是最后一个订单时,创建分区
            if (currentPartitionTotal >= TARGET_RECORDS_PER_PARTITION || orderCount == orderCounts.get(orderCounts.size() - 1)) {
                ExecutionContext context = new ExecutionContext();
                context.put("orderIds", currentOrderIds);
                partitions.put("partition" + partitionIndex, context);
                partitionIndex++;
                currentPartitionTotal = 0;
                currentOrderIds = new ArrayList<>();
            }
        }
        return partitions;
    }

    @Data
    private static class OrderCount {
        @Field("_id")
        private String orderId;
        private long count;
    }
}

2. 读取器(ItemReader)配置

@Configuration
public class PurchaseReaderConfig {

    @Autowired
    private MongoTemplate mongoTemplate;

    @Bean
    @StepScope
    public MongoItemReader<PurchaseRecord> purchaseItemReader(@Value("#{stepExecutionContext['orderIds']}") List<String> orderIds) {
        Query query = Query.query(Criteria.where("order_id").in(orderIds)).with(Sort.by("order_id"));

        return new MongoItemReaderBuilder<PurchaseRecord>()
                .template(mongoTemplate)
                .collection("purchase_records")
                .query(query)
                .targetType(PurchaseRecord.class)
                .build();
    }
}

3. 自定义CSV Writer实现

@Component
@StepScope
public class PurchaseCsvWriter implements ItemWriter<PurchaseRecord> {

    private static final long MAX_FILE_SIZE = 100 * 1024 * 1024; // 100MB
    private static final long THRESHOLD_SIZE = 95 * 1024 * 1024; // 95MB阈值,提前切换文件

    private BufferedWriter currentWriter;
    private long currentFileSize;
    private String currentOrderId;
    private int fileIndex = 1;

    @Value("${csv.output.path}")
    private String outputPath;

    @Override
    public void write(List<? extends PurchaseRecord> items) throws Exception {
        for (PurchaseRecord record : items) {
            String line = convertToCsvLine(record);
            long lineSize = line.getBytes(StandardCharsets.UTF_8).length;

            // 处理订单切换或文件大小达到阈值的情况
            if (!record.getOrderId().equals(currentOrderId) || (currentFileSize + lineSize) > THRESHOLD_SIZE) {
                // 如果不是第一个文件,先关闭当前writer
                if (currentWriter != null) {
                    currentWriter.flush();
                    currentWriter.close();
                }
                // 创建新文件
                String fileName = String.format("%s/purchase_records_%d.csv", outputPath, fileIndex++);
                currentWriter = new BufferedWriter(new FileWriter(fileName));
                currentFileSize = 0;
                // 写入表头(仅新文件)
                currentWriter.write("order_id,product_id,quantity,price,purchase_date");
                currentWriter.newLine();
                currentFileSize += "order_id,product_id,quantity,price,purchase_date\n".getBytes(StandardCharsets.UTF_8).length;
            }

            // 写入记录
            currentWriter.write(line);
            currentWriter.newLine();
            currentFileSize += lineSize + System.lineSeparator().getBytes(StandardCharsets.UTF_8).length;
            currentOrderId = record.getOrderId();
        }
    }

    // 自定义CSV行转换逻辑
    private String convertToCsvLine(PurchaseRecord record) {
        return String.join(",",
                record.getOrderId(),
                record.getProductId(),
                String.valueOf(record.getQuantity()),
                String.valueOf(record.getPrice()),
                record.getPurchaseDate().toString()
        );
    }

    // 销毁时关闭文件
    @PreDestroy
    public void close() throws IOException {
        if (currentWriter != null) {
            currentWriter.flush();
            currentWriter.close();
        }
    }
}

4. Job与Step配置

@Configuration
@EnableBatchProcessing
public class PurchaseExportBatchConfig {

    @Autowired
    private JobBuilderFactory jobBuilderFactory;
    @Autowired
    private StepBuilderFactory stepBuilderFactory;
    @Autowired
    private PurchaseOrderPartitioner partitioner;
    @Autowired
    private TaskExecutor taskExecutor;

    @Bean
    public Step slaveStep(MongoItemReader<PurchaseRecord> purchaseItemReader,
                          PurchaseCsvWriter purchaseCsvWriter) {
        return stepBuilderFactory.get("slaveStep")
                .<PurchaseRecord, PurchaseRecord>chunk(1000) // 按块处理,提升性能
                .reader(purchaseItemReader)
                .writer(purchaseCsvWriter)
                .build();
    }

    @Bean
    public Step masterStep(Step slaveStep) {
        return stepBuilderFactory.get("masterStep")
                .partitioner(slaveStep.getName(), partitioner)
                .step(slaveStep)
                .partitionHandler(new TaskExecutorPartitionHandler() {{
                    setTaskExecutor(taskExecutor);
                    setStep(slaveStep);
                    setGridSize(10); // 可根据CPU核心数调整
                }})
                .build();
    }

    @Bean
    public Job purchaseExportJob(Step masterStep) {
        return jobBuilderFactory.get("purchaseExportJob")
                .start(masterStep)
                .build();
    }

    @Bean
    public TaskExecutor taskExecutor() {
        ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
        executor.setCorePoolSize(8);
        executor.setMaxPoolSize(12);
        executor.setQueueCapacity(20);
        executor.setThreadNamePrefix("purchase-export-");
        executor.initialize();
        return executor;
    }
}

关键注意事项

  • 索引优化:必须给order_id加索引,否则排序和查询会非常慢,尤其是百万级数据
  • 块大小调整:chunk(1000)可根据内存情况调整,太大可能导致内存溢出,太小会增加IO次数
  • 文件大小阈值:设置95MB而不是100MB,是为了避免因为最后一个订单的多条记录导致文件超过100MB
  • 异常处理:可在Writer中添加异常捕获,确保文件正常关闭;也可在Batch中配置重试、跳过策略

内容的提问来源于stack exchange,提问作者devincrisis

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 17:35:22