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

如何在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数据量很大,存入作业上下文可能导致性能或存储问题,建议用临时表存储第一步的结果:

  1. 第一步读取user表后,将id和name写入临时表(如temp_user)
  2. 第二步直接通过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,提问作者박찬오

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 11:18:43