【问题标题】:How to consume all messages from channel with Spring Integration Java DSL?如何使用 Spring Integration Java DSL 使用来自通道的所有消息?
【发布时间】:2019-08-06 10:13:23
【问题描述】:

我正在尝试为单线程处理程序定义一个流程。消息数量众多,处理程序很慢(一一处理效率低下)。所以我想让处理程序使用 Java DSL 一次性使用通道中可用的所有消息(或等到一些消息积累)。如果通道中没有消息并且处理程序已经处理了前一个组,它应该等待一段时间(超时“a”),以便在通道中积累一些消息。但是如果消息不断出现,处理程序必须在上一次执行后的一段时间后使用它们(超时“b”)。因此,处理程序执行之间的时间间隔不应超过“b”(除非没有消息到达通道)。

没有理由为这种处理程序创建多个实例:它为接口生成数据。下面的代码描述了一些基本配置。我的问题是,一旦处理程序执行完成,我无法想出去抖动(超时“b”)和释放组。

    @Configuration
    public class SomeConfig {

        private AtomicBoolean someHandlerBusy = new AtomicBoolean(false);

        @Bean
        StandardIntegrationFlow someFlow() {
            return IntegrationFlows
                    .from("someChannel")
                    .aggregate(aggregatorSpec -> aggregatorSpec
                                    //The only rule to release a group:
                                    //wait 500ms after last message and have a free someHandler
                                    .groupTimeout(500)
                                    .sendPartialResultOnExpiry(true) //if 500ms expired - send group
                                    .expireGroupsUponCompletion(true) //group should be filled again
                                    .correlationStrategy(message -> true) //one group key, all messages in oe group
                                    .releaseStrategy(message -> false) //never release messages, only with timeout

                                    //Send messages one by one. This is not part of this task.
                                    //I just want to know how to do that. Like splitter.
                                    //.outputProcessor(MessageGroup::getMessages)
                    )
                    .handle("someHandler")
                    .get();
        }
    }

我有纯 Java (kotlin) 代码的解决方案:https://pastebin.com/mti3Y5tD


更新

下面的配置不会删除组。该小组正在不断壮大,但最终却犯了错误。

错误:

*** java.lang.instrument ASSERTION FAILED ***: "!errorOutstanding" with message transform method call failed at JPLISAgent.c line: 844

配置:

    @Configuration
    public class InterfaceHandlerConfigJava {


        @Bean
        MessageChannel interfaceAggregatorFlowChannel() {
            return MessageChannels.publishSubscribe("interfaceAggregatorFlowChannel").get();
        }

        @EventListener(ApplicationReadyEvent.class)
        public void initTriggerPacket(ApplicationReadyEvent event) {
            MessageChannel channel = event.getApplicationContext().getBean("interfaceAggregatorFlowChannel", MessageChannel.class);
            channel.send(MessageBuilder.withPayload(new InterfaceHandler.HandlerReadyMessage()).build());
        }

        @Bean
        StandardIntegrationFlow someFlow(
                InterfaceHandler interfaceHandler
        ) {
            long lastMessageTimeout = 10L;
            return IntegrationFlows
                    .from("interfaceAggregatorFlowChannel")
                    .aggregate(aggregatorSpec -> aggregatorSpec
                            .groupTimeout(messageGroup -> {
                                if (haveInstance(messageGroup, InterfaceHandler.HandlerReadyMessage.class)) {
                                    System.out.println("case HandlerReadyMessage");
                                    if (haveInstance(messageGroup, DbChangeStreamConfiguration.InitFromDbMessage.class)) {
                                        System.out.println("case InitFromDbMessage");
                                        return 0L;
                                    } else if (messageGroup.size() > 1) {
                                        long groupCreationTimeout =
                                                messageGroup.getTimestamp() + 500L - System.currentTimeMillis();
                                        long timeout = Math.min(groupCreationTimeout, lastMessageTimeout);

                                        System.out.println("case messageGroup.size() > 1, timeout: " + timeout);
                                        return timeout;
                                    }
                                }
                                System.out.println("case Handler NOT ReadyMessage");
                                return null;
                            })
                            .sendPartialResultOnExpiry(true)
                            .expireGroupsUponCompletion(true)
                            .expireGroupsUponTimeout(true)
                            .correlationStrategy(message -> true)
                            .releaseStrategy(message -> false)
                    )
                    .handle(interfaceHandler, "handle")
                    .channel("interfaceAggregatorFlowChannel")
                    .get();
        }

        private boolean haveInstance(MessageGroup messageGroup, Class clazz) {
            for (Message<?> message : messageGroup.getMessages()) {
                if (clazz.isInstance(message.getPayload())) {
                    return true;
                }
            }
            return false;
        }
    }

我想强调:这个流程是循环的。有 IN 没有 OUT。消息进入 IN 但处理程序最后发出 HandlerReadyMessage。 也许应该有一些线程中断通道?


最终变体

