【问题标题】:How to set the override on a compound trigger?如何在复合触发器上设置覆盖?
【发布时间】:2017-02-22 23:44:20
【问题描述】:

我有一个 Spring 集成应用程序,它通常使用 cron 触发器通过 SFTP 每天轮询文件。但是,如果它没有找到它期望的文件,它应该通过周期性触发器每 x 分钟轮询一次,直到 y 尝试。为此,我使用以下组件:

@Component
public class RetryCompoundTriggerAdvice extends AbstractMessageSourceAdvice {

    private final static Logger logger = LoggerFactory.getLogger(RetryCompoundTriggerAdvice.class);

    private final CompoundTrigger compoundTrigger;

    private final Trigger override;

    private final ApplicationProperties applicationProperties;

    private final Mail mail;

    private int attempts = 0;

    public RetryCompoundTriggerAdvice(CompoundTrigger compoundTrigger, 
            @Qualifier("secondaryTrigger") Trigger override, 
            ApplicationProperties applicationProperties,
            Mail mail) {
        this.compoundTrigger = compoundTrigger;
        this.override = override;
        this.applicationProperties = applicationProperties;
        this.mail = mail;
    }

    @Override
    public boolean beforeReceive(MessageSource<?> source) {
        return true;
    }

    @Override
    public Message<?> afterReceive(Message<?> result, MessageSource<?> source) {
        final int  maxOverrideAttempts = applicationProperties.getMaxFileRetry();
        attempts++;
        if (result == null && attempts < maxOverrideAttempts) {
            logger.info("Unable to find load file after " + attempts + " attempt(s). Will reattempt");
            this.compoundTrigger.setOverride(this.override);
        } else if (result == null && attempts >= maxOverrideAttempts) {
            mail.sendAdminsEmail("Missing File");
            attempts = 0;
            this.compoundTrigger.setOverride(null);
        }
        else {
            attempts = 0;
            this.compoundTrigger.setOverride(null);
            logger.info("Found load file");
        }
        return result;
    }

    public void setOverrideTrigger() {
        this.compoundTrigger.setOverride(this.override);
    }

    public CompoundTrigger getCompoundTrigger() {
        return compoundTrigger;
    }
}

如果文件不存在,这很好用。也就是说,覆盖(即周期性触发)生效并每 x 分钟轮询一次,直到 y 尝试。

但是,如果文件确实存在但它不是预期的文件(例如,数据的日期错误),另一个类(读取文件)调用 RetryCompoundTriggerAdvice 类的 setOverrideTrigger。但是afterReceive 随后不会每隔 x 分钟调用一次。为什么会这样?

更多应用代码如下:

SftpInboundFileSynchronizer:

@Bean
public SftpInboundFileSynchronizer sftpInboundFileSynchronizer() {
    SftpInboundFileSynchronizer fileSynchronizer = new SftpInboundFileSynchronizer(sftpSessionFactory());
    fileSynchronizer.setDeleteRemoteFiles(false);
    fileSynchronizer.setRemoteDirectory(applicationProperties.getSftpDirectory());
    CompositeFileListFilter<ChannelSftp.LsEntry> compositeFileListFilter = new CompositeFileListFilter<ChannelSftp.LsEntry>();
    compositeFileListFilter.addFilter(new SftpPersistentAcceptOnceFileListFilter(store, "sftp"));
    compositeFileListFilter.addFilter(new SftpSimplePatternFileListFilter(applicationProperties.getLoadFileNamePattern()));
    fileSynchronizer.setFilter(compositeFileListFilter);
    fileSynchronizer.setPreserveTimestamp(true);
    return fileSynchronizer;
}

会话工厂是:

@Bean
public SessionFactory<LsEntry> sftpSessionFactory() {
    DefaultSftpSessionFactory sftpSessionFactory = new DefaultSftpSessionFactory();
    sftpSessionFactory.setHost(applicationProperties.getSftpHost());
    sftpSessionFactory.setPort(applicationProperties.getSftpPort());
    sftpSessionFactory.setUser(applicationProperties.getSftpUser());
    sftpSessionFactory.setPassword(applicationProperties.getSftpPassword());
    sftpSessionFactory.setAllowUnknownKeys(true);
    return new CachingSessionFactory<LsEntry>(sftpSessionFactory);
}

