【问题标题】:Spring Integration: poller with fixedDelay for the whole flow executionSpring Integration:具有固定延迟的轮询器,用于整个流程执行
【发布时间】:2018-03-24 21:35:57
【问题描述】:

在我的集成流程中,我使用 JdbcPollingChannelAdapter 来提供实体键列表、处理这些实体并将下一个 sql 查询的时间窗口设置为下一次迭代(基于实体的最后修改时间)。 如果出现处理延迟或处理错误,则不应发生下一次迭代,并且 MessageSource 应等待当前流执行成功或发生超时。所以在每个时刻最大。应该处理一个键列表。

有没有优雅的方法来设置 Pollers.fixedDelay(...) 不是为 MessageSource 而是为整个流程执行,以便在当前完成后开始下一个流程执行?

    @Test
    public void testDelayedExecutionSequence() {
    final Queue<List<Integer>> inQueue = new LinkedBlockingQueue<>();
    MessageSource<List<Integer>> inMs = new AbstractMessageSource<List<Integer>>() {
        @Override
        public String getComponentType() {
            return null;
        }

        @Override
        protected Object doReceive() {
            return inQueue.poll();
        }
    };

    final int messagesPerStep = 100;
    final int maxIterations = 10;
    for(int iteration = 1, from=1, to=messagesPerStep; iteration <= maxIterations; iteration++, from += messagesPerStep, to += messagesPerStep) {
        System.out.println(String.format("add list from=%d, to=%d", from, to));
        inQueue.add(IntStream.rangeClosed(from, to).boxed().collect(Collectors.toList()));
    }

    final AtomicInteger amqpSendCounter = new AtomicInteger();
    final AtomicInteger iterationCounter = new AtomicInteger();
    final List<Integer> resultSequence = new ArrayList();

    IntegrationFlow integrationFlow = IntegrationFlows
            .from(inMs, c->c.poller(Pollers.fixedDelay(50).maxMessagesPerPoll(1)))
            .split()
            .channel(c -> c.executor(Executors.newFixedThreadPool(10)))
            .<Integer, Integer>transform(p -> {
            if(p == 405) {
                try {
                    Thread.currentThread().sleep(1000);
                } catch (InterruptedException e) {
                }
            }
            return p;})             
            .handle((p, h) -> {amqpSendCounter.incrementAndGet(); return p;})
            .aggregate()
            .log(l -> "!!!!!!!!!!!!!!!!!!!!!!!!!!!!!!!!!!! after agregate(): "+l)
            .<List<Integer>>handle(m -> {iterationCounter.incrementAndGet();
                Integer firstIntInList = ((List<Integer>)m.getPayload()).get(0);
                resultSequence.add(firstIntInList / 100);
            })
            .get();

    IntegrationFlowRegistration registration = this.flowContext.registration(integrationFlow).register();

    try {
        Thread.currentThread().sleep(3000);
    } catch (InterruptedException e) {
    }

    assertThat(iterationCounter.get()).as("iterationCounter.get").isEqualTo(maxIterations);
    assertThat(resultSequence).as("result sequence").isEqualTo(IntStream.rangeClosed(0, 9).boxed().collect(Collectors.toList()));
}

【问题讨论】:

    标签: spring spring-integration


    【解决方案1】:

    您需要查看Conditional polling,并且在更改某些状态之前不要让调用receive()PollSkipAdvice 应该是您的不错选择。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2012-11-03
      • 1970-01-01
      • 2023-03-31
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多