无需新增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
相关产品推荐
相关产品推荐

