Symfony Messenger 与延迟消息处理
Symfony Messenger 组件支持延迟消息处理,主要通过 Transport 和 Middleware 实现,以下是完整的实现方案:

基本配置
# config/packages/messenger.yaml
framework:
messenger:
transports:
# 同步传输
sync: 'sync://'
# 使用 Doctrine 传输(支持延迟)
async:
dsn: '%env(MESSENGER_TRANSPORT_DSN)%'
options:
# 延迟队列表
table_name: messenger_messages
# 队列名称
queue_name: default
# 自动设置延迟
auto_setup: true
# 重试策略
retry_strategy:
max_retries: 3
delay: 1000
routing:
'App\Message\YourMessage': async
创建延迟消息
// src/Message/DelayedNotification.php
namespace App\Message;
class DelayedNotification
{
public function __construct(
private string $content,
private \DateTimeInterface $delayedUntil
) {}
public function getContent(): string
{
return $this->content;
}
public function getDelayedUntil(): \DateTimeInterface
{
return $this->delayedUntil;
}
}
使用内置延迟功能
AMQP Transport(RabbitMQ)
# config/packages/messenger.yaml
framework:
messenger:
transports:
async:
dsn: '%env(MESSENGER_TRANSPORT_DSN)%'
options:
# RabbitMQ 原生延迟支持
delay:
# 延迟队列交换机
exchange:
name: delays
type: direct
# 延迟队列
queue:
name: delayed
arguments:
x-dead-letter-exchange: messages
x-message-ttl: 60000
Doctrine Transport(数据库)
// src/MessageHandler/DelayedNotificationHandler.php
namespace App\MessageHandler;
use App\Message\DelayedNotification;
use Symfony\Component\Messenger\Attribute\AsMessageHandler;
use Symfony\Component\Messenger\Bridge\Doctrine\Transport\DoctrineTransport;
#[AsMessageHandler]
class DelayedNotificationHandler
{
public function __invoke(DelayedNotification $message)
{
// 处理延迟消息
$content = $message->getContent();
// 业务逻辑...
}
}
发送延迟消息
// src/Service/NotificationService.php
namespace App\Service;
use App\Message\DelayedNotification;
use Symfony\Component\Messenger\MessageBusInterface;
use Symfony\Component\Messenger\Stamp\DelayStamp;
class NotificationService
{
public function __construct(
private MessageBusInterface $bus
) {}
public function sendDelayedNotification(string $content, int $delayMs = 5000): void
{
$message = new DelayedNotification($content, new \DateTimeImmutable("+{$delayMs}ms"));
// 方法1:使用 DelayStamp
$this->bus->dispatch(
$message,
[new DelayStamp($delayMs)]
);
// 方法2:使用自定义 Middleware
// $this->bus->dispatch($message);
}
}
自定义延迟 Middleware
// src/Messenger/Middleware/DelayMiddleware.php
namespace App\Messenger\Middleware;
use App\Message\DelayedNotification;
use Symfony\Component\Messenger\Envelope;
use Symfony\Component\Messenger\Middleware\MiddlewareInterface;
use Symfony\Component\Messenger\Middleware\StackInterface;
use Symfony\Component\Messenger\Stamp\DelayStamp;
use Symfony\Component\Messenger\Stamp\ReceivedStamp;
class DelayMiddleware implements MiddlewareInterface
{
public function handle(Envelope $envelope, StackInterface $stack): Envelope
{
$message = $envelope->getMessage();
if ($message instanceof DelayedNotification) {
$now = new \DateTimeImmutable();
$delay = $message->getDelayedUntil()->getTimestamp() - $now->getTimestamp();
if ($delay > 0) {
// 添加延迟戳
$envelope = $envelope->with(new DelayStamp($delay * 1000));
}
}
return $stack->next()->handle($envelope, $stack);
}
}
处理延迟消息的 Worker
# 启动消息消费者 php bin/console messenger:consume async --limit=10 --time-limit=3600 # 启动多个消费者 php bin/console messenger:consume async & php bin/console messenger:consume async & # 启动守护进程模式 php bin/console messenger:consume async --sleep=1
高级延迟策略
// src/Messenger/Stamp/CustomDelayStamp.php
namespace App\Messenger\Stamp;
use Symfony\Component\Messenger\Stamp\DelayStamp;
class CustomDelayStamp extends DelayStamp
{
public function __construct(
int $delay,
private string $reason = ''
) {
parent::__construct($delay);
}
public function getReason(): string
{
return $this->reason;
}
}
使用 Redis 实现延迟
# config/packages/messenger.yaml
framework:
messenger:
transports:
async:
dsn: 'redis://localhost:6379/messages'
options:
# Redis 延迟支持
redeliver_timeout: 3600
claim_timeout: 86400
# 自动清理过期消息
auto_setup: true
最佳实践
// 1. 使用 DTO 模式
class OrderTimeoutMessage
{
public function __construct(
private int $orderId,
private \DateTimeImmutable $timeoutAt
) {}
}
// 2. 处理结果
#[AsMessageHandler]
class OrderTimeoutHandler
{
public function __invoke(OrderTimeoutMessage $message): void
{
// 检查订单状态
$order = $this->orderRepository->find($message->getOrderId());
if ($order && $order->getStatus() === OrderStatus::PENDING) {
// 执行超时逻辑
$this->orderService->handleTimeout($order);
}
}
}
// 3. 批量发送延迟消息
class BatchDelaySender
{
public function sendDelayedBatch(array $messages, int $delayMs): void
{
foreach ($messages as $message) {
$this->bus->dispatch($message, [new DelayStamp($delayMs)]);
}
}
}
监控与调试
# 查看队列状态 php bin/console messenger:debug # 监控传输 php bin/console messenger:stats # 清理失败的消息 php bin/console messenger:failed:remove # 重试失败的消息 php bin/console messenger:failed:retry --force
注意事项
- 传输选择:AMQP 传输原生支持延迟,Doctrine 传输通过轮询实现
- 性能考量:大量延迟消息可能影响数据库/队列性能
- 持久化:确保延迟消息不丢失(使用持久化传输)
- 时间精度:延迟时间可能有几秒的误差
- 资源清理:定期清理过期的延迟消息
Symfony Messenger 的延迟功能灵活且可扩展,适合各种异步任务调度场景。