Azure Service Bus消息重复入队问题排查求助
问题描述
作为Azure Service Bus新手,开发了一个简单的Azure Function,通过EventBusService向Service Bus队列发送消息。调试确认端点仅被触发一次,但每次运行测试后,队列中都会出现两条完全相同的消息。
相关代码如下:
Azure Function 端点代码
public class TestFunction(IEventBusService eventBusService) { [Function("TestFunction")] public async Task<HttpResponseData> Run([HttpTrigger(AuthorizationLevel.Function, "get")] HttpRequestData req) { try { await eventBusService.SendAsync(new TestObject { Message = "Message for service bus", SentAtDateTimeOffset = DateTimeOffset.UtcNow }, Queues.Test); var response = req.CreateResponse(HttpStatusCode.OK); await response.WriteAsJsonAsync("Message successfully published to service bus", "application/json"); return response; } catch (ServiceBusOperationFailedException) { var response = req.CreateResponse(HttpStatusCode.InternalServerError); await response.WriteAsJsonAsync("Unable to publish message to service bus", "application/json"); return response; } } } public class TestObject { public string Message { get; set; } public DateTimeOffset SentAtDateTimeOffset { get; set; } }
消息发送服务实现
public class EventBusService(ServiceBusClient serviceBusClient) : IEventBusService { public async Task SendAsync<T>(T dataToSend, string sendTo, CancellationToken cancellationToken = default) { ValidateParams(dataToSend, sendTo); var serviceBusSender = serviceBusClient.CreateSender(sendTo); var messageJson = new ServiceBusMessage(JsonSerializer.Serialize(dataToSend)); try { await serviceBusSender.SendMessageAsync(messageJson, cancellationToken); } catch (Exception ex) { throw new ServiceBusOperationFailedException($"Unable to publish message to '{sendTo}'.", ex); } finally { await serviceBusSender.DisposeAsync(); } } private void ValidateParams<T>(T dataToSend, string sendTo) { if (dataToSend is null) throw new ArgumentNullException(nameof(dataToSend)); if (string.IsNullOrEmpty(sendTo)) throw new ArgumentNullException(nameof(sendTo)); } }
测试代码
[TestMethod] public async Task Should_OnSuccessfulPublishToEventBus_ReceiveConfirmationMessage() { // Arrange var testFunction = new TestFunction(EventBusHelper.CreateService()); var mockFunctionContext = EventBusHelper.GetMockFunctionContext(); //Act var response = await testFunction.Run(new TestHttpRequestData(mockFunctionContext)); //Assert Assert.IsNotNull(response); Assert.AreEqual(HttpStatusCode.OK, response.StatusCode); response.Body.Position = 0; using var streamReader = new StreamReader(response.Body); var responseBody = await streamReader.ReadToEndAsync(); Assert.AreEqual("\"Message successfully published to service bus\"", responseBody); }
排查思路
检查ServiceBusClient实例管理
ServiceBusClient是线程安全的,需以单例模式复用。如果EventBusHelper.CreateService()每次都创建新的ServiceBusClient实例,可能引发底层连接问题导致重复发送。验证测试代码执行逻辑
确认测试框架是否存在重复执行机制,或TestHttpRequestData/EventBusHelper的实现是否意外触发多次发送。排查SDK默认重试策略
Azure Service Bus SDK默认启用指数退避重试,若第一次发送因网络延迟超时但服务端已接收消息,重试会导致重复发送。检查ServiceBusSender使用方式
频繁创建并销毁Sender可能引入不稳定因素,虽安全但非最优实践,建议复用Sender实例(Sender同样线程安全)。查看服务端诊断日志
在Azure Portal的Service Bus资源中查看诊断日志,确认服务端是否收到两次发送请求,区分是客户端重复发送还是服务端异常导致的重复。
解决方案
确保ServiceBusClient单例注册
在依赖注入容器中将ServiceBusClient注册为单例:builder.Services.AddSingleton<ServiceBusClient>(sp => new ServiceBusClient(Environment.GetEnvironmentVariable("ServiceBusConnectionString")));启用消息重复检测
- 给每个消息设置唯一
MessageId:var messageJson = new ServiceBusMessage(JsonSerializer.Serialize(dataToSend)) { MessageId = Guid.NewGuid().ToString() }; - 在Azure Portal中开启队列的重复检测功能,设置合理的重复检测窗口(如5分钟),Service Bus会自动丢弃相同
MessageId的重复消息。
- 给每个消息设置唯一
调整SDK重试策略
根据业务需求限制或禁用重试:var serviceBusClient = new ServiceBusClient(connectionString, new ServiceBusClientOptions { RetryOptions = new ServiceBusRetryOptions { MaxRetries = 0 // 禁用重试,或设置合理次数 } });复用ServiceBusSender实例
修改EventBusService缓存Sender实例:public class EventBusService(ServiceBusClient serviceBusClient) : IEventBusService, IAsyncDisposable { private readonly Dictionary<string, ServiceBusSender> _senders = new(); private readonly object _lock = new(); public async Task SendAsync<T>(T dataToSend, string sendTo, CancellationToken cancellationToken = default) { ValidateParams(dataToSend, sendTo); ServiceBusSender sender; lock (_lock) { if (!_senders.TryGetValue(sendTo, out sender)) { sender = serviceBusClient.CreateSender(sendTo); _senders.Add(sendTo, sender); } } var messageJson = new ServiceBusMessage(JsonSerializer.Serialize(dataToSend)) { MessageId = Guid.NewGuid().ToString() }; try { await sender.SendMessageAsync(messageJson, cancellationToken); } catch (Exception ex) { throw new ServiceBusOperationFailedException($"Unable to publish message to '{sendTo}'.", ex); } } public async ValueTask DisposeAsync() { foreach (var sender in _senders.Values) { await sender.DisposeAsync(); } } private void ValidateParams<T>(T dataToSend, string sendTo) { if (dataToSend is null) throw new ArgumentNullException(nameof(dataToSend)); if (string.IsNullOrEmpty(sendTo)) throw new ArgumentNullException(nameof(sendTo)); } }
内容的提问来源于stack exchange,提问作者marcuthh

