Spring Batch分区模式下事务与线程同步问题求助
问题背景
基于Spring Batch的ItemReader/Processor/ItemWriter架构开发应用,实现了Partitioner进行数据分区,读取数据处理后写回数据库。但作业完成后存在部分数据丢失情况,分区执行失败具有随机性,单线程或无分区场景下运行正常。
遇到的异常
java.lang.RuntimeException: java.lang.reflect.UndeclaredThrowableExceptionis mapped to a primary key column in the database. Updates are not allowed.org.eclipse.persistence.exceptions.DatabaseException Internal Exception: java.sql.SQLException: Closed Resultset: getObject
疑问
是否需要维护线程同步或事务机制?
示例场景
- 总记录数:1000
- Chunk大小:100
- 分区1:500条
- 分区2:500条
示例代码
@Configuration @EnableBatchProcessing @EnableTransactionManagement @EnableAspectJAutoProxy(proxyTargetClass = true) public class BatchConfig { @Autowired private JobBuilderFactory jobBuilderFactory; @Autowired private StepBuilderFactory stepBuilderFactory; @Bean public SimpleJobLauncher jobLauncher(JobRepository jobRepository) { SimpleJobLauncher launcher = new SimpleJobLauncher(); launcher.setJobRepository(jobRepository); return launcher; } @Bean(name = "customerJob") public Job prepareBatch1() { return jobBuilderFactory.get("customerJob").incrementer(new RunIdIncrementer()).start(masterStep()).listener(listener()) .build(); } @Bean public Step masterStep() { return stepBuilderFactory.get("masterStep"). partitioner(slaveStep().getName(), partitioner()) .partitionHandler(partitionHandler()) .build(); } @Bean public BatchListener listener() { return new BatchListener(); } @Bean @JobScope public BatchPartitioner partitioner() { return new BatchPartitioner(); } @Bean @StepScope public PartitionHandler partitionHandler() { TaskExecutorPartitionHandler taskExecutorPartitionHandler = new TaskExecutorPartitionHandler(); taskExecutorPartitionHandler.setGridSize(2); taskExecutorPartitionHandler.setTaskExecutor(taskExecutor()); taskExecutorPartitionHandler.setStep(slaveStep()); try { taskExecutorPartitionHandler.afterPropertiesSet(); } catch (Exception e) { } return taskExecutorPartitionHandler; } @Bean @StepScope public Step slaveStep() { return stepBuilderFactory.get("slaveStep").<Customer, CustomerWrapperDTO>chunk(100) .reader(getReader()) .processor(processor()) .writer(writer()) .build(); } @Bean @StepScope public BatchWriter writer() { return new BatchWriter(); } @Bean @StepScope public BatchProcessor processor() { return new BatchProcessor(); } @Bean @StepScope public BatchReader getReader() { return new BatchReader(); } @Bean public TaskExecutor taskExecutor() { ThreadPoolTaskExecutor taskExecutor = new ThreadPoolTaskExecutor(); taskExecutor.setMaxPoolSize(Runtime.getRuntime().availableProcessors()); taskExecutor.setCorePoolSize(Runtime.getRuntime().availableProcessors()); taskExecutor.afterPropertiesSet(); return taskExecutor; } }
class CustomerWrapperDTO { private Address address; private Customer customer; // setter getter for address, customer }
Entity
class Customer { String processStatus; // "U" : unprocessed, "C" : completed, "F" : Failed }
public class BatchListener implements JobExecutionListener { @Autowired private CustomerRepo customerRepo; public BatchListener() { } @Override public void beforeJob(JobExecution jobExecution) { List<Customer> customers; try { customers = customerRepo.getAllUnprocessedCustomer(); } catch (Exception e) { throw new CustomerException("failed in BatchListener", e); } jobExecution.getExecutionContext().put("customers",customers); } @Override public void afterJob(JobExecution jobExecution) { } }
public class BatchPartitioner implements Partitioner { @Value("#{jobExecutionContext[customers]}") private List<Customer> customers; @Override public Map<String, ExecutionContext> partition(int gridSize) { Map<String, ExecutionContext> result = new HashMap<>(); int size = customers.size() / gridSize; List<List<Customer>> lists = IntStream.range(0, customers.size()).boxed() .collect(Collectors.groupingBy(i -> i / size, Collectors.mapping(customers::get, Collectors.toList()))) .values().stream().collect(Collectors.toList()); for (int i = 0; i < gridSize; i++) { ExecutionContext executionContext = new ExecutionContext(); executionContext.putString("name", "Thread_" + i); executionContext.put("customers", lists.get(i)); result.put("partition" + i, executionContext); } return result; } }
@Component @StepScope class BatchReader implements ItemReader<Customer> { private int index; @Value("#{stepExecutionContext[customers]}") private List<Customer> customers; @Override public Customer read() { Customer customer = null; if (index < customers.size()) { customer = customers.get(index); index++; } else { index = 0; } return customer; } }
@Component @StepScope public class BatchProcessor implements ItemProcessor<Customer, CustomerWrapperDTO> { public BatchProcessor() { } @Override public CustomerWrapperDTO process(Customer item) { CustomerWrapperDTO customerWrapper = new CustomerWrapperDTO(); try { // logic to get address Address address = // API call or some business logic. item.setAddress(address); item.setProcessStatus("C"); // Completed } catch(Exception e) { item.setProcessStatus("F");// failed } customerWrapper.setCustomer(item); customerWrapper.setAddress(address); return customerWrapper; } }
@Component @StepScope public class BatchWriter implements ItemWriter<CustomerWrapperDTO> { @Autowired private CustomerRepo customerRepo; @Autowired private AddressRepo addressRepo; public BatchWriter() { } @Override public void write(List<? extends CustomerWrapperDTO> items) { items.forEach(item -> { try { if(item.getCustomer() != null) { customerRepo.merge(item.getCustomer()); } if(item.getAddress() != null) { addressRepo.save(item.getAddress()); } } catch (Exception e) { throw new RuntimeException(e); } }); } }
核心问题诊断
你的问题根源在于分区场景下的线程安全问题和事务配置/持久化操作不当,具体点如下:
实体对象跨线程共享:在
BatchListener中查询的List<Customer>被放入JobExecutionContext,多个分区线程直接复用同一批Customer实例,Processor中同时修改这些实例的属性(processStatus、address)会引发数据竞争,导致脏写、属性值覆盖,这是数据丢失和随机异常的核心原因。事务与持久化错误:
- Spring Batch的Chunk步骤默认会为每个Chunk绑定事务,但你在Writer中逐个调用Repo方法,若单个操作失败会抛出RuntimeException,可能导致整个Chunk回滚不彻底;
- 异常
is mapped to a primary key column...说明你可能尝试修改主键字段,或对已存在的实体使用了错误的持久化方法(比如用save而非merge); Closed Resultset异常是因为Listener中查询的Customer带有懒加载关联,原Session关闭后线程访问这些属性触发异常。
Reader线程状态问题:
BatchReader中的index变量虽因@StepScope每个分区有独立实例,但重置逻辑(index=0)会导致重复读取,可能引发数据重复或异常。
具体修复方案
1. 实现线程隔离的数据源访问
不要在JobExecutionContext中传递实体对象,改为传递分区查询条件:
- 修改
BatchListener:仅记录未处理数据的标记(如processStatus='U'),不查询具体数据; - 修改
BatchPartitioner:按ID范围划分分区,将起始/结束ID放入StepExecutionContext; - 修改
BatchReader:根据StepExecutionContext中的ID范围,独立查询当前分区的Customer数据,确保每个线程操作独立的实体实例。
2. 修正事务与持久化逻辑
- 改用批量持久化方法:将Writer中的forEach循环改为调用
customerRepo.mergeAll()和addressRepo.saveAll(),减少事务提交次数,保证Chunk操作的原子性; - 检查实体映射:确保主键字段不可修改,对已存在的实体使用
merge,新实体使用save; - 关闭实体懒加载:或在Reader查询时使用
fetch join加载关联属性,避免Session关闭后访问懒加载字段。
3. 修复Reader逻辑
移除index=0的重置逻辑,当读取到列表末尾时返回null,符合Spring Batch Reader的规范:
@Override public Customer read() { if (index < customers.size()) { return customers.get(index++); } return null; }
4. 修正代码错误
- 修复
BatchProcessor的返回值错误(原代码返回BatchProcessor,应返回CustomerWrapperDTO); - 修正Repo类名拼写错误:
CustmerRepo→CustomerRepo,addessRepo→addressRepo; - 确保
BatchReader实现ItemReader<Customer>接口。
总结
必须保证线程安全(每个分区线程操作独立的数据和对象),同时依赖Spring Batch的事务机制保证Chunk操作的原子性。不需要手动添加线程同步锁,核心是避免跨线程共享可变状态,正确配置组件作用域与持久化逻辑。
内容的提问来源于stack exchange,提问作者vaibhav

