迁移至Azure.Messaging.ServiceBus后,ServiceBusTrigger无法绑定ServiceBusReceiver
问题分析与解决方案
核心问题
从旧版Microsoft.Azure.ServiceBus迁移到Azure.Messaging.ServiceBus后,遇到两个关键问题:
- 新库的Service Bus触发器不支持直接注入
ServiceBusReceiver作为参数,旧版MessageReceiver的绑定方式失效。 - 手动创建的
ServiceBusReceiver无法操作函数已锁定的消息,触发MessageLockLost错误。 - 核心诉求:精准控制消息的完成/死信逻辑,确保跨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
相关产品推荐
相关产品推荐

