You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

NetMQ Dealer-Router模式单元测试间歇性超时故障排查

解决NetMQ Dealer-Router模式下单元测试时而超时的问题

问题根源分析

  1. TCP TIME_WAIT端口占用:测试结束后TCP连接关闭会进入TIME_WAIT状态(默认约2分钟),原端口无法立即被重新绑定,导致后续测试的Router Socket无法正常监听,回复无法送达Dealer。
  2. NetMQ资源清理不彻底:NetMQConfig.Cleanup()调用时机过早,且Worker线程未及时终止,导致Router Socket未正确释放资源。
  3. 线程同步缺陷:使用布尔变量_isRunning控制Worker循环存在线程可见性问题,无法保证Stop命令立即生效,导致Router仍在运行但无法正常发送回复。
  4. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.09 21:44:50