如何在Spring Batch ItemReader中实现数据库表的过滤与关联查询
实现方案
1. 作业参数传递
首先配置Spring Batch作业参数转换器支持LocalDateTime类型,避免参数类型转换异常。如果通过命令行触发作业,可以将时间参数以ISO标准字符串格式传入,框架会自动完成转换。
在需要使用参数的组件(如ItemReader)上添加@StepScope注解,通过SpEL表达式注入作业参数:
@StepScope @Bean public JdbcPagingItemReader<EntityA> entityAReader( @Value("#{jobParameters['timeParam1']}") LocalDateTime timeParam1, @Value("#{jobParameters['timeParam2']}") LocalDateTime timeParam2, DataSource dataSource ) { // Reader配置逻辑 }
2. ItemReader实现(核心逻辑)
直接使用JdbcPagingItemReader分页读取数据,避免全量加载导致内存溢出。查询语句直接使用你提供的去重SQL(补齐语法缺失的右括号),在数据库层面完成关联、时间过滤、去重操作,性能远高于Java内存层处理:
SELECT * FROM TableA WHERE nr IN (SELECT DISTINCT a.nr FROM TableA a JOIN TableB b ON a.ID = b.ID WHERE b.TIME BETWEEN ? AND ?)
Reader完整配置示例:
JdbcPagingItemReader<EntityA> reader = new JdbcPagingItemReader<>(); reader.setDataSource(dataSource); reader.setPageSize(1000); // 按需调整分页大小 // 配置查询参数 Map<String, Object> paramMap = new HashMap<>(); paramMap.put("timeParam1", timeParam1); paramMap.put("timeParam2", timeParam2); reader.setParameterValues(paramMap); // 配置行映射,将查询结果转为EntityA实体 reader.setRowMapper(new BeanPropertyRowMapper<>(EntityA.class)); // 配置分页查询SQL SqlPagingQueryProviderFactoryBean queryProvider = new SqlPagingQueryProviderFactoryBean(); queryProvider.setDataSource(dataSource); queryProvider.setSelectClause("*"); queryProvider.setFromClause("TableA"); queryProvider.setWhereClause("nr IN (SELECT DISTINCT a.nr FROM TableA a JOIN TableB b ON a.ID = b.ID WHERE b.TIME BETWEEN :timeParam1 AND :timeParam2)"); queryProvider.setSortKey("nr"); // 选择TableA的唯一键作为排序字段,保证分页一致性 reader.setQueryProvider(queryProvider.getObject()); return reader;
该Reader输出的每条EntityA记录天然去重,不需要额外处理,完全符合第一阶段需求。
3. ItemProcessor实现
实现ItemProcessor<EntityA, List<String>>接口,完成单个EntityA到字符串列表的转换逻辑:
public class EntityAProcessor implements ItemProcessor<EntityA, List<String>> { @Override public List<String> process(EntityA item) throws Exception { // 补充你的转换逻辑,返回对应字符串列表 return List.of("line1", "line2"); } }
4. ItemWriter实现
可以直接使用Spring Batch自带的FlatFileItemWriter配合装饰器处理列表输出,实现逐行写入文件的需求:
@Bean public FlatFileItemWriter<String> flatFileItemWriter() { FlatFileItemWriter<String> writer = new FlatFileItemWriter<>(); writer.setResource(new FileSystemResource("output.txt")); writer.setLineAggregator(new PassThroughLineAggregator<>()); return writer; } // 包装Writer,支持处理List<String>类型的输入 @Bean public ItemWriter<List<String>> listItemWriter(FlatFileItemWriter<String> delegate) { return items -> { List<String> allLines = items.stream() .flatMap(Collection::stream) .toList(); delegate.write(allLines); }; }
5. 作业组装
将上述组件组装为Step和Job即可:
@Bean public Step processStep(JdbcPagingItemReader<EntityA> reader, EntityAProcessor processor, ItemWriter<List<String>> writer, JobRepository jobRepository, PlatformTransactionManager transactionManager) { return new StepBuilder("processStep", jobRepository) .<EntityA, List<String>>chunk(1000, transactionManager) .reader(reader) .processor(processor) .writer(writer) .build(); } @Bean public Job batchJob(Step processStep, JobRepository jobRepository) { return new JobBuilder("batchJob", jobRepository) .start(processStep) .build(); }
内容的提问来源于stack exchange,提问作者Manu
相关产品推荐
相关产品推荐

