NServiceBus:如何让父Saga捕获子Saga处理的事件
解决方案
针对你遇到的Saga_X需要捕获Saga_Y触发事件的场景,以下几种方案可以避免重新发布事件或手动传递ParentCorrelationId的繁琐操作:
1. 配置Saga_X的事件关联规则
核心思路是利用Saga框架的关联映射机制,让消息总线自动通过事件中的标识字段匹配到对应的Saga_X实例。
具体步骤:
- 在Saga_X发送创建Saga_Y的命令时,将自身的CorrelationId(比如业务场景中的
OrderId)携带在命令中 - Saga_Y初始化时,将这个
OrderId保存到自身的状态数据中 - Saga_Y发布目标事件时,把
OrderId作为事件字段一并发出 - 在Saga_X中,为目标事件配置关联映射,指定事件的
OrderId对应Saga_X状态中的OrderId
示例代码(以NServiceBus为例):
// Saga_X的定义 public class Saga_X : Saga<Saga_XState>, IAmStartedByMessages<CreateOrderItemsCommand>, IHandleMessages<OrderItemProcessed> { // 处理创建Saga_Y的命令 public async Task Handle(CreateOrderItemsCommand message, IMessageHandlerContext context) { // 保存自身CorrelationId到状态 Data.OrderId = message.OrderId; // 发送创建多个Saga_Y的命令,携带OrderId foreach (var itemId in message.ItemIds) { await context.Send(new CreateOrderItemCommand { OrderId = message.OrderId, ItemId = itemId }).ConfigureAwait(false); } } // 配置事件关联规则 protected override void ConfigureHowToFindSaga(SagaPropertyMapper<Saga_XState> mapper) { mapper.ConfigureMapping<OrderItemProcessed>(msg => msg.OrderId) .ToSaga(saga => saga.OrderId); } // 处理Saga_Y发布的事件 public async Task Handle(OrderItemProcessed message, IMessageHandlerContext context) { // 处理逻辑 Data.ProcessedItemIds.Add(message.ItemId); if (Data.ProcessedItemIds.Count == Data.TotalItemCount) { MarkAsComplete(); } } } // Saga_Y发布事件的逻辑 public class Saga_Y : Saga<Saga_YState>, IAmStartedByMessages<CreateOrderItemCommand>, IHandleMessages<ItemProcessedEvent> { public async Task Handle(ItemProcessedEvent message, IMessageHandlerContext context) { // 发布事件,携带从命令中获取的OrderId await context.Publish(new OrderItemProcessed { OrderId = Data.OrderId, ItemId = Data.ItemId }).ConfigureAwait(false); MarkAsComplete(); } }
这种方式下,消息总线会自动根据OrderId匹配到对应的Saga_X实例,无需手动转发事件。
2. 使用消息订阅过滤规则
如果你的消息传输层支持订阅过滤(比如NServiceBus的SQL Transport、RabbitMQ Transport),可以让Saga_X在订阅目标事件时,只接收包含自身CorrelationId的事件。
具体操作:
- Saga_X在启动或发送创建Saga_Y的命令后,调用订阅API并添加过滤条件,例如只接收
OrderId = '{当前Saga_X的CorrelationId}'的事件 - Saga_Y发布事件时携带
OrderId,总线会自动将事件路由到符合过滤条件的Saga_X实例
示例代码(NServiceBus RabbitMQ Transport为例):
// 在Saga_X中订阅事件并添加过滤 public async Task Handle(CreateOrderItemsCommand message, IMessageHandlerContext context) { Data.OrderId = message.OrderId; // 订阅OrderItemProcessed事件,只接收当前OrderId的消息 await context.Subscribe<OrderItemProcessed>( filter: $"OrderId = '{message.OrderId}'").ConfigureAwait(false); // 发送创建Saga_Y的命令... }
这种方式适合不需要全局订阅事件,仅特定Saga实例需要接收对应事件的场景。
3. 直接发送事件到Saga_X实例
在Saga_Y处理完业务逻辑后,不使用Publish发布事件,而是直接Send事件到Saga_X所在的终结点,利用消息的关联标识让总线找到目标Saga_X实例。
示例代码:
// 在Saga_Y的事件处理方法中 public async Task Handle(ItemProcessedEvent message, IMessageHandlerContext context) { var processedEvent = new OrderItemProcessed { OrderId = Data.OrderId, ItemId = Data.ItemId }; // 直接发送到处理OrderItemProcessed的终结点(即Saga_X所在终结点) await context.Send(processedEvent).ConfigureAwait(false); MarkAsComplete(); }
结合Saga_X中配置的关联规则,总线会自动将消息路由到对应的Saga_X实例,避免了事件的全局发布。
内容的提问来源于stack exchange,提问作者user1696503
相关产品推荐
相关产品推荐

