You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spring Batch分区模式下事务与线程同步问题求助

Spring Batch分区作业数据丢失与随机异常问题

问题背景

基于Spring Batch的ItemReader/Processor/ItemWriter架构开发应用,实现了Partitioner进行数据分区,读取数据处理后写回数据库。但作业完成后存在部分数据丢失情况,分区执行失败具有随机性,单线程或无分区场景下运行正常。

遇到的异常

  • java.lang.RuntimeException: java.lang.reflect.UndeclaredThrowableException
  • is 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);
            }
        });
        
    }
}

问题分析与解决方案

核心问题诊断

你的问题根源在于分区场景下的线程安全问题和事务配置/持久化操作不当,具体点如下:

  1. 实体对象跨线程共享:在BatchListener中查询的List<Customer>被放入JobExecutionContext,多个分区线程直接复用同一批Customer实例,Processor中同时修改这些实例的属性(processStatus、address)会引发数据竞争,导致脏写、属性值覆盖,这是数据丢失和随机异常的核心原因。

  2. 事务与持久化错误:

    • Spring Batch的Chunk步骤默认会为每个Chunk绑定事务,但你在Writer中逐个调用Repo方法,若单个操作失败会抛出RuntimeException,可能导致整个Chunk回滚不彻底;
    • 异常is mapped to a primary key column...说明你可能尝试修改主键字段,或对已存在的实体使用了错误的持久化方法(比如用save而非merge);
    • Closed Resultset异常是因为Listener中查询的Customer带有懒加载关联,原Session关闭后线程访问这些属性触发异常。
  3. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.18 21:40:40