PHP项目Symfony messenger与延迟

wen PHP项目 1

Symfony Messenger 与延迟消息处理

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

PHP项目Symfony messenger与延迟

基本配置

# 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

注意事项

  1. 传输选择:AMQP 传输原生支持延迟,Doctrine 传输通过轮询实现
  2. 性能考量:大量延迟消息可能影响数据库/队列性能
  3. 持久化:确保延迟消息不丢失(使用持久化传输)
  4. 时间精度:延迟时间可能有几秒的误差
  5. 资源清理:定期清理过期的延迟消息

Symfony Messenger 的延迟功能灵活且可扩展,适合各种异步任务调度场景。

抱歉,评论功能暂时关闭!