【发布时间】:2018-02-11 23:44:51
【问题描述】:
我想一次处理多条消息,例如一次从频道获取 10 条消息,并一次将它们写入日志文件。
考虑到这种情况,我可以编写一个服务激活器来获取预定义集中的消息,即 5 或 10 条消息并处理它吗?如果这是不可能的,那么如何使用 Spring Integration 来实现这一点。
【问题讨论】:
我想一次处理多条消息,例如一次从频道获取 10 条消息,并一次将它们写入日志文件。
考虑到这种情况,我可以编写一个服务激活器来获取预定义集中的消息,即 5 或 10 条消息并处理它吗?如果这是不可能的,那么如何使用 Spring Integration 来实现这一点。
【问题讨论】:
这正是您可以通过Aggregator 获得的。您可以使用size() == 10 等简单表达式向群组收集多条消息。当组完成时,DefaultAggregatingMessageGroupProcessor 会发出一条消息,其中包含组中消息的有效负载列表。您可以将结果发送给服务激活器以立即处理该批次。
更新
类似这样的:
.aggregate(aggregator -> aggregator
.correlationStrategy(message -> 1)
.releaseStrategy(group -> group.size() == 10)
.outputProcessor(g -> new GenericMessage<Collection<Message<?>>>(g.getMessages()))
.expireGroupsUponCompletion(true))
因此,我们通过静态 1 键关联消息(对它们进行分组或缓冲)。
该组(或缓冲区大小为10),当我们到达它时,我们发出一条消息,其中包含来自该组的所有消息。发出结果后,我们从该组中清理存储,以便为新的消息序列形成一个新的存储。
【讨论】:
这首先取决于创建消息的原因;如果是消息驱动的通道适配器,那么该适配器中的并发性是关键。
对于其他消息源,您可以使用ExecutorChannel 作为服务激活器的输入通道,执行器的池大小为 10。
根据发送消息的内容,您需要小心在服务器发生故障时丢失消息。
如果没有关于您的应用程序的更多信息,很难提供一般性的答案。
【讨论】:
size() == 10)。