NetMQ Dealer-Router模式单元测试间歇性超时故障排查
解决NetMQ Dealer-Router模式下单元测试时而超时的问题
问题根源分析
- TCP TIME_WAIT端口占用:测试结束后TCP连接关闭会进入TIME_WAIT状态(默认约2分钟),原端口无法立即被重新绑定,导致后续测试的Router Socket无法正常监听,回复无法送达Dealer。
- NetMQ资源清理不彻底:
NetMQConfig.Cleanup()调用时机过早,且Worker线程未及时终止,导致Router Socket未正确释放资源。 - 线程同步缺陷:使用布尔变量
_isRunning控制Worker循环存在线程可见性问题,无法保证Stop命令立即生效,导致Router仍在运行但无法正常发送回复。 - Dealer Socket使用错误:原测试中Dealer使用
Bind而非Connect,违背了Dealer作为客户端连接Router的角色逻辑,引发连接异常。
具体解决方案
1. 启用端口重用
给Router和Dealer Socket添加ReuseAddress配置,允许端口在TIME_WAIT状态下被立即重用:
// Router Socket配置 routerSocket.Options.ReuseAddress = true; routerSocket.Bind(address); // Dealer Socket配置 dealer.Options.ReuseAddress = true; dealer.Connect(address);
2. 使用CancellationToken控制Worker线程终止
替换_isRunning布尔变量为CancellationTokenSource,确保Stop命令能立即终止Worker循环:
private CancellationTokenSource _cts; public void Start(string address) { _cts = new CancellationTokenSource(); Task.Run(() => Worker(address, _cts.Token), _cts.Token); } public void Stop() { _cts.Cancel(); _cts.Dispose(); // 可选:等待Worker线程完全退出 } private void Worker(string address, CancellationToken token) { using var routerSocket = new RouterSocket(); routerSocket.Options.ReuseAddress = true; routerSocket.Bind(address); while (!token.IsCancellationRequested) { var msg = new NetMQMessage(); try { if (!routerSocket.TryReceiveMultipartMessage(TimeSpan.FromMilliseconds(100), ref msg, 2)) continue; } catch (NetMQException) { /* 异常处理 */ } catch (OperationCanceledException) { break; } if (msg == null || msg.FrameCount == 0) continue; if (msg.FrameCount != 2) throw new InvalidOperationException("Unexpected msg received..."); var identity = msg.Pop().ConvertToString(); var content = msg.Pop().ConvertToString(); Reply result; try { result = HandleMessage(content); } catch (Exception e) { result = new Reply { State = State.Error, Content = e.ToString() }; } if (token.IsCancellationRequested) break; routerSocket.SendMoreFrame(identity).SendFrame(JsonConvert.SerializeObject(result)); } }
3. 测试中使用随机端口
彻底规避端口冲突问题,每次测试生成随机端口:
[Fact] public void InvalidRequest_ReturnsErrorReply() { var port = new Random().Next(50000, 60000); var address = $"tcp://localhost:{port}"; var service = _serviceProvider.GetService<ISocketPulseReceiver>(); service?.Start(address); using var dealer = new DealerSocket(); dealer.Options.ReuseAddress = true; dealer.Connect(address); try { dealer.SendFrame("invalid data"); var received = dealer.TryReceiveFrameString(TimeSpan.FromSeconds(2), out var replyStr); Assert.True(received, "Did not receive a reply in the expected time"); var reply = JsonConvert.DeserializeObject<Reply>(replyStr!); Assert.Equal(State.Error, reply?.State); } finally { service?.Stop(); dealer.Close(); NetMQConfig.Cleanup(); // 短暂等待确保NetMQ IO线程完成清理 Thread.Sleep(100); } }
4. 调整NetMQ清理时机
在测试的finally块中,先停止服务并等待线程结束,再关闭Socket,最后调用NetMQConfig.Cleanup(),确保资源完全释放。
内容的提问来源于stack exchange,提问作者moccajoghurt
相关产品推荐
相关产品推荐

