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

基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.30 05:59:08