基于Azure Service Bus实现Blazor应用与微服务API间通信:接收消息无法用于SQL查询的问题
问题分析与解决方案
你的核心问题是错误使用了Azure Service Bus的消息监听模式,同时代码还存在异步阻塞、SQL注入、线程安全等多个隐患,导致无法正确拿到筛选条件并用于SQL查询。下面我一步步拆解问题并给出修复方案:
1. 最关键的逻辑错配:Service Bus消息接收方式不对
你的ReceiveMsg方法是用来注册持续监听队列的消息处理器,而非一次性获取消息的方法。但你在同步的GetAll方法里用GetAwaiter().GetResult()阻塞调用它,这会导致:
- 死锁风险:同步阻塞异步代码是.NET环境下的常见坑,尤其在ASP.NET场景中
- 监听刚启动就被终止:
ReceiveMsg的finally里调用了await queueClient.CloseAsync(),会直接停止监听,根本等不到消息 - 逻辑矛盾:你期望调用
ReceiveMsg后立即拿到content1,但消息是异步到达的,注册处理器后不会立即返回结果
修复方案:重构为后台持续监听模式
我们需要把消息监听改成后台常驻服务,收到消息后把筛选条件缓存起来,API接口再读取缓存的条件进行查询:
第一步:在应用启动时初始化单例QueueClient
// 在Program.cs或Startup.cs的启动逻辑里初始化(注册为单例) public static QueueClient queueClient; // 用线程安全容器存储最新筛选条件 private static readonly ConcurrentDictionary<string, string[]> _filterCache = new ConcurrentDictionary<string, string[]>(); public static void InitializeServiceBusListener(string sbConnectionString, string sbQueueName) { queueClient = new QueueClient(sbConnectionString, sbQueueName); var messageHandlerOptions = new MessageHandlerOptions(ExceptionReceivedHandler) { MaxConcurrentCalls = 1, AutoComplete = false }; queueClient.RegisterMessageHandler(ReceiveMessagesAsync, messageHandlerOptions); Console.WriteLine("Service Bus 监听服务已启动"); } public static async Task ReceiveMessagesAsync(Message message, CancellationToken token) { try { var receivedMsg = Encoding.UTF8.GetString(message.Body); var filterMsg = JsonConvert.DeserializeObject<ServiceBusMessage>(receivedMsg); // 缓存最新的筛选条件,可根据业务场景调整key规则 _filterCache.AddOrUpdate("LatestFilter", filterMsg.Content, (key, oldVal) => filterMsg.Content); await queueClient.CompleteAsync(message.SystemProperties.LockToken); Debug.WriteLine($"已接收并缓存筛选条件:{string.Join(",", filterMsg.Content)}"); } catch (Exception ex) { Console.WriteLine($"处理消息失败:{ex.Message}"); await queueClient.AbandonAsync(message.SystemProperties.LockToken); } } static Task ExceptionReceivedHandler(ExceptionReceivedEventArgs exceptionReceivedEventArgs) { Console.WriteLine($"Service Bus 异常:{exceptionReceivedEventArgs.Exception}"); return Task.CompletedTask; }
2. 修复API接口的查询逻辑
现在GetAll方法只需读取缓存的筛选条件,不用再调用消息接收方法:
public class UserController : ControllerBase { public IEnumerable<dynamic> Get() { // 读取最新筛选条件,无条件时返回全量数据 if (!_filterCache.TryGetValue("LatestFilter", out var filterContent)) { return userRepository.GetAllUsers(); } var startDate = filterContent[0]; var endDate = filterContent[1]; return userRepository.GetUsersByBirthDateRange(startDate, endDate); } }
3. 修复严重的SQL注入风险
原始代码直接拼接SQL字符串,这是高危安全漏洞,必须改用参数化查询(Dapper原生支持):
public IEnumerable<dynamic> GetUsersByBirthDateRange(string startDate, string endDate) { using (IDbConnection dbConnection = connection) { var result = dbConnection.Query(@" select * from [User] where DateofBirth between @StartDate and @EndDate", new { StartDate = startDate, EndDate = endDate }); return result; } } public IEnumerable<dynamic> GetAllUsers() { using (IDbConnection dbConnection = connection) { var result = dbConnection.Query("select * from [User]"); return result; } }
4. 其他优化点
- 移除
ReceiveMsg().GetAwaiter().GetResult()这种同步阻塞异步代码的写法,彻底避免死锁 - 用
ConcurrentDictionary替代静态变量content1,解决多请求下的线程安全问题 - 确保
queueClient是单例生命周期,避免重复创建和销毁资源 - 如果业务场景是「前端请求时需要等待Service Bus的筛选条件」,可以考虑改用Azure Service Bus的请求-响应模式,而非单向队列
内容的提问来源于stack exchange,提问作者Sameer Gaonkar
相关产品推荐
相关产品推荐

