You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Symfony Messenger:派发新消息时如何判断是否处于消息处理器中

Symfony Messenger 消息层级追溯:获取父消息ID的方案

需求背景

我需要记录Symfony Messenger中所有消息的全生命周期信息:

  • 消息在<timestamp>时刻派发
  • 消息在<timestamp>时刻处理/失败/重试

目前已实现基础的事件监听逻辑:

已创建的事件订阅者

final readonly class MessengerAuditSubscriber implements EventSubscriberInterface
{
    public static function getSubscribedEvents(): array
    {
        return [
            SendMessageToTransportsEvent::class => 'onSendMessageToTransportsEvent',
            WorkerMessageFailedEvent::class => 'onWorkerMessageFailedEvent',
            WorkerMessageHandledEvent::class => 'onWorkerMessageHandledEvent',
            WorkerMessageReceivedEvent::class => 'onWorkerMessageReceivedEvent',
            WorkerMessageRetriedEvent::class => 'onWorkerMessageRetriedEvent',
        ];
    }

    // ...
}

为派发消息添加自定义ID Stamp

public function onSendMessageToTransportsEvent(SendMessageToTransportsEvent $event): void
{
    $envelope = $event->getEnvelope();
    $envelope = $envelope->with(new MessageIdStamp(uniqid()));
    $event->setEnvelope($envelope);

    // ...记录事件日志
}

在后续事件中获取消息ID

public function onWorkerMessageHandledEvent(WorkerMessageHandledEvent $event): void
{
    $envelope = $event->getEnvelope();
    $messageIdStamp = $envelope->last(MessageIdStamp::class);

    // ...记录事件日志
}

当前待解决问题

现在需要实现:当某个消息处理器内部派发新消息时,为新消息记录父消息ID,形成可追溯的消息层级:

  • ✉️ 消息 1 (父ID null)
    • 处理器
      • ✉️ 消息 2 (父ID 1)
        • 处理器
      • ✉️ 消息 3 (父ID 1)
        • 处理器
    • 处理器
      • ✉️ 消息 4 (父ID 1)
        • 处理器
          • ✉️ 消息 5 (父ID 4)
            • 处理器

核心疑问:

  1. 如何在onSendMessageToTransportsEvent()方法中判断当前是否处于另一个消息处理器的调用过程中?
  2. Symfony Messenger是否有类似RequestStack的机制来追踪当前正在处理的消息?

解决方案:利用WorkerStack追踪当前处理消息

Symfony Messenger 内置了WorkerStack类,作用类似RequestStack,可以追踪当前正在处理的消息信封。

步骤1:注入WorkerStack到订阅者

修改订阅者构造函数,注入WorkerStack依赖:

use Symfony\Component\Messenger\Worker\WorkerStack;

final readonly class MessengerAuditSubscriber implements EventSubscriberInterface
{
    public function __construct(
        private WorkerStack $workerStack
    ) {}

    // ... 其余代码保持不变
}

步骤2:在派发消息时获取父消息ID

在onSendMessageToTransportsEvent()方法中,通过WorkerStack获取当前正在处理的消息,提取父ID并添加到新消息的Stamp中:

public function onSendMessageToTransportsEvent(SendMessageToTransportsEvent $event): void
{
    $envelope = $event->getEnvelope();
    
    // 生成当前消息的唯一ID
    $messageId = uniqid();
    $envelope = $envelope->with(new MessageIdStamp($messageId));
    
    // 获取父消息ID
    $parentMessageId = null;
    $currentProcessingEnvelope = $this->workerStack->getCurrentEnvelope();
    if ($currentProcessingEnvelope !== null) {
        $parentIdStamp = $currentProcessingEnvelope->last(MessageIdStamp::class);
        if ($parentIdStamp !== null) {
            $parentMessageId = $parentIdStamp->getId(); // 假设MessageIdStamp包含getId()方法
        }
    }
    
    // 添加父ID Stamp(需自定义该Stamp类)
    if ($parentMessageId !== null) {
        $envelope = $envelope->with(new ParentMessageIdStamp($parentMessageId));
    }
    
    $event->setEnvelope($envelope);
    
    // 记录日志,包含当前消息ID和父消息ID
    // ...
}

步骤3:自定义ParentMessageIdStamp类

创建用于存储父消息ID的Stamp:

use Symfony\Component\Messenger\Stamp\StampInterface;

final class ParentMessageIdStamp implements StampInterface
{
    public function __construct(
        private string $parentMessageId
    ) {}

    public function getParentMessageId(): string
    {
        return $this->parentMessageId;
    }
}

逻辑验证

  • 当在处理器内部派发新消息时,WorkerStack会返回当前正在处理的消息信封,此时能正确获取父消息ID;
  • 当从外部直接派发消息(如控制器中),WorkerStack返回null,父ID设为null,符合根消息的定义。

内容的提问来源于stack exchange,提问作者BenMorel

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.10 09:57:41