【问题标题】:Spring integration aggregate messages every 10 secondsSpring 集成每 10 秒聚合一次消息
【发布时间】:2021-02-02 08:09:56
【问题描述】:

在我的一个流程中,我正在查看每 10 秒聚合一次,并将这些有效负载写入文件共享。我对使用聚合器不是很清楚。

@Bean
public IntegrationFlow errorHandlingQueueFlow() {
    return IntegrationFlows.from(ERROR_QUEUE_CHANNEL)
            .bridge(e -> e.poller(Pollers.fixedDelay(1000).maxMessagesPerPoll(MAX_MSG_PER_POLL)))                
            .aggregate(a -> a.groupTimeout(10000))// How do i make it collect every 10 seconds.
            .transform(objectToCSVTransformer, "transform")//Converts payload to a CSV
            .handle(smbErrorMessageHandler())// Takes care of writing into Fileshare
            .get();
}

由于这是用于错误处理,因此只有少数出错的消息会进入此 ERROR_QUEUE_CHANNEL。所以我想每 10 秒收集一次,而不是等待收到来自一个组的所有消息。当我使用 grouptimeout 时,每 10 秒将收集到的所有消息发送到 nullchannel。

【问题讨论】:

    标签: spring spring-integration spring-integration-dsl


    【解决方案1】:

    groupTimeout() 的默认用途是清理过期的组。如果您想正常释放它们而不是丢弃它们,您应该考虑使用sendPartialResultOnExpiry = true。当然,如果您在这些消息中确实有相关详细信息标头,那么这一切都是有意义的。否则,您需要考虑 correlationStrategy 来对这些错误消息进行分组。

    请阅读文档中有关聚合器及其选项的更多信息:https://docs.spring.io/spring-integration/docs/current/reference/html/message-routing.html#aggregator

    【讨论】:

    • .aggregate(a -> a.groupTimeout(10000).sendPartialResultOnExpiry(true)) 有效。我实际上不必关联这些消息。我只想汇总所有出错的记录,并定期将它们写入文件。如果任何消息在 groupTimeout 之后到达,它们会被丢弃还是会在下一个周期中聚合?
    • 好吧,如果没有相关性,那么聚合器中就没有原因。您的消息将作为一组发布。如果您在存储到文件之前谈到一些延迟,请考虑改用delay()
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2013-11-27
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2011-10-21
    相关资源
    最近更新 更多