如何通过SignalR实现ASP.NET Web API的消息收发及响应回传
解决方案
核心思路是通过唯一请求追踪ID关联Web API请求和SignalR异步响应,结合TaskCompletionSource实现Web API对SignalR回调的异步等待,最终将设备返回的数据回传给客户端。
具体实现步骤与代码示例
生成唯一追踪ID并关联请求
为每个Web API请求生成唯一的traceId,将其与TaskCompletionSource绑定,用来等待SignalR的响应结果。Web API控制器实现
控制器中发送SignalR消息后,等待设备端的响应,超时后返回错误,最终清理资源:
[ApiController] [Route("api/device")] public class DeviceController : ControllerBase { private readonly IHubContext<DeviceHub> _hubContext; // 单实例部署用静态字典,多实例需替换为分布式缓存(如Redis) private static readonly Dictionary<string, TaskCompletionSource<string>> _pendingRequests = new(); public DeviceController(IHubContext<DeviceHub> hubContext) { _hubContext = hubContext; } [HttpGet("{deviceId}/data")] public async Task<IActionResult> GetDeviceData(string deviceId) { var traceId = Guid.NewGuid().ToString(); var tcs = new TaskCompletionSource<string>(); // 绑定追踪ID与等待任务 lock (_pendingRequests) { _pendingRequests[traceId] = tcs; } try { // 向指定设备分组发送数据请求,携带追踪ID await _hubContext.Clients.Group(deviceId).SendAsync("FetchData", traceId); // 设置30秒超时,避免客户端无限等待 var timeoutTask = Task.Delay(TimeSpan.FromSeconds(30)); var completedTask = await Task.WhenAny(tcs.Task, timeoutTask); if (completedTask == timeoutTask) { return StatusCode(StatusCodes.Status504GatewayTimeout, "设备响应超时"); } var rawData = await tcs.Task; var deviceData = JsonSerializer.Deserialize<DeviceData>(rawData); return Ok(deviceData); } finally { // 清理绑定关系,防止内存泄漏 lock (_pendingRequests) { _pendingRequests.Remove(traceId); } } } }
- SignalR Hub处理设备响应
Hub中提供ReturnData方法,接收设备端传回的数据,通过追踪ID找到对应的TaskCompletionSource并完成任务:
public class DeviceHub : Hub { // 设备端调用此方法返回数据 public async Task ReturnData(string traceId, string data) { if (_pendingRequests.TryGetValue(traceId, out var tcs)) { tcs.SetResult(data); } } // 设备连接时加入对应分组,便于按deviceId定向发送消息 public async Task JoinGroup(string deviceId) { await Groups.AddToGroupAsync(Context.ConnectionId, deviceId); } }
- 防火墙后设备端SignalR客户端实现
设备端监听云端的FetchData指令,处理后通过ReturnData方法将数据传回云端:
// 设备端SignalR客户端初始化 var connection = new HubConnectionBuilder() .WithUrl("https://your-cloud-domain/deviceHub") .Build(); // 监听云端的数据请求 connection.On<string>("FetchData", async (traceId) => { // 替换为实际的设备数据获取逻辑 var deviceData = new DeviceData { DeviceId = "device1", SensorValue = "25.6°C", Timestamp = DateTime.UtcNow }; var dataJson = JsonSerializer.Serialize(deviceData); // 将数据传回云端,携带追踪ID await connection.SendAsync("ReturnData", traceId, dataJson); }); // 断线重连逻辑 connection.Closed += async (error) => { await Task.Delay(new Random().Next(0, 5) * 1000); await connection.StartAsync(); // 重连后重新加入分组 await connection.SendAsync("JoinGroup", "device1"); }; // 启动连接并加入分组 await connection.StartAsync(); await connection.SendAsync("JoinGroup", "device1");
关键注意事项
- 多实例部署适配:如果云端API是多实例集群,静态字典无法跨实例共享,需改用Redis等分布式缓存存储追踪ID与响应任务的映射,或通过消息队列(如RabbitMQ)转发设备响应到对应实例。
- 错误处理:设备端处理失败时,应调用
ReturnData返回错误信息,云端可通过tcs.SetException或返回特定错误状态码告知客户端。 - 超时与资源清理:必须设置请求超时,同时在finally块中清理绑定关系,避免内存泄漏;分布式场景下需给缓存条目设置过期时间。
内容的提问来源于stack exchange,提问作者Ben
相关产品推荐
相关产品推荐

