【发布时间】:2016-07-14 21:12:28
【问题描述】:
我在管道中有两个应用程序(我们称它们为 A 和 B)。
应用 A 是一个 Spring Boot / Spring Integration 应用,它从队列 1 读取消息 X,执行一些工作,然后基于队列 2 向队列 2 发出大量消息 Y——对于来自队列 1 的每个 X,大约 300 Y发布到 2。单个线程处理每个 X 中的工作并发布各个消息。优化应用程序 A 使我达到了单个实例每秒可以确认大约 50 X 的地方;因此每秒向队列 2 发布大约 15k Y。
App B 也是一个 Spring Boot / Spring Integration 应用程序,它从队列 2 中读取 Ys,并聚合它们,完成管道。 B 的单个实例每秒处理大约 7-8k Y。
因此,总而言之,使用 1 个 A、2 个 B 和一个相当大的 RabbitMQ 服务器(AWS r3.4xlarge),我每秒达到大约 50 X 和 15000 Y。
我一直在尝试扩大这个过程;至少,我想达到每秒 100 X / 30000 Y。因为这些应用程序中的逻辑适合水平扩展,所以我一直在尝试将部署加倍;即 2 As 和 4 Bs。
但是,放大 As 并没有达到预期的效果; Xs 的确认保持大致稳定在 50/s,Ys 也保持稳定在大约每秒 15k,队列 2 或多或少地保持为空。
仔细检查发现 A 的发布通道处于flow 模式,大概是因为每秒 15k 的速度要塞进队列 2 中。但是,限制似乎不仅仅是一个队列每秒接收 15k Y ;如果我设置更改绑定,以便将 Y 发布到没有任何消费者的队列中,我会达到预期的 100 X / 30k Y 每秒。
为什么我似乎无法以每秒 30k Y 的速度进入队列 2?
其他细节:
- 队列 2 使用 DLX 声明为持久的
- Y 消息发布为非持久性
发布者声明如下:
@Bean
public IntegrationFlow outboundFlow(ConnectionFactory connectionFactory, MessageConverter jsonNodeMessageConverter,
AmqpHeaderMapper headerMapper) {
RabbitTemplate outboundTemplate = new RabbitTemplate(connectionFactory);
outboundTemplate.setMessageConverter(jsonNodeMessageConverter);
return IntegrationFlows
.from(BeanNames.OUTPUT_CHANNEL)
.split()
.handle(Amqp.outboundAdapter(outboundTemplate)
.exchangeName(OUTBOUND_EXCHANGE_NAME)
.headerMapper(headerMapper)
.routingKeyExpression("headers." + ApplicationHeaders.DESTINATION_ROUTING_KEY))
.get();
}
@Bean
public AmqpHeaderMapper headerMapper() {
return new DefaultAmqpHeaderMapper() {
@Override
protected void populateStandardHeaders(Map<String, Object> headers, MessageProperties amqpMessageProperties) {
super.populateStandardHeaders(headers, amqpMessageProperties);
amqpMessageProperties.setDeliveryMode(MessageDeliveryMode.NON_PERSISTENT);
}
};
}
更新
我可以在这里充实我的更多设置。
对于应用程序 A:
- 消费者/线程:16
- 在应用程序 A 中预取 X:50
- 没有发布者确认
- 非持久性
对于应用 B:
- 消费者/线程:16
- 在应用程序 B 中预取 Y:250
- tx 大小:250
在 RabbitMQ(我已经升级到 r3.8xlarge)中:
- Erlang VM 线程高达 384(32 核 x 12)
- tcp recbuf 和 sndbuf 设置为 192kb 每 https://www.rabbitmq.com/networking.html
- sysctl net.ipv4.tcp_max_syn_backlog=8192
此时我什至已经拆分了流程;两个应用程序 A 从一个队列 1 中读取,但发布到两个队列 2——2a 和 2b。队列 1 上的 acks/s 完全没有变化,2a 和 2b 上的 acks/s 之和与只有一个队列 2 时相同。确实,从应用程序 A 到 rabbit 和从应用程序 A 到队列 2a 和 2b 的通道仍然经常进入flow。
所有的改变是我使用的硬件资源比以前更少(因为我升级了 rabbit 服务器)—— rabbit 的 CPU 从未超过 30% 并且内存使用率保持低得离谱。
【问题讨论】:
标签: spring rabbitmq spring-integration