SftpInboundFileSynchronizingMessageSource 设置为使用复合触发器进行轮询。

@Bean
@InboundChannelAdapter(autoStartup="true", channel = "sftpChannel", poller = @Poller("pollerMetadata"))
public SftpInboundFileSynchronizingMessageSource sftpMessageSource() {
    SftpInboundFileSynchronizingMessageSource source =
            new SftpInboundFileSynchronizingMessageSource(sftpInboundFileSynchronizer());
    source.setLocalDirectory(applicationProperties.getScheduledLoadDirectory());
    source.setAutoCreateLocalDirectory(true);
    CompositeFileListFilter<File> compositeFileFilter = new CompositeFileListFilter<File>();
    compositeFileFilter.addFilter(new LastModifiedFileListFilter());
    compositeFileFilter.addFilter(new FileSystemPersistentAcceptOnceFileListFilter(store, "dailyfilesystem"));
    source.setLocalFilter(compositeFileFilter);
    source.setCountsEnabled(true);
    return source;
}

@Bean
public PollerMetadata pollerMetadata(RetryCompoundTriggerAdvice retryCompoundTriggerAdvice) {
    PollerMetadata pollerMetadata = new PollerMetadata();
    List<Advice> adviceChain = new ArrayList<Advice>();
    adviceChain.add(retryCompoundTriggerAdvice);
    pollerMetadata.setAdviceChain(adviceChain);
    pollerMetadata.setTrigger(compoundTrigger());
    pollerMetadata.setMaxMessagesPerPoll(1);
    return pollerMetadata;
}

@Bean
public CompoundTrigger compoundTrigger() {
    CompoundTrigger compoundTrigger = new CompoundTrigger(primaryTrigger());
    return compoundTrigger;
}

@Bean
public CronTrigger primaryTrigger() {
    return new CronTrigger(applicationProperties.getSchedule());
}

@Bean
public PeriodicTrigger secondaryTrigger() {
    return new PeriodicTrigger(applicationProperties.getRetryInterval());
}

更新

这是消息处理程序:

@Bean
@ServiceActivator(inputChannel = "sftpChannel")
public MessageHandler dailyHandler(SimpleJobLauncher jobLauncher, Job job, Mail mail) {
    JobRunner jobRunner = new JobRunner(jobLauncher, job, store, mail);
    jobRunner.setDaily("true");
    jobRunner.setOverwrite("false");
    return jobRunner;
}

JobRunner 启动 Spring Batch 作业。处理完作业后,我的应用程序会查看文件是否包含当天预期的数据。如果不是,它正在设置覆盖触发器。

【问题讨论】:

    标签: spring spring-integration


    【解决方案1】:

    这就是触发器的工作方式 - 您只有在触发器触发时才有机会更改触发器。

    由于您重置为 cron 触发器,下一个更改机会是触发器触发时(如果轮询线程在更改触发器之前被下游流释放)。

    您是否将文件移交给另一个线程(队列通道或执行程序)?如果没有,我希望应该应用对触发器的任何更改,因为在下游流返回之前不会调用 nextExecutionTime()。

    如果存在线程切换,则您没有机会更改触发器。

    【讨论】:

    • 谢谢。是的,这与我所知道的不同。所以,这就解释了。在 OP 更新部分添加了相关代码。虽然不调用决定从RetryCompoundTriggerAdvice 中重试的逻辑,但不确定如何实现我的目标。逻辑虽然基于几个因素,包括读取和比较文件中的所有行(这就是我在 Spring Batch 作业完成后应用它的原因)。
    • CompoundTriggerAdvice 提供smart poller 功能。一种可能性是,在处理消息时,将短期轮询器保持在原位,直到评估文件状态。让服务“确认”一切正常的建议,并让它在下一次投票时更改触发器。在处理文件时(以及在确认之前),让 beforeReceive() 返回 false 取消当前轮询
    猜你喜欢
    • 2011-01-15
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多