
本文详解如何在 Spring Batch 中实现自定义 Reader(读取 ID 列表)、Processor(按 ID 批量查出并处理对象列表)、Writer(写入对象列表),并重点纠正因误用 @StepScope 导致的重复写入问题。
本文详解如何在 spring batch 中实现自定义 reader(读取 id 列表)、processor(按 id 批量查出并处理对象列表)、writer(写入对象列表),并重点纠正因误用 `@stepscope` 导致的重复写入问题。
在 Spring Batch 中构建高度定制化的数据流时,需严格遵循其生命周期与作用域契约。您当前的实现虽逻辑清晰,但存在一个关键性设计偏差:错误地将 @StepScope 应用于 CustomItemReader,这直接导致了状态丢失与重复写入问题。
? 问题根源分析
@StepScope 表示 Bean 在每次 Step 启动时全新实例化。您的 CustomItemReader 中依赖 bookingIds 这一本地缓存字段来维持“已读 ID 清单”的状态:
private List<string> bookingIds; // ← 状态变量
@Override
public String read() {
if (Objects.isNull(bookingIds)) {
bookingIds = bookingInfoRepository.findDistinctId(); // ← 每次 step 启动都重新查!
}
return bookingIds.isEmpty() ? null : bookingIds.remove(0);
}</string>
由于 @StepScope 的语义,每次 Step 执行(哪怕同一 Job 的多次重试或分片)都会创建新 Reader 实例 → bookingIds 始终为 null → 每次都重新查询全部 ID → 同一批 ID 被反复读取、处理、写入 → 出现数据重复插入。
此外,@StepScope 对 Processor 和 Writer 并无必要:它们是无状态的纯函数组件(接收输入、返回/消费输出),无需绑定到 Step 生命周期。
✅ 正确的作用域配置方案
| 组件 | 推荐作用域 | 原因 |
|---|---|---|
| CustomItemReader | @JobScope | 需在整个 Job 生命周期内共享状态(如 bookingIds 列表),避免重复加载;JobScope 确保同一 Job 实例中复用同一个 Reader 实例。 |
| CorrectionProcessor | 默认 Singleton(无需注解) | 无内部状态,仅依赖注入的 Repository,线程安全且高效。 |
| CustomItemWriter | 默认 Singleton(无需注解) | 同上;写入逻辑不依赖 Step 上下文状态。 |
⚠️ 注意:启用 @JobScope 前,必须在任意 @Configuration 类上添加 @EnableBatchProcessing(Spring Boot 2.5+ 默认启用,但显式声明更稳妥)。
?️ 修正后的代码示例
@Component
@JobScope // ✅ 改为 @JobScope
@StepScope // ❌ 移除
@Slf4j
public class CustomItemReader implements ItemReader<string> {
@Autowired
private BookingInfoRepository bookingInfoRepository;
private List<string> bookingIds; // 状态由 JobScope 保障复用
@Override
public String read() {
if (bookingIds == null) {
bookingIds = bookingInfoRepository.findDistinctId();
log.info("Loaded {} distinct booking IDs", bookingIds.size());
}
return bookingIds.isEmpty() ? null : bookingIds.remove(0);
}
}</string></string>
@Component // ✅ 移除 @StepScope,使用默认 singleton
@Slf4j
public class CorrectionProcessor implements ItemProcessor<string list>> {
@Autowired
private BookingInfoRepository bookingInfoRepository;
@Override
public List<bookinginfo> process(String bookingId) {
List<bookinginfo> list = bookingInfoRepository.findById(bookingId);
// ✅ 示例:对每个 BookingInfo 执行业务修正
list.forEach(info -> {
info.setStatus("CORRECTED");
info.setLastModified(Instant.now());
});
return list;
}
}</bookinginfo></bookinginfo></string>
@Component // ✅ 移除 @StepScope
public class CustomItemWriter implements ItemWriter<list>> {
@Autowired
private BookingInfoRepository bookingInfoRepository;
@Override
public void write(Chunk extends List<bookinginfo>> chunk) throws Exception {
// ✅ 扁平化嵌套列表:Chunk<list>> → List<bookinginfo>
List<bookinginfo> allItems = chunk.stream()
.flatMap(List::stream)
.collect(Collectors.toList());
if (!allItems.isEmpty()) {
bookingInfoRepository.saveAll(allItems); // 使用 JPA saveAll 批量写入
log.debug("Written {} BookingInfo objects in current chunk", allItems.size());
}
}
}</bookinginfo></bookinginfo></list></bookinginfo></list>
? 关键注意事项与最佳实践
-
Reader 必须是线程安全的:@JobScope 下 Reader 是单例,若 Job 并行执行(如 Split 或远程分片),需确保 read() 方法无竞态条件。当前实现中 bookingIds.remove(0) 是线程不安全的 —— 若需并行,应改用 ConcurrentLinkedQueue
或加锁。 -
Processor 返回类型需与 Writer 入参严格匹配:ItemProcessor
> → ItemWriter - > 是合法的,但 Spring Batch 的 Chunk 机制会将 List
- Chunk 大小(chunkSize=10)指 10 个 List
,而非 10 个 BookingInfo; - 若单个 bookingId 查出 50 条记录,则一个 chunk 可能包含 10 × 50 = 500 条物理记录;
- 写入时务必 flatMap 扁平化,否则会尝试将 List
当作单个实体保存(报错)。
视为一个“项”(item)。这意味着: - Chunk 大小(chunkSize=10)指 10 个 List
- 考虑性能优化:findDistinctId() 若返回海量 ID,建议配合分页(如 Pageable)或游标式读取,避免内存溢出。
- 启用日志与监控:在 read()、process()、write() 中添加结构化日志,便于追踪数据流与定位重复问题。
通过正确配置作用域并明确各组件职责边界,您即可稳定运行该“ID 驱动、批量查-改-存”的定制流程,彻底规避重复写入风险。











