C# 如何区分连接同一socket服务的本地不同进程并返回对应响应
C# 本地Socket服务请求批量转发与响应路由实现方案
核心逻辑
你提出的通过标识绑定请求与发起方的思路是可行的,实际实现时建议按单个请求维度分配唯一标识,而非仅绑定进程ID——因为同一个进程可能发起多个并行请求,仅用进程ID会出现响应匹配错乱的问题。
具体实现步骤
- 第一步:定义统一的前后端通信结构,要求本地进程发请求时只需携带业务参数,Socket服务端收到请求后自动为该请求生成全局唯一的
RequestId(可用Guid或自增长整型实现) - 第二步:维护两个线程安全的存储结构:
pendingRequests:ConcurrentDictionary类型,键为RequestId,值存储请求对应的Socket连接对象、原始请求内容等上下文信息,用于后续收到API响应后匹配发起方batchBuffer:ConcurrentQueue类型,临时缓存待批量发送的请求,达到设定的数量阈值(比如20个)或者超时阈值(比如100毫秒无新请求)时触发批量API调用
- 第三步:调用外部API时,将所有待发请求的
RequestId和业务参数一起提交,要求API返回响应时必须带回对应的RequestId。如果第三方API不支持传递自定义标识,可在提交批量请求时记录每个位置对应的RequestId,收到响应后按索引一一对应即可 - 第四步:收到API批量响应后,遍历每个响应项,用返回的
RequestId从pendingRequests中取出对应的Socket连接,将响应序列化后通过该连接发回给发起请求的进程,发送完成后删除pendingRequests中的对应记录,避免内存泄漏
核心代码示例
// 请求上下文,存储请求与对应连接的绑定关系 public class RequestContext { public Guid RequestId { get; set; } public System.Net.Sockets.Socket ClientSocket { get; set; } public int? ProcessId { get; set; } // 可选,需要按进程统计时加上,由客户端上报 public object OriginalRequest { get; set; } } // Socket服务核心逻辑示例 public class SocketBatchForwardService { // 待响应的请求映射,线程安全 private readonly ConcurrentDictionary<Guid, RequestContext> _pendingRequests = new(); // 批量发送缓存队列 private readonly ConcurrentQueue<KeyValuePair<Guid, object>> _batchBuffer = new(); // 批量发送阈值:满20条发一次 private const int BatchThreshold = 20; // 批量发送超时:100ms没新请求也发 private readonly TimeSpan _batchTimeout = TimeSpan.FromMilliseconds(100); private Timer _batchFlushTimer; // 客户端发请求后的处理入口 public async Task HandleIncomingRequest(Socket clientSocket, object requestBody, int? clientProcessId = null) { var requestId = Guid.NewGuid(); // 绑定请求与客户端连接 _pendingRequests.TryAdd(requestId, new RequestContext { RequestId = requestId, ClientSocket = clientSocket, ProcessId = clientProcessId, OriginalRequest = requestBody }); // 加入批量缓存 _batchBuffer.Enqueue(new KeyValuePair<Guid, object>(requestId, requestBody)); // 达到阈值立即触发批量发送 if (_batchBuffer.Count >= BatchThreshold) { await FlushBatchRequestsAsync(); } // 重置超时定时器 _batchFlushTimer?.Change(_batchTimeout, Timeout.InfiniteTimeSpan); } // 执行批量API调用与响应回写 private async Task FlushBatchRequestsAsync() { var currentBatch = new List<KeyValuePair<Guid, object>>(); // 取出当前所有待发请求 while (_batchBuffer.TryDequeue(out var item)) { currentBatch.Add(item); } if (!currentBatch.Any()) return; // 构造API请求参数 var apiRequest = currentBatch.Select(x => new { RequestId = x.Key, BizData = x.Value }).ToList(); // 调用外部API,假设返回结果每个都带RequestId字段 List<ApiResponseItem> apiResponses = await PostToExternalApiAsync(apiRequest); // 响应回写客户端 foreach (var resp in apiResponses) { if (_pendingRequests.TryRemove(resp.RequestId, out var context)) { byte[] respBytes = SerializeResponseToBytes(resp); await context.ClientSocket.SendAsync(respBytes, SocketFlags.None); } } } // 此处为演示用的占位方法,实际替换为自己的API调用、序列化逻辑即可 private Task<List<ApiResponseItem>> PostToExternalApiAsync(object request) => throw new NotImplementedException(); private byte[] SerializeResponseToBytes(object resp) => throw new NotImplementedException(); } // 示例API响应结构 public class ApiResponseItem { public Guid RequestId { get; set; } public object BizResult { get; set; } public int Code { get; set; } }
注意事项
- 必须处理超时场景:给
pendingRequests中的每个请求添加过期时间,定时清理超时的请求并给客户端返回错误信息,避免内存泄漏 - 所有公共存储结构必须使用线程安全的实现,避免多客户端并发请求时出现竞态异常
- 如果不需要按进程做统计、限流等逻辑,完全可以不用存储进程ID,仅靠
RequestId就能实现准确的响应路由,适配同进程多并行请求的场景 - 批量发送的阈值和超时时间可以根据业务对延迟的要求灵活调整,延迟要求高就调小阈值和超时,追求更高的批量合并率就调大对应参数
内容的提问来源于stack exchange,提问作者Kim Sandberg
相关产品推荐
相关产品推荐

