【发布时间】:2016-11-05 17:28:13
【问题描述】:
我正在为一个应用程序使用 Symfony 和 RabbitMQ 包并遇到以下问题:当消费者服务抛出未捕获的异常/错误(例如:内存不足)时,消息被重新发布并一次又一次地使用,直到它得到拒绝或确认信号。我想更改此行为,以便在第一次使用消息时发生任何未捕获的异常/错误时丢弃消息。
这可能吗?如果可以,怎么做?谢谢!
【问题讨论】:
我正在为一个应用程序使用 Symfony 和 RabbitMQ 包并遇到以下问题:当消费者服务抛出未捕获的异常/错误(例如:内存不足)时,消息被重新发布并一次又一次地使用,直到它得到拒绝或确认信号。我想更改此行为,以便在第一次使用消息时发生任何未捕获的异常/错误时丢弃消息。
这可能吗?如果可以,怎么做?谢谢!
【问题讨论】:
是的,您必须确认消息。您可以通过将自动确认标志设置为 true(取决于您使用的语言/API/库)或手动/显式确认消息来执行此操作。确认无法处理的消息是完全正常的,否则就像你说的message is republished and consumed again and again。
如果你愿意,你也可以set requeue to false。我不使用 PHP 来处理 RabbitMQ,所以我不知道 API 等价物是什么,这就是 nack 的实现位置/方式 - 在这种情况下(不是重新排队)它可能配置dead letter exchange 是个好主意(引用链接):
来自队列的消息可能是“死信”;也就是说,重新发布到 发生以下任何事件时的另一次交换:
邮件被拒绝(basic.reject 或 basic.nack) 重新排队=假,
...
【讨论】:
查看下面的示例,它应该可以解决您的问题。在这种情况下:
解决方案
class OrderCreateConsumer implements ConsumerInterface
{
public function execute(AMQPMessage $message)
{
$body = json_decode($message->body, true);
try {
// Do whatever you want with $body
} catch (Exception $e) {
return ConsumerInterface::MSG_REJECT;
}
}
}
或完整的 symfony + RabbitMQ 示例细节在这里:The proper ways of handling errors in symfony RabbitMQ consumer。看起来选项 2,3 和 4 适用于您,这是您通过问题的声音试图避免的。
【讨论】:
到目前为止,我能找到的最接近的解决方案是扩展 RabbitMq 包。在BaseConsumer 类(命名空间OldSound\RabbitMqBundle\RabbitMq)中,有一个名为setupConsumer 的方法,如下所示:
protected function setupConsumer()
{
if ($this->autoSetupFabric) {
$this->setupFabric();
}
$this->getChannel()->basic_consume($this->queueOptions['name'], $this->getConsumerTag(), false, false, false, false, array($this, 'processMessage'));
}
basic_consume 方法的第四个参数称为$no_ack,并设置为false。当此参数设置为true 时,消息在处理后将被丢弃,无论它是否抛出错误或抛出异常或一切顺利。因此,无论哪种方式,消息都会被丢弃。
请记住,当 $no_ack 参数设置为 true 时,消费者返回的状态并不重要,因此返回 ConsumerInterface::MSG_REJECT_REQUEUE 将没有任何效果 - 消息仍会被丢弃。
【讨论】:
AMQP 协议为我们提供了一条消息redelivered property。
我不知道如何在 RabbitMQBundle 中执行此操作(但我完全确信这是可能的)。虽然我可以向您展示如何使用 Enqueue lib:
<?php
use Enqueue\Psr\PsrContext;
use Enqueue\Psr\PsrMessage;
use Enqueue\Psr\PsrProcessor;
class FooProcessor implements PsrProcessor
{
public function process(PsrMessage $message, PsrContext $context)
{
if ($message->isRedelivered()) {
// we already tried to process this message and failed.
return self::REJECT;
}
// this is the new message we've never seen before.
// do the job
return self::ACK;
}
}
解决此类问题的另一种方法是延迟重新传递的消息,这同样可以使用 RabbitMQBundle 完成,但您必须手动配置每个部分,其中 enqueue bundle 执行此操作。您只需要设置延迟插件并打开配置选项。更多内容请关注post
try catch 并不总是可靠的。您可能会收到无法获取的错误,例如致命错误。
自动确认模式也不好,因为您可能会在没有任何通知或警告或恢复能力的情况下丢失失败的消息。
【讨论】: