Azure Service Bus同主题下不同消息的差异化处理方案咨询
同一Azure Service Bus主题处理创建/更新消息的简洁实现方案
不用新增主题,有两种优雅的实现方式,核心都是通过消息属性区分操作类型,避免解析payload判断的冗余逻辑:
方案一:单ServiceBusTrigger接收,基于消息属性分发处理
这种方式无需修改ServiceBus的订阅配置,在代码里统一接收消息后根据属性路由到不同处理逻辑:
步骤1:发送消息时添加操作类型属性
修改Inbound项目的两个HttpTrigger,在发送消息时给Service Bus消息添加OperationType自定义属性:
更新消息的Trigger修改:
// Inbound proj - PostItemUpdateHttpTrigger public async Task<HttpResponseData> Run([HttpTrigger(AuthorizationLevel.Anonymous, "post", Route = FunctionRoute)] HttpRequestData req) { var itemUpdateDto = await req.ReadFromJsonAsync<ItemUpdateDto>() ?? throw new Exception("Error deserialising data."); // 发送时添加操作类型属性 await eventBusService.SendAsync(itemUpdateDto, Topics.Items, new Dictionary<string, object> { ["OperationType"] = "Update" }); return req.CreateResponse(HttpStatusCode.OK); }
新增创建消息的Trigger示例:
// Inbound proj - PostItemCreateHttpTrigger public class PostItemCreateHttpTrigger(IEventBusService eventBusService) { public const string FunctionRoute = "ItemCreate"; [Function(FunctionRoute)] [OpenApiOperation(operationId: FunctionRoute, tags: ["Items"])] [OpenApiRequestBody(nameof(ItemCreateDto), typeof(ItemCreateDto))] [OpenApiResponseWithoutBody(statusCode: HttpStatusCode.OK, Description = "Publishes create item data to the 'items' service bus topic.")] public async Task<HttpResponseData> Run([HttpTrigger(AuthorizationLevel.Anonymous, "post", Route = FunctionRoute)] HttpRequestData req) { var itemCreateDto = await req.ReadFromJsonAsync<ItemCreateDto>() ?? throw new Exception("Error deserialising data."); // 发送时添加操作类型属性 await eventBusService.SendAsync(itemCreateDto, Topics.Items, new Dictionary<string, object> { ["OperationType"] = "Create" }); return req.CreateResponse(HttpStatusCode.OK); } }
扩展EventBusService支持自定义属性:
// EventBusService 新增重载方法 public async Task SendAsync<T>(T message, string topicName, IDictionary<string, object>? properties = null) { using var sender = _serviceBusClient.CreateSender(topicName); var jsonPayload = JsonConvert.SerializeObject(message); var serviceBusMessage = new ServiceBusMessage(jsonPayload); // 添加自定义属性 if (properties != null) { foreach (var prop in properties) { serviceBusMessage.ApplicationProperties[prop.Key] = prop.Value; } } await sender.SendMessageAsync(serviceBusMessage); }
步骤2:修改BackEnd的Trigger统一处理并分发
把原来的Trigger参数从string改为ServiceBusReceivedMessage,以便读取消息属性:
// BackEnd proj - 统一处理的ServiceBusTrigger public class ItemServiceBusTrigger(IMediator mediator, ServiceBusReceiver messageReceiver) { [Function(nameof(ItemServiceBusTrigger))] public async Task Run([ServiceBusTrigger( "Items", "%ServiceBusSettings:SubscriptionName%", Connection = "ServiceBusSettings:ConnectionString")] ServiceBusReceivedMessage message) { // 读取操作类型属性 if (!message.ApplicationProperties.TryGetValue("OperationType", out var opTypeObj) || !(opTypeObj is string operationType)) { // 无效消息直接死信 await messageReceiver.DeadLetterMessageAsync(message, "Missing or invalid OperationType property"); return; } var payload = message.Body.ToString(); switch (operationType) { case "Update": var updateDto = JsonConvert.DeserializeObject<ItemUpdateDto>(payload); await mediator.Send(UpdateCommand.FromItemUpdateDto(updateDto!)); break; case "Create": var createDto = JsonConvert.DeserializeObject<ItemCreateDto>(payload); await mediator.Send(CreateCommand.FromItemCreateDto(createDto!)); break; default: await messageReceiver.DeadLetterMessageAsync(message, $"Unknown OperationType: {operationType}"); break; } } }
方案二:创建带过滤规则的独立订阅(推荐)
利用Azure Service Bus的订阅过滤功能,在服务端直接把消息分流到不同订阅,每个订阅对应独立的Trigger,职责更单一:
步骤1:创建两个订阅并添加过滤规则
在Azure Portal或通过Azure CLI给Items主题创建两个订阅:
- 订阅名:
Items-Update,添加SQL过滤规则:OperationType = 'Update' - 订阅名:
Items-Create,添加SQL过滤规则:OperationType = 'Create'
步骤2:创建两个独立的ServiceBusTrigger
分别处理更新和创建消息,逻辑和原来的Trigger一致,只是订阅名不同:
处理更新的Trigger:
// BackEnd proj - ItemUpdateServiceBusTrigger public class ItemUpdateServiceBusTrigger(IMediator mediator) { [Function(nameof(ItemUpdateServiceBusTrigger))] public async Task Run([ServiceBusTrigger( "Items", "Items-Update", Connection = "ServiceBusSettings:ConnectionString")] string inputEvent) { var updateDto = JsonConvert.DeserializeObject<ItemUpdateDto>(inputEvent); await mediator.Send(UpdateCommand.FromItemUpdateDto(updateDto!)); } }
处理创建的Trigger:
// BackEnd proj - ItemCreateServiceBusTrigger public class ItemCreateServiceBusTrigger(IMediator mediator) { [Function(nameof(ItemCreateServiceBusTrigger))] public async Task Run([ServiceBusTrigger( "Items", "Items-Create", Connection = "ServiceBusSettings:ConnectionString")] string inputEvent) { var createDto = JsonConvert.DeserializeObject<ItemCreateDto>(inputEvent); await mediator.Send(CreateCommand.FromItemCreateDto(createDto!)); } }
方案对比
- 方案一:适合小型项目或逻辑简单的场景,代码集中管理,但单Trigger可能成为高并发下的瓶颈。
- 方案二:符合Service Bus的设计理念,消息在服务端分流,每个Trigger职责单一,扩展性更好(后续新增操作只需加新订阅和过滤规则),是更推荐的生产级方案。
关键注意点
- 必须通过**消息属性(而非payload)**传递操作类型,这样既能避免解析payload的性能损耗,也能支持Service Bus的订阅过滤功能。
- 对于无效消息(比如缺少操作类型),一定要死信处理,避免消息循环重试。
内容的提问来源于stack exchange,提问作者marcuthh
相关产品推荐
相关产品推荐

