基于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
相关产品推荐
相关产品推荐

