同一进程内能否使用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
相关产品推荐
相关产品推荐

