如何在Spring Batch中构建多读取器的作业流程?
在Spring Batch中实现跨步骤数据传递与关联查询的作业流程
完全可以实现你想要的流程,核心思路是跨步骤共享数据,结合自定义读取器完成关联查询。下面分两种场景给出具体实现方案:
方案一:小数据量场景 - 利用Job Execution Context传递数据
适合数据量不大的情况,把第一步读取的userList存入作业执行上下文(Job Execution Context),第二步从上下文取出数据后关联查询历史表。
第一步:读取并收集User数据
第一步不需要写入持久化存储,而是把处理后的User数据收集到作业上下文中:
@Configuration public class BatchConfig { @Autowired private DataSource dataSource; @Autowired private JdbcTemplate jdbcTemplate; // 第一步:读取user表,收集数据到作业上下文 @Bean public Step step1(JobRepository jobRepository, PlatformTransactionManager transactionManager) { return new StepBuilder("step1", jobRepository) .<User, User>chunk(10, transactionManager) .reader(userReader()) .processor(userProcessor()) .writer(userListCollector()) .listener(new StepExecutionListener() { @Override public void beforeStep(StepExecution stepExecution) { // 初始化存储容器 stepExecution.getJobExecution().getExecutionContext().put("userList", new ArrayList<User>()); } @Override public ExitStatus afterStep(StepExecution stepExecution) { return ExitStatus.COMPLETED; } }) .build(); } // 读取user表的id和name @Bean public JdbcCursorItemReader<User> userReader() { JdbcCursorItemReader<User> reader = new JdbcCursorItemReader<>(); reader.setDataSource(dataSource); reader.setSql("SELECT id, name FROM user"); reader.setRowMapper((rs, rowNum) -> new User(rs.getLong("id"), rs.getString("name"))); return reader; } // 自定义User处理器(按需处理数据) @Bean public ItemProcessor<User, User> userProcessor() { return user -> { // 这里添加你的处理逻辑,比如数据校验、字段转换等 return user; }; } // 把User数据收集到作业上下文 @Bean public ItemWriter<User> userListCollector() { return items -> { ExecutionContext jobContext = StepSynchronizationManager.getContext() .getStepExecution().getJobExecution().getExecutionContext(); List<User> userList = (List<User>) jobContext.get("userList"); userList.addAll(items); }; }
第二步:读取上下文数据并关联查询历史表
第二步通过自定义读取器,从作业上下文取出userList,逐个关联查询user_service_history表的数据:
// 第二步:关联查询历史数据并写入目标表 @Bean public Step step2(JobRepository jobRepository, PlatformTransactionManager transactionManager) { return new StepBuilder("step2", jobRepository) .<UserServiceHistory, ProcessedResult>chunk(10, transactionManager) .reader(historyAssociatedReader()) .processor(historyProcessor()) .writer(targetDataWriter()) .build(); } // 自定义读取器:关联user和user_service_history @Bean public ItemReader<UserServiceHistory> historyAssociatedReader() { return new ItemReader<UserServiceHistory>() { private List<User> userList; private Iterator<User> userIterator; private Iterator<UserServiceHistory> historyIterator; @Override public UserServiceHistory read() throws Exception { // 初始化userList if (userList == null) { ExecutionContext jobContext = StepSynchronizationManager.getContext() .getStepExecution().getJobExecution().getExecutionContext(); userList = (List<User>) jobContext.get("userList"); userIterator = userList.iterator(); } // 切换下一个用户的历史数据 if (historyIterator == null || !historyIterator.hasNext()) { if (!userIterator.hasNext()) { return null; // 所有数据处理完毕 } User currentUser = userIterator.next(); // 查询当前用户的服务历史 List<UserServiceHistory> historyList = jdbcTemplate.query( "SELECT id, user_id, status FROM user_service_history WHERE user_id = ?", new Object[]{currentUser.getId()}, (rs, rowNum) -> new UserServiceHistory( rs.getLong("id"), rs.getLong("user_id"), rs.getString("status") ) ); historyIterator = historyList.iterator(); } return historyIterator.next(); } }; } // 处理关联后的数据,转换为目标格式 @Bean public ItemProcessor<UserServiceHistory, ProcessedResult> historyProcessor() { return history -> { // 根据user_id找到对应的User信息 User user = userList.stream() .filter(u -> u.getId().equals(history.getUserId())) .findFirst() .orElseThrow(() -> new IllegalArgumentException("User not found: " + history.getUserId())); // 组装最终结果 return new ProcessedResult(user.getId(), user.getName(), history.getStatus()); }; } // 写入最终目标数据 @Bean public ItemWriter<ProcessedResult> targetDataWriter() { return items -> { // 批量写入目标表(示例为数据库,可替换为文件、消息队列等) jdbcTemplate.batchUpdate( "INSERT INTO target_table (user_id, user_name, status) VALUES (?, ?, ?)", new BatchPreparedStatementSetter() { @Override public void setValues(PreparedStatement ps, int i) throws SQLException { ProcessedResult result = items.get(i); ps.setLong(1, result.getUserId()); ps.setString(2, result.getUserName()); ps.setString(3, result.getStatus()); } @Override public int getBatchSize() { return items.size(); } } ); }; } // 定义作业,串联两个步骤 @Bean public Job userHistoryJob(JobRepository jobRepository, Step step1, Step step2) { return new JobBuilder("userHistoryJob", jobRepository) .start(step1) .next(step2) .build(); } // 实体类需实现Serializable接口,才能存入Execution Context public static class User implements Serializable { private Long id; private String name; // 构造器、getter、setter public User(Long id, String name) { this.id = id; this.name = name; } // getter和setter省略 } public static class UserServiceHistory { private Long id; private Long userId; private String status; // 构造器、getter、setter public UserServiceHistory(Long id, Long userId, String status) { this.id = id; this.userId = userId; this.status = status; } // getter和setter省略 } public static class ProcessedResult { private Long userId; private String userName; private String status; // 构造器、getter、setter public ProcessedResult(Long userId, String userName, String status) { this.userId = userId; this.userName = userName; this.status = status; } // getter和setter省略 } }
方案二:大数据量场景 - 临时表存储中间结果
如果userList数据量很大,存入作业上下文可能导致性能或存储问题,建议用临时表存储第一步的结果:
- 第一步读取
user表后,将id和name写入临时表(如temp_user) - 第二步直接通过SQL关联
temp_user和user_service_history表,用JdbcCursorItemReader读取关联结果,再处理写入目标表
示例第二步读取器:
@Bean public JdbcCursorItemReader<ProcessedResult> joinedReader() { JdbcCursorItemReader<ProcessedResult> reader = new JdbcCursorItemReader<>(); reader.setDataSource(dataSource); reader.setSql(""" SELECT u.id as user_id, u.name as user_name, h.status FROM temp_user u JOIN user_service_history h ON u.id = h.user_id """); reader.setRowMapper((rs, rowNum) -> new ProcessedResult( rs.getLong("user_id"), rs.getString("user_name"), rs.getString("status") )); return reader; }
关键注意事项
- Serializable要求:存入作业上下文的对象必须实现
Serializable接口,否则会序列化失败 - 数据量限制:作业上下文存储在JobRepository中,不适合存放大量数据,大数据量优先选临时表方案
- 事务管理:确保步骤间的事务配置符合业务需求,临时表需注意清理(可在作业结束后删除数据)
内容的提问来源于stack exchange,提问作者박찬오
相关产品推荐
相关产品推荐

