.NET 8隔离模式Azure Functions如何从单函数发送多SignalR异步消息
问题描述
将Azure Functions从.NET 6迁移至.NET 8隔离工作者模式后,原逻辑中通过SignalR向前端发送长耗时队列任务进度消息的功能失效。在.NET 8隔离模式下,SignalR输出绑定仅能返回单个SignalRMessageAction,无法在函数执行过程中发送多条消息。尝试通过单独的SendMessage函数发送消息但未成功,希望找到按顺序发送多条SignalR消息的可行方案(比如通过独立队列或调用其他SignalR消息函数)。
原函数代码:
[Function("PerformBackgroundJob")] [SignalROutput(HubName = "progress", ConnectionStringSetting = "AzureSignalRConnectionString")] public SignalRMessageAction PerformBackgroundJob( [QueueTrigger(LibraryConstants.testQueue, Connection = "connectionstring-azure-storage")] string JobId) { var msg = new SignalRMsg(); _logger.LogInformation("PerformBackgroundJob started"); Stopwatch sw1 = new Stopwatch(); sw1.Start(); List<ProcessFileResponse> response = new List<ProcessFileResponse>(); msg = new SignalRMsg() { UserId = JobId, Target = "taskStarted", Arguments = new object[] { "PerformBackgroundJob started" } }; //Send a message here //SendMessage(msg); _logger.LogInformation("PerformBackgroundJob data reviewed"); for (int i = 0; i < 100; i++) { msg = new SignalRMsg() { UserId = JobId, Target = "taskProgressChanged", Arguments = new object[] { i + 1 } }; //Send a message here //SendMessage(msg); Thread.Sleep(200); } sw1.Stop(); TimeSpan ts = sw1.Elapsed; // Format and display the TimeSpan value. string elapsedTime = String.Format("{0:00}:{1:00}:{2:00}.{3:00}", ts.Hours, ts.Minutes, ts.Seconds, ts.Milliseconds / 10); ProcessFileResponse outResponse = new ProcessFileResponse(); outResponse.Errors = 0; outResponse.Warnings = 0; outResponse.Status = 1; outResponse.TabName = "Background action"; outResponse.Message = $"RunTime: {elapsedTime}"; response.Add(outResponse); _logger.LogInformation("PerformBackgroundJob completed"); //this is The only message that is visible return new SignalRMessageAction("taskEnded") { Arguments = new object[] { response }, UserId = JobId }; }
尝试的SendMessage函数(未生效):
[Function(nameof(SendMessage))] [SignalROutput(HubName = "progress", ConnectionStringSetting = "AzureSignalRConnectionString")] public static SignalRMessageAction SendMessage( [HttpTrigger(AuthorizationLevel.Anonymous, "post")] SignalRMsg req) { return new SignalRMessageAction(req.Target) { Arguments = req.Arguments, UserId = req.UserId }; } public class SignalRMsg { public string UserId { get; set; } public string Target { get; set; } public object[] Arguments { get; set; } }
解决方案
方法1:使用SignalR输出绑定返回消息集合
.NET 8隔离模式的SignalR输出绑定支持返回IEnumerable<SignalRMessageAction>类型,可一次性返回多条消息。注意这种方式是批量发送,所有消息会在函数执行完成后统一推送,无法实时反映任务进度,适合提前确定所有消息内容的场景。
修改后的代码示例:
[Function("PerformBackgroundJob")] [SignalROutput(HubName = "progress", ConnectionStringSetting = "AzureSignalRConnectionString")] public IEnumerable<SignalRMessageAction> PerformBackgroundJob( [QueueTrigger(LibraryConstants.testQueue, Connection = "connectionstring-azure-storage")] string JobId) { var messages = new List<SignalRMessageAction>(); _logger.LogInformation("PerformBackgroundJob started"); Stopwatch sw1 = new Stopwatch(); sw1.Start(); // 添加任务开始消息 messages.Add(new SignalRMessageAction("taskStarted") { Arguments = new object[] { "PerformBackgroundJob started" }, UserId = JobId }); _logger.LogInformation("PerformBackgroundJob data reviewed"); // 添加所有进度消息 for (int i = 0; i < 100; i++) { messages.Add(new SignalRMessageAction("taskProgressChanged") { Arguments = new object[] { i + 1 }, UserId = JobId }); Thread.Sleep(200); } sw1.Stop(); TimeSpan ts = sw1.Elapsed; string elapsedTime = String.Format("{0:00}:{1:00}:{2:00}.{3:00}", ts.Hours, ts.Minutes, ts.Seconds, ts.Milliseconds / 10); List<ProcessFileResponse> response = new List<ProcessFileResponse>(); ProcessFileResponse outResponse = new ProcessFileResponse(); outResponse.Errors = 0; outResponse.Warnings = 0; outResponse.Status = 1; outResponse.TabName = "Background action"; outResponse.Message = $"RunTime: {elapsedTime}"; response.Add(outResponse); // 添加任务结束消息 messages.Add(new SignalRMessageAction("taskEnded") { Arguments = new object[] { response }, UserId = JobId }); _logger.LogInformation("PerformBackgroundJob completed"); return messages; }
方法2:直接调用Azure SignalR服务客户端发送消息
通过依赖注入获取SignalR服务客户端,在任务执行过程中实时推送消息,这种方式支持实时更新进度,不依赖输出绑定的返回值,是长耗时任务推送进度的最优方案。
步骤1:配置依赖注入
在Program.cs中注册SignalR服务:
var host = new HostBuilder() .ConfigureFunctionsWorkerDefaults() .ConfigureServices(services => { services.AddAzureSignalR(); }) .Build(); host.Run();
步骤2:修改函数代码,注入SignalR客户端
private readonly IHubContext<ProgressHub> _hubContext; private readonly ILogger<YourFunctionClass> _logger; public YourFunctionClass(IHubContext<ProgressHub> hubContext, ILogger<YourFunctionClass> logger) { _hubContext = hubContext; _logger = logger; } [Function("PerformBackgroundJob")] public async Task PerformBackgroundJob( [QueueTrigger(LibraryConstants.testQueue, Connection = "connectionstring-azure-storage")] string JobId) { _logger.LogInformation("PerformBackgroundJob started"); Stopwatch sw1 = new Stopwatch(); sw1.Start(); // 发送任务开始消息 await _hubContext.Clients.User(JobId).SendAsync("taskStarted", "PerformBackgroundJob started"); _logger.LogInformation("PerformBackgroundJob data reviewed"); // 实时发送进度消息 for (int i = 0; i < 100; i++) { await _hubContext.Clients.User(JobId).SendAsync("taskProgressChanged", i + 1); await Task.Delay(200); // 用Task.Delay替代Thread.Sleep,避免阻塞线程池 } sw1.Stop(); TimeSpan ts = sw1.Elapsed; string elapsedTime = String.Format("{0:00}:{1:00}:{2:00}.{3:00}", ts.Hours, ts.Minutes, ts.Seconds, ts.Milliseconds / 10); List<ProcessFileResponse> response = new List<ProcessFileResponse>(); ProcessFileResponse outResponse = new ProcessFileResponse(); outResponse.Errors = 0; outResponse.Warnings = 0; outResponse.Status = 1; outResponse.TabName = "Background action"; outResponse.Message = $"RunTime: {elapsedTime}"; response.Add(outResponse); // 发送任务结束消息 await _hubContext.Clients.User(JobId).SendAsync("taskEnded", response); _logger.LogInformation("PerformBackgroundJob completed"); } // 定义SignalR Hub类 public class ProgressHub : Hub { }
方法3:使用中间队列触发独立的SignalR消息函数
将每个进度消息写入独立队列,再用另一个Azure Function监听该队列并发送对应的SignalR消息。这种方式解耦了任务执行与消息发送逻辑,避免长耗时任务占用函数资源,适合高并发场景。
步骤1:修改主函数,将消息写入中间队列
[Function("PerformBackgroundJob")] public async Task PerformBackgroundJob( [QueueTrigger(LibraryConstants.testQueue, Connection = "connectionstring-azure-storage")] string JobId, [QueueOutput("signalr-messages", Connection = "connectionstring-azure-storage")] IAsyncCollector<SignalRMsg> messageQueue) { _logger.LogInformation("PerformBackgroundJob started"); Stopwatch sw1 = new Stopwatch(); sw1.Start(); // 写入任务开始消息到队列 await messageQueue.AddAsync(new SignalRMsg() { UserId = JobId, Target = "taskStarted", Arguments = new object[] { "PerformBackgroundJob started" } }); _logger.LogInformation("PerformBackgroundJob data reviewed"); // 写入所有进度消息到队列 for (int i = 0; i < 100; i++) { await messageQueue.AddAsync(new SignalRMsg() { UserId = JobId, Target = "taskProgressChanged", Arguments = new object[] { i + 1 } }); await Task.Delay(200); } sw1.Stop(); TimeSpan ts = sw1.Elapsed; string elapsedTime = String.Format("{0:00}:{1:00}:{2:00}.{3:00}", ts.Hours, ts.Minutes, ts.Seconds, ts.Milliseconds / 10); List<ProcessFileResponse> response = new List<ProcessFileResponse>(); ProcessFileResponse outResponse = new ProcessFileResponse(); outResponse.Errors = 0; outResponse.Warnings = 0; outResponse.Status = 1; outResponse.TabName = "Background action"; outResponse.Message = $"RunTime: {elapsedTime}"; response.Add(outResponse); // 写入任务结束消息到队列 await messageQueue.AddAsync(new SignalRMsg() { UserId = JobId, Target = "taskEnded", Arguments = new object[] { response } }); _logger.LogInformation("PerformBackgroundJob completed"); }
步骤2:创建监听中间队列的SignalR消息函数
[Function("SendSignalRMessage")] [SignalROutput(HubName = "progress", ConnectionStringSetting = "AzureSignalRConnectionString")] public SignalRMessageAction SendSignalRMessage( [QueueTrigger("signalr-messages", Connection = "connectionstring-azure-storage")] SignalRMsg msg) { return new SignalRMessageAction(msg.Target) { Arguments = msg.Arguments, UserId = msg.UserId }; } public class SignalRMsg { public string UserId { get; set; } public string Target { get; set; } public object[] Arguments { get; set; } }
内容的提问来源于stack exchange,提问作者Ivan Cortez
相关产品推荐
相关产品推荐

