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

同一进程内能否使用ZeroMQ(NetMQ) TCP传输实现发布订阅通信?

搞定NetMQ TCP订阅收不到消息的问题

嘿,我帮你捋清楚这个问题——你遇到的情况其实是NetMQ里TCP和inproc通信的典型差异导致的,咱们一步步来解决:

问题根源

inproc://是进程内的内存通信,连网络握手都不需要,发布者和订阅者的连接、订阅几乎是瞬间完成的,所以你的代码直接跑没问题。但换成TCP协议时,需要完成三次握手建立连接,而且订阅者还要给发布者发送订阅指令,这个过程需要一点时间,你现在的代码是连接完立刻发消息,这时候订阅者还没完全“注册”到发布者那里,消息自然发不到订阅者手里。

靠谱的解决办法

这里给你几个适合测试场景的方案,按推荐程度来:

1. 用Poller等待连接完成(最可靠)

NetMQ自带的Poller可以监听发布者的Accept事件,确保订阅者真的连上了再发消息,测试用这个最稳妥:

string address = "tcp://localhost:1026";

var pubSocket = new PublisherSocket();
pubSocket.Bind(address);

var subSocket = new SubscriberSocket();
subSocket.Connect(address);
subSocket.SubscribeToAnyTopic();

// 用Poller等发布者收到订阅者的连接
var poller = new Poller();
pubSocket.AcceptReady += (s, e) => poller.Stop();
poller.Add(pubSocket);
poller.StartAsync();

// 再等一小会儿,确保订阅指令也完成了
Task.Delay(100).Wait();

pubSocket.SendFrame("Hello world!", false);
Console.WriteLine(subSocket.ReceiveFrameString());

// 别忘了清理资源
subSocket.Dispose();
pubSocket.Dispose();
poller.Dispose();

2. 加个短暂延迟(简单但凑合用)

如果是快速测试,不想写复杂逻辑,也可以在发消息前加个小延迟,给TCP连接和订阅留够时间:

string address = "tcp://localhost:1026"; 

var pubSocket = new PublisherSocket();
pubSocket.Bind(address);

var subSocket = new SubscriberSocket();
subSocket.Connect(address);
subSocket.SubscribeToAnyTopic();

// 给TCP连接和订阅流程留200毫秒时间
System.Threading.Thread.Sleep(200); 

pubSocket.SendFrame("Hello world!", false);
Console.WriteLine(subSocket.ReceiveFrameString());

subSocket.Dispose();
pubSocket.Dispose();

提醒下:这个方法不够严谨,不同机器的网络延迟不一样,适合临时测试,别用到生产代码里哈。

3. 监听订阅者的Connected事件

也可以盯着订阅者的Connected事件,确认连上后再发消息:

string address = "tcp://localhost:1026";
bool isConnected = false;

var pubSocket = new PublisherSocket();
pubSocket.Bind(address);

var subSocket = new SubscriberSocket();
subSocket.Connected += (s, e) => isConnected = true;
subSocket.Connect(address);
subSocket.SubscribeToAnyTopic();

// 循环等连接完成
while (!isConnected)
{
    System.Threading.Thread.Sleep(10);
}
// 再等一下确保订阅生效
System.Threading.Thread.Sleep(50);

pubSocket.SendFrame("Hello world!", false);
Console.WriteLine(subSocket.ReceiveFrameString());

subSocket.Dispose();
pubSocket.Dispose();

内容的提问来源于stack exchange,提问作者jamesmus

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 09:48:45