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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 04:57:33