由于聚合器和处理程序不应相互阻塞,也不应尝试产生 stackoverflow 异常,它们应该在不同的线程中运行。在上面的配置中,这是通过队列通道实现的。看起来发布-订阅频道没有在不同的线程中运行订阅者(至少对于一个订阅者而言)。

    @Configuration
    public class InterfaceHandlerConfigJava {

        // acts as thread breaker too
        @Bean
        MessageChannel interfaceAggregatorFlowChannel() {
            return MessageChannels.queue("interfaceAggregatorFlowChannel").get();
        }

        @Bean
        MessageChannel threadBreaker() {
            return MessageChannels.queue("threadBreaker").get();
        }

        @EventListener(ApplicationReadyEvent.class)
        public void initTriggerPacket(ApplicationReadyEvent event) {
            MessageChannel channel = event.getApplicationContext().getBean("interfaceAggregatorFlowChannel", MessageChannel.class);
            channel.send(MessageBuilder.withPayload(new InterfaceHandler.HandlerReadyMessage()).build());
        }

        @Bean
        StandardIntegrationFlow someFlow(
                InterfaceHandler interfaceHandler
        ) {
            long lastMessageTimeout = 10L;
            return IntegrationFlows
                    .from("interfaceAggregatorFlowChannel")
                    .aggregate(aggregatorSpec -> aggregatorSpec
                            .groupTimeout(messageGroup -> {
                                if (haveInstance(messageGroup, InterfaceHandler.HandlerReadyMessage.class)) {
                                    System.out.println("case HandlerReadyMessage");
                                    if (haveInstance(messageGroup, DbChangeStreamConfiguration.InitFromDbMessage.class)) {
                                        System.out.println("case InitFromDbMessage");
                                        return 0L;
                                    } else if (messageGroup.size() > 1) {
                                        long groupCreationTimeout =
                                                messageGroup.getTimestamp() + 500L - System.currentTimeMillis();
                                        long timeout = Math.min(groupCreationTimeout, lastMessageTimeout);

                                        System.out.println("case messageGroup.size() > 1, timeout: " + timeout);
                                        return timeout;
                                    }
                                }
                                System.out.println("case Handler NOT ReadyMessage");
                                return null;
                            })
                            .sendPartialResultOnExpiry(true)
                            .expireGroupsUponCompletion(true)
                            .expireGroupsUponTimeout(true)
                            .correlationStrategy(message -> true)
                            .releaseStrategy(message -> false)
                            .poller(pollerFactory -> pollerFactory.fixedRate(1))
                    )
                    .channel("threadBreaker")
                    .handle(interfaceHandler, "handle", spec -> spec.poller(meta -> meta.fixedRate(1)))
                    .channel("interfaceAggregatorFlowChannel")
                    .get();
        }

        private boolean haveInstance(MessageGroup messageGroup, Class clazz) {
            for (Message<?> message : messageGroup.getMessages()) {
                if (clazz.isInstance(message.getPayload())) {
                    return true;
                }
            }
            return false;
        }
    }

【问题讨论】:

    标签: spring-integration spring-integration-dsl


    【解决方案1】:

    不清楚您所说的计时器b 是什么意思,但是您可以使用.groupTimeoutExpression(...) 来动态确定组超时时间。

    您无需担心一一发送消息;当输出处理器返回Message&lt;?&gt; 的集合时,它们会一次发送一个。

    【讨论】:

    • 超时 "a" - 组的最后一个元素之后的时间跨度。超时 "b" - 组的第一个元素之后的时间跨度。因此,消息组是在时间跨度“b”内或超时“a”发生之前到达的所有消息。
    • MessageGroupStoreReaper(默认)将根据组创建时间过期。 groupTimeout 根据空闲时间过期。您可以配置两者以满足您的需求。 Reaper docs here。配置两者有点不寻常,但我不明白为什么它不起作用。据推测,b 小于 a。
    • 应该更像“a”小于“b”吗?我刚刚提供了另一个版本的配置。我还没有读到MessageGroupStoreReaper,我认为“如何告诉 Reaper 处理程序已准备好”也会有同样的麻烦。新配置应该有助于理解我的目标和问题。
    • 看起来没问题,但您应该返回1L 而不是0L - 否则版本将在同一个线程上运行,您最终可能会出现堆栈溢出。处理程序如何知道何时发送HandlerReadyMessage?
    • 哦;您现在正在使用默认的 outputProcessor,因此有效负载是有效负载的集合;使用原始问题中的自定义处理器,您将获得发布的离散消息。是的,这样做没问题 - 你不会旋转 CPU,因为队列通道中的默认 receiveTimeout 是 1000。不过,你将有效地永久占用一个调度程序线程(默认情况下有 10 个)所以如果你经常这样做,你可能需要更多的调度线程。
    猜你喜欢
    • 2019-06-29
    • 2015-01-29
    • 2017-12-31
    • 1970-01-01
    • 1970-01-01
    • 2018-04-15
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多