如何向Kestrel中已连接的指定Socket发送消息?
解决方案
首先纠正你代码里的几个基础错误:
- 接口属性拼写错误:
[HttPost]应该改为[HttpPost] - 发送消息时的编码错误:
socket.Send(Encoding.UTF8.GetString(msg))应该用GetBytes而非GetString,因为Send方法需要字节数组
接下来是核心的连接管理和消息发送实现:
1. 实现TCP连接的Socket管理
Kestrel的ConnectionHandler会处理每个新接入的TCP连接,我们需要在自定义的MyTcpHandler中维护已连接Socket的映射关系,使用线程安全的集合来应对并发场景:
using System.Collections.Concurrent; using System.Net; using System.Net.Sockets; using Microsoft.AspNetCore.Connections; public class MyTcpHandler : ConnectionHandler { // 线程安全的集合,存储IPEndpoint与对应Socket的映射 public static ConcurrentDictionary<IPEndPoint, Socket> ConnectedSockets { get; } = new(); public override async Task OnConnectedAsync(ConnectionContext context) { // 获取当前连接的Socket对象 var socket = context.Transport.GetType().GetProperty("Socket")?.GetValue(context.Transport) as Socket; if (socket == null) { await context.DisposeAsync(); return; } // 获取远程设备的IPEndpoint var remoteEndPoint = socket.RemoteEndPoint as IPEndPoint; if (remoteEndPoint != null) { // 将Socket加入映射集合 ConnectedSockets.TryAdd(remoteEndPoint, socket); } try { // 保持连接存活(Kestrel需要持续读取数据,否则会断开连接) await context.Transport.Input.CopyToAsync(context.Transport.Output); } finally { // 连接断开时从集合中移除 if (remoteEndPoint != null) { ConnectedSockets.TryRemove(remoteEndPoint, out _); } await context.DisposeAsync(); } } }
2. 完善HttpPost接口的实现
在你的Controller中,通过MyTcpHandler.ConnectedSockets集合获取对应Socket,并完成消息发送,同时处理异常情况:
using System.Net; using System.Net.Sockets; using System.Text; using Microsoft.AspNetCore.Mvc; [ApiController] [Route("[controller]")] public class DevicesController : ControllerBase { // 假设这是你维护的设备ID到IPEndpoint的映射表 private readonly Dictionary<int, IPEndPoint> _ipEndPoints = new() { { 1, new IPEndPoint(IPAddress.Parse("192.168.1.100"), 1234) }, { 2, new IPEndPoint(IPAddress.Parse("192.168.1.101"), 1234) } }; [HttpPost("devices/{id}")] public async Task<IActionResult> Send(int id, [FromBody] string msg) { // 验证设备ID是否存在 if (!_ipEndPoints.TryGetValue(id, out var ipEndPoint)) { return NotFound($"设备ID {id} 不存在"); } // 从映射集合中获取对应的Socket if (!MyTcpHandler.ConnectedSockets.TryGetValue(ipEndPoint, out var socket)) { return BadRequest($"设备 {ipEndPoint} 未连接"); } try { // 将消息转为字节数组并发送 var buffer = Encoding.UTF8.GetBytes(msg); await socket.SendAsync(buffer, SocketFlags.None); return Ok("消息发送成功"); } catch (SocketException ex) { // 发送失败时移除无效连接 MyTcpHandler.ConnectedSockets.TryRemove(ipEndPoint, out _); return StatusCode(500, $"消息发送失败:{ex.Message}"); } } }
关键说明
- 线程安全集合:使用
ConcurrentDictionary确保多线程场景下(同时有TCP连接建立/断开、HTTP请求调用)的操作安全 - 连接存活处理:在
OnConnectedAsync中执行Input.CopyToAsync(Output)是为了让Kestrel保持连接,否则Kestrel会因为没有数据交互而主动断开TCP连接 - 异常处理:发送消息时捕获
SocketException,并移除无效连接,避免后续请求使用已断开的Socket
内容的提问来源于stack exchange,提问作者Vincent
相关产品推荐
相关产品推荐

