【问题标题】:Rabbit MQ doesn't flush acks?Rabbitmq 不刷新确认?
【发布时间】:2021-09-08 17:11:15
【问题描述】:

问题出现在日志中:Consumer failed to start in 60000 milliseconds; does the task executor have enough threads to support the container concurrency?

我们尝试通过SimpleMessageListenerContainer.addQueueNames() 动态打开类似 50 个队列的处理程序,然后启动应用程序。它消耗了一些消息,但 RabbitMQ 管理面板显示它们未被确认。经过相当长的一段时间后,每个队列的消息最多堆叠 6 条未确认的消息(队列每分钟的消息量相当低),总共有 300 条消息,发生了一些事情,它们都被消费和确认了。当消息未被确认时,容器似乎正在尝试加载另一个消费者,直到它遇到限制。

我们现在是用AUTO的确认方式,在MANUAL的时候就可以了。

有两个问题:

  1. 未确认消息的原因可能是什么?是否有不经常触发的刷新机制?

  2. 如何处理“没有足够的线程”消息?

这两者似乎真的相互关联。

设置如下:

@Bean
fun queueMessageListenerContainer(
    connectionFactory: ConnectionFactory,
    retryOperationsInterceptor: RetryOperationsInterceptor,
    vehicleQueueListenerFactory: QueueListenerFactory,
): SimpleMessageListenerContainer {
    return SimpleMessageListenerContainer().also {
        it.connectionFactory = connectionFactory
        it.setConsumerTagStrategy { queueName -> consumerTag(queueName) }
        it.setMessageListener(vehicleQueueListenerFactory.create())
        it.setConcurrentConsumers(2)
        it.setMaxConcurrentConsumers(5)
        it.setListenerId("queue-consumer")
        it.setAdviceChain(retryOperationsInterceptor)
        it.setRecoveryInterval(RABBIT_HEARTH_BEAT.toMillis())
        //had 10-100 threads, didn't help
        it.setTaskExecutor(rabbitConsumersExecutorService)
        // AUTO suppose to set ack for the messages, right?
        it.acknowledgeMode = AcknowledgeMode.AUTO
    }
}

@Bean
fun connectionFactory(rabbitProperties: RabbitProperties): AbstractConnectionFactory {
    val rabbitConnectionFactory = com.rabbitmq.client.ConnectionFactory().also { connectionFactory ->
        connectionFactory.isAutomaticRecoveryEnabled = true
        connectionFactory.isTopologyRecoveryEnabled = true
        connectionFactory.networkRecoveryInterval = RABBIT_HEARTH_BEAT.toMillis()
        connectionFactory.requestedHeartbeat = RABBIT_HEARTH_BEAT.toSeconds().toInt()
        // was up to 100 connections, didn't help
        connectionFactory.setSharedExecutor(rabbitConnectionExecutorService)
        connectionFactory.host = rabbitProperties.host
        connectionFactory.port = rabbitProperties.port ?: connectionFactory.port
    }
    return CachingConnectionFactory(rabbitConnectionFactory)
        .also {
            it.cacheMode = rabbitProperties.cache.connection.mode
            it.connectionCacheSize = rabbitProperties.cache.connection.size
            it.setConnectionNameStrategy { "simulation-gateway:${springProfiles.firstOrNull()}:event-consumer" }
        }
}

class QueueListenerFactory { 
  fun create(){
     return MessageListener {
        try {
            // no ack, rely on AUTO acknowledgement mode
            handle()
        } catch (e: Throwable) {
            ...
        }
     }
  }
}

【问题讨论】:

  • 有 50 个队列,每个队列有 10 个并发消费者,对我来说似乎很重 (50*10=500)。你确定你需要它吗?您的系统能够处理它吗?您是否尝试减少此数字?
  • @yoni 哎呀,数字错误,感谢您的注意。目前实际上是 2 和 5,我试图将最小值降低到 1,但没有帮助。到目前为止,我的假设是由于未确认的消息,它会尝试使用超过 1 个。
  • 你用的是什么版本?此外,Spring AMQP 不支持 amqp-client 中的自动恢复——它有自己的恢复机制,比自动恢复早很多年。在大多数情况下,此类问题是由于消费者线程“卡”在用户代码中;进行线程转储以查看容器线程在做什么。凝视的延迟可能是由于其他原因;任务执行者问题是最有可能的;再次,线程转储应该会有所帮助。
  • @GaryRussell 版本是 spring-boot 2.3.11.RELEASE 上的 2.2.17.RELEASE。我们有这样的恢复是因为它可以在由于网络问题与兔子失去连接后恢复。你会推荐使用另一种方式吗?在线程转储中,大多数线程都停留在LockSupport.park...SynchronousQueue.take,这似乎意味着它们正在等待消息到达......
  • 是的,它会在没有这些设置的情况下恢复,我们已经进行了十多年的连接恢复。我们通过立即关闭恢复的连接来有效地禁用自动恢复。它给从未消费过任何东西的悬空消费者带来了各种各样的问题。我们早在 2018 年就添加了该代码,从那时起就没有任何恢复问题。除非消费者线程被卡住,否则我看不到如何在 AUTO ack 模式下获得未确认的消息。您需要显示更多信息 - 线程转储、日志、屏幕截图等。

标签: java spring kotlin rabbitmq spring-amqp


【解决方案1】:

好的,我发现了问题所在。基本上,它无法及时启动所有的队列消费者,因为它不仅是 SimpleMessageListenerContainer 太慢的过程慢,而且我们尝试addQueueNames 一个一个。

   userRepository.findAll()
        .map { user -> queueName(user) }
        .onEach { queueName ->
            simpleContainerListener.addQueueNames(queueName)
        }

但SimpleMessageListenerContainer 的以下文档行仍未引起注意:

The existing consumers will be cancelled after they have processed any pre-fetched messages and new consumers will be created 

这意味着实际发生的是 (1, 2, ... N) 消费者的娱乐。更糟糕的是,如果请求来自 API,我们在处理完请求后做了完全相同的simpleContainerListener.addQueueNames(queueName),然后重新创建了所有消费者。

此外,消费者的娱乐是AUTO 确认不起作用的原因:线程挂起试图在超时之前建立足够的消费者。

我通过添加DirectMessageListenerContainer 来处理最近添加的队列来解决此问题,与仅添加一个新消费者的特定情况下的SimpleMessageListenerContainer 相比,这非常快。

DirectMessageListenerContainer(connectionFactory).also {
        it.setConsumerTagStrategy { queueName -> consumerTag(queueName, RECENT_CONSUMER_TAG) }
        it.setMessageListener(ListenerFactory.create())
        it.setListenerId("queue-consumer-recent")
        it.setAdviceChain(retryOperationsInterceptor)
        it.setTaskExecutor(recentQueuesTaskExecutor)
        it.acknowledgeMode = AcknowledgeMode.AUTO
    }

缺点是DirectMessageListenerContainer 在每个实例上每个队列使用 1 个线程。这正是我一开始不想使用它的原因,但是将DirectMessageListenerContainer 用于最近的队列和SimpleContainerListener 用于已经存在的队列显着减少了处理这些队列所需的线程数量。据我了解,DirectMessageListenerContainer 的大量使用最终会导致 OOM,因此下一步可以将队列从直接转移到简单的容器监听器。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2015-09-21
    • 2016-05-09
    • 1970-01-01
    • 2023-03-22
    • 2014-09-30
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多