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

迁移至Azure.Messaging.ServiceBus后,ServiceBusTrigger无法绑定ServiceBusReceiver

问题分析与解决方案

核心问题

从旧版Microsoft.Azure.ServiceBus迁移到Azure.Messaging.ServiceBus后,遇到两个关键问题:

  1. 新库的Service Bus触发器不支持直接注入ServiceBusReceiver作为参数,旧版MessageReceiver的绑定方式失效。
  2. 手动创建的ServiceBusReceiver无法操作函数已锁定的消息,触发MessageLockLost错误。
  3. 核心诉求:精准控制消息的完成/死信逻辑,确保跨Service Bus的消息仅发送一次,避免重复执行。

解决方案

1. 使用ServiceBusMessageActions替代ServiceBusReceiver

新库的Service Bus触发器提供了ServiceBusMessageActions类型的绑定参数,它与当前触发器持有的消息锁上下文绑定,可安全完成/死信当前消息,完全替代旧版MessageReceiver的功能。

修改函数参数

将原来的ServiceBusReceiver替换为ServiceBusMessageActions:

[FunctionName("Test")]
public async Task Run(
    [ServiceBusTrigger("%InternalInformation:Topic%", "%InternalInformation:Subscription%",
        Connection = "InternalInformation:ListenKey")]
    ServiceBusReceivedMessage sbMsg,
    ServiceBusMessageActions messageActions, // 替换原ServiceBusReceiver参数
    ExecutionContext context,
    ILogger log)

使用messageActions操作消息

在完成或死信消息的逻辑中,直接调用messageActions的对应方法,无需手动创建ServiceBusReceiver:

try
{
    DoSomething2(sbMsg);
    // 标记消息为完成,从订阅中移除
    await messageActions.CompleteMessageAsync(sbMsg);
}
catch (Exception e)
{
    // 将消息移入死信队列,避免重试导致重复发送跨Service Bus消息
    await messageActions.DeadLetterMessageAsync(sbMsg);
    throw e;
}

2. 优化ServiceBusClient的实例化方式

当前代码每次函数调用都创建新的ServiceBusClient,会造成性能损耗和连接资源浪费。ServiceBusClient是线程安全的,应注册为单例依赖注入:

注册单例客户端

在Function的启动类(或Program.cs,根据Function版本)中添加:

builder.Services.AddSingleton(new ServiceBusClient(builder.Configuration["InternalInformation:SendKey"]));

构造函数注入客户端

在myClass中注入单例ServiceBusClient:

private readonly ServiceBusClient _sendClient;

public myClass(IOptions<AppConfig.InternalInformation> info, IOptions<AppConfig.EventInformation> info2,
    IConfiguration configuration, ServiceBusClient sendClient)
{
    this._configuration = configuration;
    this._info = info.Value;
    this._info2 = info2.Value;
    this._sendClient = sendClient;
}

使用注入的客户端发送消息

ServiceBusMessage newMsg = new ServiceBusMessage(Encoding.UTF8.GetBytes(someMsg))
{
    // 设置MessageId,方便目标Service Bus接收端做幂等校验
    MessageId = sbMsg.MessageId
};

ServiceBusSender sender = _sendClient.CreateSender(_info2.Topic);
await sender.SendMessageAsync(newMsg);
// 无需手动Close/Dispose,单例客户端会自动管理连接

3. 保障跨Service Bus消息仅发送一次

即使函数执行到步骤2后失败,重试时仍可能重复发送消息,最可靠的解决方案是在目标Service Bus的接收端实现幂等校验:

  • 使用消息的MessageId或自定义业务唯一标识作为幂等键
  • 接收端处理消息前,先检查该标识是否已处理过,避免重复执行

修正后的完整代码

using System;
using System.Text;
using System.Threading.Tasks;
using Microsoft.Azure.WebJobs;
using Microsoft.Extensions.Configuration;
using Microsoft.Extensions.Logging;
using myNameSpace.Models;
using Microsoft.Extensions.Options;
using Azure.Messaging.ServiceBus;

namespace myNameSpace
{
    public class myClass
    {
        private readonly AppConfig.InternalInformation _info;
        private readonly AppConfig.EventInformation _info2;
        private readonly IConfiguration _configuration;
        private readonly ServiceBusClient _sendClient;

        public myClass(IOptions<AppConfig.InternalInformation> info, IOptions<AppConfig.EventInformation> info2,
            IConfiguration configuration, ServiceBusClient sendClient)
        {
            this._configuration = configuration;
            this._info = info.Value;
            this._info2 = info2.Value;
            this._sendClient = sendClient;
        }

        [FunctionName("Test")]
        public async Task Run(
            [ServiceBusTrigger("%InternalInformation:Topic%", "%InternalInformation:Subscription%",
                Connection = "InternalInformation:ListenKey")]
            ServiceBusReceivedMessage sbMsg,
            ServiceBusMessageActions messageActions,
            ExecutionContext context,
            ILogger log)
        {
            // 1. 执行第一步业务逻辑
            string someMsg = DoSomething(sbMsg);

            // 2. 发送消息到目标Service Bus
            ServiceBusMessage newMsg = new ServiceBusMessage(Encoding.UTF8.GetBytes(someMsg))
            {
                MessageId = sbMsg.MessageId // 复用原消息ID,用于目标端幂等校验
            };
            
            ServiceBusSender sender = _sendClient.CreateSender(_info2.Topic);
            await sender.SendMessageAsync(newMsg);

            // 3. 执行第二步业务逻辑
            try
            {
                DoSomething2(sbMsg);
                // 完成当前消息
                await messageActions.CompleteMessageAsync(sbMsg);
            }
            catch (Exception e)
            {
                // 死信当前消息,避免重试重复发送跨Service Bus消息
                await messageActions.DeadLetterMessageAsync(sbMsg);
                throw e;
            }
        }
        
        private string DoSomething(ServiceBusReceivedMessage sbMsg)
        {
            // ... 业务逻辑实现
            return "Something";
        }
        
        private void DoSomething2(ServiceBusReceivedMessage sbMsg)
        {
            // ... 业务逻辑实现
        }
    }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 08:21:57