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

如何将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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 09:55:20