【发布时间】: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