如何将Spring Integration消息负载传递给Spring Batch作业?
解决方案:Spring Integration FTP流传递至Spring Batch作业
针对你遇到的问题——无法将FTP流消息负载直接传入Spring Batch的JobParameters(仅支持基本类型),可通过以下两种方案实现消息负载的传递、解析及数据库存储:
方案一:小文件场景——将负载存入JobExecution上下文
适合文件内容较小、可直接转为字符串存入内存的场景:
1. 修正JobLaunchRequest转换器
将消息头中的FTP文件名作为JobParameters的唯一标识参数(避免重复执行),同时利用Spring Batch Integration的特性自动将消息负载存入JobExecution的ExecutionContext:
@Transformer(inputChannel = "BATCH_ALARM_CHANNEL", outputChannel = "jobLaunchRequestChannel") public JobLaunchRequest toRequest(Message<String> message) { JobParametersBuilder jobParametersBuilder = new JobParametersBuilder(); // 从消息头提取FTP文件名,作为作业实例的唯一标识 String fileName = message.getHeaders().get(FileHeaders.FILENAME, String.class); jobParametersBuilder.addString("fileName", fileName, true); return new JobLaunchRequest(customJob(), jobParametersBuilder.toJobParameters()); }
2. 配置JobLaunchingGateway
确保网关将原始消息的上下文传递至JobExecution:
@Bean @ServiceActivator(inputChannel = "jobLaunchRequestChannel") public JobLaunchingGateway jobLaunchingGateway(JobLauncher jobLauncher) { JobLaunchingGateway gateway = new JobLaunchingGateway(jobLauncher); gateway.setHeaderMapper(new DefaultJobHeadersMapper()); return gateway; }
3. 自定义ItemReader读取负载
从JobExecution的ExecutionContext中获取字符串负载,用于后续解析:
@Component @StepScope public class StringPayloadItemReader implements ItemStreamReader<String> { @Value("#{jobExecution.executionContext['" + IntegrationMessageHeaderAccessor.INTEGRATION_MESSAGE_PAYLOAD + "']}") private String fileContent; private boolean readCompleted = false; @Override public String read() { if (!readCompleted) { readCompleted = true; return fileContent; } return null; // 标记读取完成 } @Override public void open(ExecutionContext executionContext) { readCompleted = false; } @Override public void update(ExecutionContext executionContext) {} @Override public void close() {} }
4. 实现解析与存储逻辑
- ItemProcessor:将字符串解析为业务实体
@Component public class FileContentProcessor implements ItemProcessor<String, YourEntity> { @Override public YourEntity process(String content) { // 示例:解析CSV格式内容 String[] fields = content.split(","); YourEntity entity = new YourEntity(); entity.setField1(fields[0]); entity.setField2(fields[1]); // 填充其他字段 return entity; } }
- ItemWriter:将实体保存至数据库
@Component public class EntityWriter implements ItemWriter<YourEntity> { private final YourEntityRepository repository; public EntityWriter(YourEntityRepository repository) { this.repository = repository; } @Override public void write(List<? extends YourEntity> items) { repository.saveAll(items); } }
5. 配置Batch作业与步骤
@Bean public Job customJob(Step processingStep) { return jobBuilderFactory.get("customFileProcessingJob") .start(processingStep) .build(); } @Bean public Step processingStep(StringPayloadItemReader reader, FileContentProcessor processor, EntityWriter writer) { return stepBuilderFactory.get("fileProcessingStep") .<String, YourEntity>chunk(1) .reader(reader) .processor(processor) .writer(writer) .build(); }
方案二:大文件场景——直接读取FTP流
适合文件较大、无法一次性存入内存的场景,保持InputStream传递并逐行读取:
1. 调整转换器与负载传递
去掉StreamTransformer(保持InputStream作为消息负载),修改转换器:
@Transformer(inputChannel = "BATCH_ALARM_CHANNEL", outputChannel = "jobLaunchRequestChannel") public JobLaunchRequest toRequest(Message<InputStream> message) { JobParametersBuilder jobParametersBuilder = new JobParametersBuilder(); String fileName = message.getHeaders().get(FileHeaders.FILENAME, String.class); jobParametersBuilder.addString("fileName", fileName, true); return new JobLaunchRequest(customJob(), jobParametersBuilder.toJobParameters()); }
2. 自定义FTP流ItemReader
根据JobParameters中的文件名重新从FTP获取流(避免上下文流失效),逐行读取解析:
@Component @StepScope public class FtpStreamItemReader implements ItemStreamReader<YourEntity> { private final FtpRemoteFileTemplate ftpRemoteFileTemplate; private final String remoteDir; private BufferedReader reader; @Value("#{jobParameters['fileName']}") private String fileName; public FtpStreamItemReader(FtpRemoteFileTemplate ftpRemoteFileTemplate, @Value("${ftp.root-dir}") String remoteDir) { this.ftpRemoteFileTemplate = ftpRemoteFileTemplate; this.remoteDir = remoteDir; } @Override public YourEntity read() throws IOException { String line = reader.readLine(); if (line == null) { return null; } // 解析行内容为实体 String[] fields = line.split(","); YourEntity entity = new YourEntity(); entity.setField1(fields[0]); entity.setField2(fields[1]); return entity; } @Override public void open(ExecutionContext executionContext) throws IOException { // 从FTP重新获取文件流 InputStream stream = ftpRemoteFileTemplate.get(remoteDir + "/" + fileName); reader = new BufferedReader(new InputStreamReader(stream, Charset.defaultCharset())); } @Override public void update(ExecutionContext executionContext) {} @Override public void close() throws IOException { if (reader != null) { reader.close(); } } }
3. 复用方案一中的Processor、Writer及作业配置
仅需将Step中的Reader替换为FtpStreamItemReader即可。
关键注意事项
- 作业唯一性:通过
fileName作为JobParameter的标识参数(第三个参数设为true),确保同一文件不会重复执行。 - 资源释放:大文件场景需确保InputStream和Reader在Step结束后正确关闭,避免资源泄漏。
- 作业重启:大文件场景通过文件名重新获取FTP流,避免上下文流失效导致重启失败。
内容的提问来源于stack exchange,提问作者Jekshenov Chingiz
相关产品推荐
相关产品推荐

