【问题标题】:dealing with flow while scaling rabbitmq在缩放rabbitmq时处理流量
【发布时间】: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)中:

此时我什至已经拆分了流程;两个应用程序 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


    【解决方案1】:

    我假设您没有在生产者端使用交易渠道。

    尝试在 Y 消费者侦听器容器上增加 prefetchCounttxSize;第一个将在容器中缓冲消息;第二个将减少 ack 流量。

    编辑

    有关 RabbitMQ 的详细信息,请参阅 Finding bottlenecks with RabbitMQ

    【讨论】:

    • 我已经将 prefetchCount 设置为 250 左右,但我会提高它并尝试将 tx 大小提高一大堆。
    • 这个链接刚刚发布到rabbitmq-users google 组,以回应另一个问题; rabbitmq.com/blog/2014/04/14/…
    • 是的;我一直在查看该链接,并在rabbitmq.com/networking.html 上调整吞吐量,但到目前为止还没有特别有用。我肯定处于它描述的第二种情况——流模式下的连接和通道,但不是队列。但是,服务器上的 CPU 很好,没有可察觉的磁盘负载,可能是因为消息不是持久的。
    • 嗨,加里——尽管你早先的帮助,我仍然对上述内容感到困惑;我已经澄清了我的所有设置。如果您有机会并有任何想法或提示,我将不胜感激。
    猜你喜欢
    • 2015-08-07
    • 2020-12-07
    • 2020-03-29
    • 1970-01-01
    • 2018-09-30
    • 1970-01-01
    • 2013-10-27
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多