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(父IDnull)- 处理器
- ✉️ 消息
2(父ID1)- 处理器
- ✉️ 消息
3(父ID1)- 处理器
- ✉️ 消息
- 处理器
- ✉️ 消息
4(父ID1)- 处理器
- ✉️ 消息
5(父ID4)- 处理器
- ✉️ 消息
- 处理器
- ✉️ 消息
- 处理器
核心疑问:
- 如何在
onSendMessageToTransportsEvent()方法中判断当前是否处于另一个消息处理器的调用过程中? - 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
相关产品推荐
相关产品推荐

