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

无需新增Azure Function,如何在现有函数内发送SignalR消息?

问题描述

我是一名C#学习者,现有一套基于.NET 6的Azure Function架构:当收到HTTP触发请求后,会触发编排器函数,再通过活动触发器启动QueueStore函数,将Azure存储队列中的所有消息复制到Cosmos DB。
我希望在QueueStore函数的指定代码位置(注释处),通过Azure中已有的SignalR服务和Hub,将每条消息发送给客户端。
目前关于创建SignalR和协商函数的文档很多,但我不清楚如何在已调用的函数内部直接发送SignalR消息。我尝试过从编排器调用SignalR函数,但这会新增外部函数且需重复调用队列,并不适用;也试过标准.NET SignalR代码,但未找到可行示例。
请问该如何实现此需求?是否需要新建独立函数应用并通过HTTP调用?

原QueueStore函数代码

[FunctionName(nameof(QueueStore))]
public static async Task<string> QueueStore([ActivityTrigger] QueueName queue, ILogger log)
{
// Get the connection string
string connectionString = Environment.GetEnvironmentVariable("QueueStorage");

try
{
    CosmosClient client = new CosmosClient("some info here");
    Database database = client.GetDatabase("database");
    bool databaseExists = true;
    try
    {
        var response = await database.ReadAsync();
    }
    catch (CosmosException ex)
    {
        if (ex.StatusCode.Equals(HttpStatusCode.NotFound))
        {
            // Does not exist
            databaseExists = false;
        }
    }

    //Instantiate a QueueClient which will be used to manipulate the queue
    QueueClient queueClient = new QueueClient(connectionString, queue.Name);
    QueueProperties properties = await queueClient.GetPropertiesAsync();

    bool appDisconnected = false;

    //string message = "Stored messages ";
    if (queueClient.Exists() && databaseExists)
    {
                
        Container container = await database.CreateContainerIfNotExistsAsync(id: queue.Name,
            partitionKeyPath: "/partKey", //name of the json var we want as partition key
            throughput: 400
        );

        while (appDisconnected == false)
        {
            if (queueClient.GetProperties().Value.ApproximateMessagesCount == 0)
            {
                Thread.Sleep(100);
            }
            else
            {
                QueueMessage[] retrievedMessage = await queueClient.ReceiveMessagesAsync(1);

                var fd = JsonConvert.DeserializeObject<JObject>(retrievedMessage[0].Body.ToString());

                if (!fd.ContainsKey("disconnected"))
                {


                    PartitionKey partKey = new PartitionKey(queue.PartKey);
                    // save to db

                    var createdItem = await container.CreateItemAsync<JObject>(
                        item: fd,
                        partitionKey: partKey);

                    //######## HERE IS WHERE I WANT TO SEND THE fd Object via SignalR
                            
            //######## I have tried many different things but nothing works                          
                    
                    await queueClient.DeleteMessageAsync(retrievedMessage[0].MessageId,
                        retrievedMessage[0].PopReceipt);
                }
                else
                {
                    appDisconnected = true;
                }
            }
        }
        return "Copied all Items";
    }
    else
    {
        return $"The queue peek didn't work because I can't find the queue:-(";
    }
}
catch (Exception ex)
{
    return ex.Message;
}

解决方案

不需要新建独立函数应用,直接在QueueStore活动函数中使用Azure SignalR Service的SDK即可发送消息,具体步骤如下:

1. 安装NuGet包

安装Azure SignalR Service的官方扩展包:

Install-Package Microsoft.Azure.WebJobs.Extensions.SignalRService

或通过.NET CLI:

dotnet add package Microsoft.Azure.WebJobs.Extensions.SignalRService

2. 配置SignalR连接字符串

在Function应用的应用设置中添加Azure SignalR Service的连接字符串,键名设为AzureSignalRConnectionString(默认配置键,也可自定义)。

3. 修改QueueStore函数代码

将静态函数改为实例函数,通过依赖注入获取ISignalRServiceClient,然后在指定位置发送消息:

修改后的完整代码

using Microsoft.Azure.WebJobs;
using Microsoft.Azure.WebJobs.Extensions.SignalRService;
using Microsoft.Extensions.Logging;
using Azure.Storage.Queues;
using Azure.Storage.Queues.Models;
using Microsoft.Azure.Cosmos;
using Newtonsoft.Json;
using Newtonsoft.Json.Linq;
using System;
using System.Net;
using System.Threading.Tasks;

public class QueueStoreFunctions
{
    private readonly ISignalRServiceClient _signalRClient;

    // 构造函数注入SignalR客户端
    public QueueStoreFunctions(ISignalRServiceClient signalRClient)
    {
        _signalRClient = signalRClient;
    }

    [FunctionName(nameof(QueueStore))]
    public async Task<string> QueueStore([ActivityTrigger] QueueName queue, ILogger log)
    {
        string connectionString = Environment.GetEnvironmentVariable("QueueStorage");

        try
        {
            CosmosClient client = new CosmosClient("some info here");
            Database database = client.GetDatabase("database");
            bool databaseExists = true;
            try
            {
                var response = await database.ReadAsync();
            }
            catch (CosmosException ex)
            {
                if (ex.StatusCode.Equals(HttpStatusCode.NotFound))
                {
                    databaseExists = false;
                }
            }

            QueueClient queueClient = new QueueClient(connectionString, queue.Name);
            await queueClient.GetPropertiesAsync();

            bool appDisconnected = false;

            if (queueClient.Exists() && databaseExists)
            {
                Container container = await database.CreateContainerIfNotExistsAsync(
                    id: queue.Name,
                    partitionKeyPath: "/partKey",
                    throughput: 400
                );

                while (!appDisconnected)
                {
                    if (queueClient.GetProperties().Value.ApproximateMessagesCount == 0)
                    {
                        await Task.Delay(100); // 用Task.Delay替代Thread.Sleep,避免阻塞线程池
                    }
                    else
                    {
                        QueueMessage[] retrievedMessage = await queueClient.ReceiveMessagesAsync(1);
                        if (retrievedMessage.Length == 0) continue;

                        var fd = JsonConvert.DeserializeObject<JObject>(retrievedMessage[0].Body.ToString());

                        if (!fd.ContainsKey("disconnected"))
                        {
                            PartitionKey partKey = new PartitionKey(queue.PartKey);
                            var createdItem = await container.CreateItemAsync<JObject>(
                                item: fd,
                                partitionKey: partKey);

                            // ######## 通过SignalR发送消息到客户端
                            // 发送给所有连接的客户端
                            await _signalRClient.SendToAllAsync("ReceiveMessage", fd);
                            // 若需定向发送,可使用以下方法:
                            // await _signalRClient.SendToUserAsync("target-user-id", "ReceiveMessage", fd);
                            // await _signalRClient.SendToGroupAsync("target-group-name", "ReceiveMessage", fd);

                            await queueClient.DeleteMessageAsync(retrievedMessage[0].MessageId,
                                retrievedMessage[0].PopReceipt);
                        }
                        else
                        {
                            appDisconnected = true;
                        }
                    }
                }
                return "Copied all Items";
            }
            else
            {
                return $"The queue peek didn't work because I can't find the queue:-(";
            }
        }
        catch (Exception ex)
        {
            return ex.Message;
        }
    }
}

// 自定义QueueName模型
public class QueueName
{
    public string Name { get; set; }
    public string PartKey { get; set; }
}

关键说明

  • 依赖注入:将静态函数改为实例函数,通过构造函数注入ISignalRServiceClient,这是Azure Functions中使用SignalR的标准方式。
  • 消息发送:SendToAllAsync可向所有连接Hub的客户端推送消息;如果需要定向发送,可使用SendToUserAsync(指定用户ID)或SendToGroupAsync(指定组名)。
  • 线程优化:用Task.Delay替代Thread.Sleep,避免阻塞线程池线程,符合异步编程最佳实践。
  • 配置验证:确保AzureSignalRConnectionString已正确配置,值为你的Azure SignalR Service实例的连接字符串。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 14:41:44