Spring Batch主从表读取与分层结构文件写入实现咨询
针对Spring Batch大订单量发票文件生成的Step设计方案
核心思路
由于单发票下订单量可达40万+,必须避免将整个发票的所有订单加载到内存,采用流式处理+按发票分组的策略,结合Spring Batch的Chunk处理和自定义ItemWriter实现结构化输出,从根源上避免内存溢出问题。
Step拆分与组件设计
1. 读取阶段:分页读取发票,避免全量加载
- 使用
JdbcPagingItemReader分页读取发票主表数据,每次仅读取1个发票(适配单发票大订单量场景),而非一次性拉取所有发票。 - 不在Reader中关联订单及详情数据,避免提前加载海量数据占用内存。
示例Reader配置(简化):
@Bean public JdbcPagingItemReader<Invoice> invoiceReader(DataSource dataSource) { return new JdbcPagingItemReaderBuilder<Invoice>() .dataSource(dataSource) .rowMapper(new BeanPropertyRowMapper<>(Invoice.class)) .selectClause("SELECT id, invoice_no, issue_date, total_amount FROM invoice") .fromClause("FROM invoice") .sortKeys(Map.of("id", Order.ASCENDING)) .pageSize(1) // 每次仅读取1个发票 .build(); }
2. 处理阶段:动态流式加载当前发票的订单及详情
针对每个发票,在写入阶段动态创建JdbcCursorItemReader,以当前发票ID为查询条件,逐行读取订单及关联的详情数据(可通过JOIN订单详情表实现,或按需单独查询详情),全程流式处理不积压数据。
3. 写入阶段:自定义ItemWriter实现结构化输出
自定义ItemWriter<Invoice>,在write方法内完成完整的文件结构输出逻辑:
- 文件头仅在首次写入时输出,用全局标志位控制;
- 依次输出当前发票的基础信息;
- 流式读取并逐行输出该发票下的所有订单及详情;
- 输出当前发票的页脚(包含订单数、总金额等统计值,流式读取时实时计算);
- 所有发票处理完成后,通过
@AfterJob注解触发文件尾的写入。
示例自定义Writer核心逻辑:
@StepScope public class InvoiceStructuredWriter implements ItemWriter<Invoice> { private final Resource outputResource; private boolean fileHeaderWritten = false; private int globalOrderCount = 0; private final DataSource dataSource; public InvoiceStructuredWriter(Resource outputResource, DataSource dataSource) { this.outputResource = outputResource; this.dataSource = dataSource; } @Override public void write(List<? extends Invoice> items) throws Exception { try (BufferedWriter writer = new BufferedWriter(new FileWriter(outputResource.getFile(), true))) { for (Invoice invoice : items) { // 写入文件头(仅第一次执行) if (!fileHeaderWritten) { writer.write("FILE_HEADER|20240520|INVOICE_BATCH_EXPORT"); writer.newLine(); fileHeaderWritten = true; } // 写入发票基础信息 writer.write(String.format("INVOICE|%s|%s|%s", invoice.getInvoiceNo(), invoice.getIssueDate(), invoice.getTotalAmount())); writer.newLine(); // 流式读取当前发票的订单及详情 JdbcCursorItemReader<OrderWithDetails> orderReader = buildOrderReader(invoice.getId()); OrderWithDetails order; int invoiceOrderCount = 0; BigDecimal invoiceActualTotal = BigDecimal.ZERO; while ((order = orderReader.read()) != null) { // 写入订单行 writer.write(String.format("ORDER|%s|%s", order.getOrderNo(), order.getOrderDate())); writer.newLine(); // 写入订单详情行(多详情则多行) for (OrderDetail detail : order.getDetails()) { writer.write(String.format("ORDER_DETAIL|%s|%d|%s", detail.getSku(), detail.getQuantity(), detail.getPrice())); writer.newLine(); } invoiceOrderCount++; invoiceActualTotal = invoiceActualTotal.add(order.getOrderAmount()); globalOrderCount++; } orderReader.close(); // 及时释放资源 // 写入单发票页脚 writer.write(String.format("INVOICE_FOOTER|%s|%d|%s", invoice.getInvoiceNo(), invoiceOrderCount, invoiceActualTotal)); writer.newLine(); } } } // 构建订单及详情的流式Reader private JdbcCursorItemReader<OrderWithDetails> buildOrderReader(Long invoiceId) { return new JdbcCursorItemReaderBuilder<OrderWithDetails>() .dataSource(dataSource) .rowMapper(new OrderWithDetailsRowMapper()) .sql("SELECT o.order_no, o.order_date, o.order_amount, od.sku, od.quantity, od.price " + "FROM orders o JOIN order_detail od ON o.id = od.order_id " + "WHERE o.invoice_id = ?") .preparedStatementSetter(ps -> ps.setLong(1, invoiceId)) .build(); } // Job执行完成后写入文件尾 @AfterJob public void writeFileFooter(JobExecution jobExecution) throws Exception { try (BufferedWriter writer = new BufferedWriter(new FileWriter(outputResource.getFile(), true))) { writer.write(String.format("FILE_FOOTER|%d|%s", globalOrderCount, LocalDate.now().format(DateTimeFormatter.ISO_DATE))); writer.newLine(); } } }
4. Step配置
Step采用Chunk处理模式,Chunk Size设为1(对应每次处理1个发票),控制事务范围在单个发票的处理周期内:
@Bean public Step generateInvoiceFileStep(JobRepository jobRepository, PlatformTransactionManager transactionManager, ItemReader<Invoice> invoiceReader, InvoiceStructuredWriter writer) { return new StepBuilder("generateInvoiceFileStep", jobRepository) .<Invoice, Invoice>chunk(1, transactionManager) .reader(invoiceReader) .writer(writer) .build(); }
关键优化点
- 流式数据处理:全程逐行读取订单及详情,不将海量数据加载到内存,适配40万+订单场景;
- 文件操作优化:采用追加模式写入文件,减少文件打开/关闭次数,提升性能;
- 事务控制:单个发票对应一个Chunk事务,避免大事务导致的数据库锁表或性能问题;
- 实时统计:流式读取订单时实时计算发票及全局统计数据,无需额外查询数据库。
注意事项
- 数据库索引优化:为
orders.invoice_id、order_detail.order_id字段添加索引,避免关联查询全表扫描; - 资源释放:动态创建的订单Reader必须及时关闭,防止数据库连接泄漏;
- 异常重试:可添加
SkipListener记录失败的发票ID,方便后续针对性重试。
内容的提问来源于stack exchange,提问作者Parth
相关产品推荐
相关产品推荐

