如何用Spring Batch分批处理数据并实现基于连续两行的业务逻辑?
Hey there! Let's walk through how to solve your Spring Batch requirements—both processing a fixed number of rows per chunk and implementing business logic that relies on consecutive database rows. I've got a couple of solid approaches for you.
Spring Batch is built around the chunk-oriented processing model, which makes handling fixed row counts per batch straightforward. Here's how to set it up:
Step Configuration
Define your step with a specific chunk size (this is the number of rows you want to process per batch):
@Bean public Step dataProcessingStep(StepBuilderFactory stepBuilderFactory, ItemReader<MyDataEntity> dbItemReader, ItemProcessor<MyDataEntity, MyProcessedOutput> rowProcessor, ItemWriter<MyProcessedOutput> dbItemWriter) { return stepBuilderFactory.get("dataProcessingStep") .<MyDataEntity, MyProcessedOutput>chunk(15) // Replace 15 with your desired row count per batch .reader(dbItemReader) .processor(rowProcessor) .writer(dbItemWriter) .build(); }
Ordered Database Reader
Since you need consecutive rows, critical note: your item reader must fetch data in a consistent, ordered sequence (e.g., by primary key, timestamp, or another business-relevant field). For example, using a JdbcCursorItemReader:
@Bean public JdbcCursorItemReader<MyDataEntity> dbItemReader(DataSource dataSource) { return new JdbcCursorItemReaderBuilder<MyDataEntity>() .dataSource(dataSource) .sql("SELECT id, value, created_at FROM my_table ORDER BY created_at ASC") // Ensure ordering here! .rowMapper(new BeanPropertyRowMapper<>(MyDataEntity.class)) .build(); }
By default, Spring Batch's ItemProcessor handles one row at a time. To work with consecutive rows, you have two reliable options:
Option 1: Stateful ItemProcessor
Create a processor that retains the previous row in memory, allowing you to compare it with the current row. Make sure to mark it as @StepScope to avoid state leaks across step executions or threads.
Stateful Processor Implementation
@StepScope @Component public class ConsecutiveRowProcessor implements ItemProcessor<MyDataEntity, MyProcessedOutput>, ItemStream { private MyDataEntity previousRow; private static final String PREV_ROW_CONTEXT_KEY = "previousProcessedRow"; @Override public MyProcessedOutput process(MyDataEntity currentRow) throws Exception { MyProcessedOutput output = null; // Only process if we have a previous row to compare if (previousRow != null) { output = executeConsecutiveLogic(previousRow, currentRow); } // Update previous row for next iteration previousRow = currentRow; return output; } // Your custom business logic for consecutive rows private MyProcessedOutput executeConsecutiveLogic(MyDataEntity prev, MyDataEntity curr) { MyProcessedOutput output = new MyProcessedOutput(); output.setPrevId(prev.getId()); output.setCurrId(curr.getId()); output.setValueDifference(curr.getValue() - prev.getValue()); // Add more business logic here return output; } // Implement ItemStream to persist state for job restarts @Override public void open(ExecutionContext executionContext) { if (executionContext.containsKey(PREV_ROW_CONTEXT_KEY)) { previousRow = (MyDataEntity) executionContext.get(PREV_ROW_CONTEXT_KEY); } } @Override public void update(ExecutionContext executionContext) { executionContext.put(PREV_ROW_CONTEXT_KEY, previousRow); } @Override public void close() {} }
Handling the Final Row
This processor won't generate output for the last row (since there's no next row to compare). To handle this, add a StepExecutionListener to process the final row after the step completes:
@StepScope @Component public class FinalRowListener implements StepExecutionListener { @Autowired private ConsecutiveRowProcessor rowProcessor; @Override public void beforeStep(StepExecution stepExecution) {} @Override public ExitStatus afterStep(StepExecution stepExecution) { MyDataEntity finalRow = rowProcessor.getPreviousRow(); // Add a getter in the processor if (finalRow != null) { // Handle final row logic (e.g., write to a separate table, log, etc.) processFinalRow(finalRow); } return stepExecution.getExitStatus(); } private void processFinalRow(MyDataEntity finalRow) { // Your final row handling logic } }
Don't forget to register the listener in your step:
.step(dataProcessingStep.listener(finalRowListener()))
Option 2: Custom ItemReader for Row Pairs
If you prefer a more explicit approach, create a reader that fetches pairs of consecutive rows and passes them to your processor as a single object.
Define a Row Pair Container
public class DataRowPair { private MyDataEntity previousRow; private MyDataEntity currentRow; // Constructor, getters, setters public DataRowPair(MyDataEntity previousRow, MyDataEntity currentRow) { this.previousRow = previousRow; this.currentRow = currentRow; } // Getters public MyDataEntity getPreviousRow() { return previousRow; } public MyDataEntity getCurrentRow() { return currentRow; } }
Custom Pair Reader
@StepScope @Component public class PairItemReader implements ItemReader<DataRowPair> { private final ItemReader<MyDataEntity> delegateReader; private MyDataEntity cachedRow; public PairItemReader(ItemReader<MyDataEntity> delegateReader) { this.delegateReader = delegateReader; } @Override public DataRowPair read() throws Exception { MyDataEntity nextRow = delegateReader.read(); if (nextRow == null) { // Return final row as a pair with null if we have a cached row if (cachedRow != null) { DataRowPair finalPair = new DataRowPair(cachedRow, null); cachedRow = null; return finalPair; } return null; } if (cachedRow == null) { cachedRow = nextRow; return read(); // Recurse to get the next row for the pair } else { DataRowPair pair = new DataRowPair(cachedRow, nextRow); cachedRow = nextRow; return pair; } } }
Update Step for Pair Processing
Now your step will process pairs instead of individual rows:
@Bean public Step pairProcessingStep(StepBuilderFactory stepBuilderFactory, PairItemReader pairItemReader, ItemProcessor<DataRowPair, MyProcessedOutput> pairProcessor, ItemWriter<MyProcessedOutput> dbItemWriter) { return stepBuilderFactory.get("pairProcessingStep") .<DataRowPair, MyProcessedOutput>chunk(10) // 10 pairs = 20 rows (adjust as needed) .reader(pairItemReader) .processor(pairProcessor) .writer(dbItemWriter) .build(); }
Pair Processor
Your processor can now directly work with consecutive rows:
@Component public class PairProcessor implements ItemProcessor<DataRowPair, MyProcessedOutput> { @Override public MyProcessedOutput process(DataRowPair pair) throws Exception { if (pair.getCurrentRow() != null) { return executeConsecutiveLogic(pair.getPreviousRow(), pair.getCurrentRow()); } else { return processFinalRow(pair.getPreviousRow()); } } // Reuse your existing business logic methods here }
- Ordering is Non-Negotiable: Always ensure your database query includes an
ORDER BYclause to guarantee consecutive rows are logically related. - Step Scoping: Stateful components (like the stateful processor or pair reader) must be
@StepScopeto prevent cross-thread state contamination and support job restarts. - Restart Safety: Implement
ItemStreamfor stateful components to persist/restore state between job runs, ensuring no data is missed or processed twice.
内容的提问来源于stack exchange,提问作者bhalkian

