【问题标题】:Launch Spring Batch Job from Spring Integration从 Spring Integration 启动 Spring Batch Job
【发布时间】:2018-03-10 05:27:28
【问题描述】:

我需要从远程 SFTP 服务器下载文件并使用 Spring Batch 处理它们。我已经使用 Spring Integration 实现了代码来下载文件。但我无法从 Spring 集成组件启动 Spring Batch 作业。我有以下代码:

    @Autowired
private JobLauncher jobLauncher;

public String OUTPUT_DIR = "temp_dir";

@Value("${sftp.remote.host}")
private String sftpRemoteHost;

@Value("${sftp.remote.user}")
private String sftpUsername;

@Value("${sftp.remote.password}")
private String sftpPassword;

@Value("${sftp.remote.folder}")
private String sftpFolder;

@Bean
public DefaultSftpSessionFactory sftpSessionFactory() {
    final DefaultSftpSessionFactory factory = new DefaultSftpSessionFactory();
    factory.setHost(sftpRemoteHost);
    factory.setAllowUnknownKeys(true);
    factory.setUser(sftpUsername);
    factory.setPassword(sftpPassword);
    return factory;
}

@Bean
public SftpInboundFileSynchronizer sftpInboundFileSynchronizer() {
    final SftpInboundFileSynchronizer fileSynchronizer = new SftpInboundFileSynchronizer(sftpSessionFactory());
    fileSynchronizer.setDeleteRemoteFiles(false);
    fileSynchronizer.setRemoteDirectory(sftpFolder);
    fileSynchronizer.setFilter(new SftpSimplePatternFileListFilter("*.csv"));
    return fileSynchronizer;
}

@Bean
@InboundChannelAdapter(channel = "sftpChannel", poller = @Poller(fixedDelay = "5000"))
public MessageSource<File> sftpMessageSource() {
    final SftpInboundFileSynchronizingMessageSource source =
            new SftpInboundFileSynchronizingMessageSource(sftpInboundFileSynchronizer());
    source.setLocalDirectory(new File(OUTPUT_DIR));
    source.setAutoCreateLocalDirectory(true);
    source.setLocalFilter(new AcceptOnceFileListFilter<>());
    return source;
}

@Bean
@ServiceActivator(inputChannel = "sftpChannel")
public MessageHandler handler() {
    final FileWritingMessageHandler handler = new FileWritingMessageHandler(new File(OUTPUT_DIR));
    handler.setFileExistsMode(FileExistsMode.REPLACE);
    handler.setExpectReply(true);
    handler.setOutputChannelName("parse-csv-channel");
    return handler;
}

@ServiceActivator(inputChannel = "parse-csv-channel", outputChannel = "job-channel")
public JobLaunchRequest adapt(final File file) throws Exception {
    final JobParameters jobParameters = new JobParametersBuilder().addString(
            "input.file", file.getAbsolutePath()).toJobParameters();
    return new JobLaunchRequest(batchConfiguration.job(), jobParameters);
}

@ServiceActivator(inputChannel = "job-channel", outputChannel = "finish")
public JobLaunchingMessageHandler jobHandler(JobLaunchRequest request) throws JobExecutionException {
    return new JobLaunchingMessageHandler(jobLauncher);//.launch(request);
}

@ServiceActivator(inputChannel = "finish")
public void finish() {
    System.out.println("FINISH");
}

但这不起作用(最后一个方法 adapt 中的错误),因为找不到 File 类型的 bean。我无法将这两个部分放在一起。如何接线集成和批处理?

【问题讨论】:

    标签: java spring spring-integration spring-batch


    【解决方案1】:

    您肯定只需要从您的adapt() 方法中删除@Bean 注释。如果我们真的构建MessageHandler bean,我们需要@Bean,例如JobLaunchingMessageHandler 来接受JobLaunchRequest 有效负载:https://docs.spring.io/spring-batch/trunk/reference/html/springBatchIntegration.html#launching-batch-jobs-through-messages

    在参考手册中查看有关消息注释的更多信息:https://docs.spring.io/spring-integration/docs/4.3.12.RELEASE/reference/html/configuration.html#annotations_on_beans

    更新

    @Bean
    @ServiceActivator(inputChannel = "sftpChannel")
    public MessageHandler handler() {
        final FileWritingMessageHandler handler = new FileWritingMessageHandler(new File(OUTPUT_DIR));
        handler.setFileExistsMode(FileExistsMode.REPLACE);
        handler.setExpectReply(true);
        handler.setOutputChannelName("parse-csv-channel");
        return handler;
    }
    
    @ServiceActivator(inputChannel = "parse-csv-channel", outputChannel = "job-channel")
    public JobLaunchRequest adapt(final File file) throws Exception {
        final JobParameters jobParameters = new JobParametersBuilder().addString(
                "input.file", file.getAbsolutePath()).toJobParameters();
        return new JobLaunchRequest(batchConfiguration.job(), jobParameters);
    }
    
    @Bean
    @ServiceActivator(inputChannel = "job-channel")
    public JobLaunchingGateway jobHandler() {
        JobLaunchingGateway jobLaunchingGateway = new JobLaunchingGateway(jobLauncher);
        jobLaunchingGateway.setOutputChannelName("finish");
        return jobLaunchingGateway;
    }
    

    【讨论】:

    • 以及如何结合JobLaunchingMessageHandler 和我的FileWritingMessageHandler
    • 不确定您是否担心。您应该从adapt() 方法中删除@Bean 并将outputChannel 添加到其@ServiceActivator 以将结果发送到JobLaunchingMessageHandler 端点定义。一个应该类似于你使用FileWritingMessageHandler
    • 我更新了原始问题 - 添加了新代码。正如您提到的,我尝试添加更多配置,但它不起作用。我在最后两种方法中做了调试点,这段代码从未调用过。我做错了什么?
    • FileWritingMessageHandler 必须与setExpectReply(true) 一起使用——尽管是默认值。并且文件将真正被发送到parse-csv-channel 进行JobLaunchRequest 转换。您的 launch() 方法是多余的。关于此事,您有JobLaunchingMessageHandler。就是这样
    • 不知道怎么回事。我再次更新了代码。批处理作业不会启动,但文件下载每 5 秒运行一次。并且在每个下载作业不运行之后。这让我抓狂
    猜你喜欢
    • 1970-01-01
    • 2017-04-29
    • 2015-11-03
    • 2013-12-11
    • 2017-07-02
    • 1970-01-01
    • 1970-01-01
    • 2019-10-17
    • 2018-09-19
    相关资源
    最近更新 更多