NetMQ在Unity中Pub/Sub模式文件传输异常求助
NetMQ文件传输功能异常排查请求
我在使用NetMQ时遇到功能异常:代码编译正常,但无法实现将test.wav从“Sending Folder”传输到“Receiving Folder”的目标。当前在同一台PC上测试发布端(Pub)和订阅端(Sub),通过UI Handler绑定的四个按钮控制启停,以下是相关代码,请帮忙排查问题。
发布端逻辑
using System.Threading; using NetMQ; using NetMQ.Sockets; using System.IO; public class Publisher { private readonly Thread _publisherThread; private readonly Thread _subscriberThread; public Publisher() { _publisherThread = new Thread(PublisherWork); _publisherThread.Start(); _subscriberThread = new Thread(SubscriberWork); _subscriberThread.Start(); } public void Start() { using (var pubSocket = new PublisherSocket()) { pubSocket.Bind("tcp://*:5557"); string filePath = "C:/Users/XXXXX/Desktop/0MQ Demo/Sending Folder/test.wav"; byte[] fileBytes = File.ReadAllBytes(filePath); pubSocket.SendMoreFrame("File").SendFrame(fileBytes); } } private void PublisherWork() { using (var pubSocket = new PublisherSocket()) { pubSocket.Bind("tcp://*:5557"); // [C] Port string filePath = "C:/Users/XXXXX/Desktop/0MQ Demo/Sending Folder/test.wav"; // [C] Filepath byte[] fileBytes = File.ReadAllBytes(filePath); while (true) { pubSocket.SendMoreFrame("File").SendFrame(fileBytes); // [C] Topic Thread.Sleep(1000); } } } private void SubscriberWork() { using (var subSocket = new SubscriberSocket()) { subSocket.Connect("tcp://localhost:5557"); // [C] Port subSocket.Subscribe("File"); // [C] Topic while (true) { var topic = subSocket.ReceiveFrameString(); var fileBytes = subSocket.ReceiveFrameBytes(); SaveWavFile(fileBytes); } } } private void SaveWavFile(byte[] fileBytes) { string filePath = "C:/Users/XXXXX/Desktop/0MQ Demo/Receiving Folder/test.wav"; File.WriteAllBytes(filePath, fileBytes); } }
订阅端逻辑
using System.Threading; using NetMQ; using NetMQ.Sockets; using UnityEngine; using UnityEngine.UI; using System.IO; public class Subscriber : MonoBehaviour { public Button startButton; public Button stopButton; private Thread _publisherThread; private Thread _subscriberThread; private bool _isRunning; private void Start() { startButton.onClick.AddListener(StartThreads); stopButton.onClick.AddListener(StopThreads); } private void StartThreads() { _isRunning = true; _publisherThread = new Thread(PublisherWork); _publisherThread.Start(); _subscriberThread = new Thread(SubscriberWork); _subscriberThread.Start(); } private void StopThreads() { _isRunning = false; _publisherThread.Join(); _subscriberThread.Join(); } private void PublisherWork() { using (var pubSocket = new PublisherSocket()) { pubSocket.Bind("tcp://*:5557"); // [C] Port string filePath = "C:/Users/XXXXX/Desktop/0MQ Demo/Sending Folder/test.wav"; // [C] Filepath byte[] fileBytes = File.ReadAllBytes(filePath); while (_isRunning) { pubSocket.SendMoreFrame("File").SendFrame(fileBytes); // [C] Topic Thread.Sleep(1000); } } NetMQConfig.Cleanup(); } private void SubscriberWork() { using (var subSocket = new SubscriberSocket()) { subSocket.Connect("tcp://localhost:5557"); // [C] Port subSocket.Subscribe("File"); // [C] Topic while (_isRunning) { var topic = subSocket.ReceiveFrameString(); var fileBytes = subSocket.ReceiveFrameBytes(); SaveWavFile(fileBytes); } } NetMQConfig.Cleanup(); } private void SaveWavFile(byte[] fileBytes) { string filePath = "C:/Users/XXXXX/Desktop/0MQ Demo/Receiving Folder/test.wav"; // [C] File.WriteAllBytes(filePath, fileBytes); } }
UI处理逻辑
using System.Threading; using NetMQ; using NetMQ.Sockets; using System.IO; using UnityEngine; using UnityEngine.UI; public class UIHandler : MonoBehaviour { public Button startSubscriberButton; public Button stopSubscriberButton; public Button startPublisherButton; public Button stopPublisherButton; public Text pubStartStatus; public Text pubStopStatus; public Text subStartStatus; public Text subStopStatus; private Publisher _publisher; private readonly Thread _publisherThread; private readonly Thread _subscriberThread; private bool _publisherRunning = false; private bool _subscriberRunning = false; public UIHandler() { // Initializing the threads _publisherThread = new Thread(PublisherWork); _subscriberThread = new Thread(SubscriberWork); } private void Start() { _publisher = new Publisher(); startPublisherButton.onClick.AddListener(StartPublisher); stopPublisherButton.onClick.AddListener(StopPublisher); startSubscriberButton.onClick.AddListener(StartSubscriber); stopSubscriberButton.onClick.AddListener(StopSubscriber); } public void StartPublisher() { // Starting the publisher thread _publisherThread.Start(); _publisherRunning = true; pubStartStatus.text = "pubStartStatus: Started"; } public void StopPublisher() { // Stopping the publisher thread _publisherRunning = false; _publisherThread.Join(); pubStopStatus.text = "pubStartStatus: Stopped"; } public void StartSubscriber() { // Starting the subscriber thread _subscriberThread.Start(); _subscriberRunning = true; subStartStatus.text = "subStartStatus: Stopped"; } public void StopSubscriber() { // Stopping the subscriber thread _subscriberRunning = false; _subscriberThread.Join(); subStopStatus.text = "subStartStatus: Stopped"; } private void PublisherWork() { // Creating a publisher socket and binding it to the specified address using (var pubSocket = new PublisherSocket()) { pubSocket.Bind("tcp://*:5557"); // Reading the file bytes from the specified location string filePath = "C:/Users/XXXXX/Desktop/0MQ Demo/Sending Folder/test.wav"; byte[] fileBytes = File.ReadAllBytes(filePath); // Publishing the file bytes continuously while (_publisherRunning) { pubSocket.SendMoreFrame("File").SendFrame(fileBytes); Thread.Sleep(1000); } } NetMQConfig.Cleanup(); } private void SubscriberWork() { // Creating a subscriber socket and connecting it to the specified address using (var subSocket = new SubscriberSocket()) { subSocket.Connect("tcp://localhost:5557"); subSocket.Subscribe("File"); // Receiving and saving the file bytes continuously while (_subscriberRunning) { var topic = subSocket.ReceiveFrameString(); var fileBytes = subSocket.ReceiveFrameBytes(); SaveWavFile(fileBytes); } } NetMQConfig.Cleanup(); } private void SaveWavFile(byte[] fileBytes) { // Saving the received file bytes to the specified location string filePath = "C:/Users/XXXXX/Desktop/0MQ Demo/Receiving Folder/test.wav"; File.WriteAllBytes(filePath, fileBytes); } }
核心问题排查与修复建议
1. 端口冲突与重复绑定
多个类(Publisher、Subscriber、UIHandler)都绑定了tcp://*:5557,启动时会触发端口占用错误,导致后启动的Socket无法正常工作。解决:只保留UIHandler作为唯一控制入口,删除冗余的Publisher和Subscriber类。
2. 线程重复启动与状态控制异常
UIHandler构造函数提前初始化线程,多次点击启动按钮会抛出“线程已启动”异常;解决:在StartPublisher/StartSubscriber方法内动态创建线程,启动前先判断线程状态。StartSubscriber方法中错误设置状态为subStartStatus: Stopped;解决:改为subStartStatus: Started。
3. 阻塞接收导致线程无法终止
订阅端使用ReceiveFrameString()阻塞接收,即使设置_subscriberRunning = false,线程也无法及时响应终止信号;解决:使用TryReceiveFrameString(TimeSpan.FromMilliseconds(500), out topic)带超时的接收方法,让线程有机会检查终止标志。
4. NetMQ清理时机错误
单个线程内调用NetMQConfig.Cleanup()会破坏全局Socket上下文;解决:在所有线程终止后,在OnDestroy()中统一调用一次清理方法。
5. Pub/Sub启动顺序问题
Pub/Sub模式下,订阅端必须先完成订阅,发布端再发送消息,否则初始消息会丢失;解决:提示用户先启动订阅端,再启动发布端,或在代码中添加同步逻辑。
修正后的UIHandler示例代码
using System.Threading; using System.Threading.Tasks; using NetMQ; using NetMQ.Sockets; using System.IO; using UnityEngine; using UnityEngine.UI; public class UIHandler : MonoBehaviour { public Button startSubscriberButton; public Button stopSubscriberButton; public Button startPublisherButton; public Button stopPublisherButton; public Text pubStatus; public Text subStatus; private Thread _publisherThread; private Thread _subscriberThread; private volatile bool _publisherRunning = false; private volatile bool _subscriberRunning = false; private void Start() { startPublisherButton.onClick.AddListener(StartPublisher); stopPublisherButton.onClick.AddListener(StopPublisher); startSubscriberButton.onClick.AddListener(StartSubscriber); stopSubscriberButton.onClick.AddListener(StopSubscriber); pubStatus.text = "Publisher: Stopped"; subStatus.text = "Subscriber: Stopped"; } public void StartPublisher() { if (_publisherRunning || (_publisherThread != null && _publisherThread.IsAlive)) return; _publisherRunning = true; _publisherThread = new Thread(PublisherWork); _publisherThread.Start(); pubStatus.text = "Publisher: Started"; } public void StopPublisher() { _publisherRunning = false; if (_publisherThread != null && _publisherThread.IsAlive) { _publisherThread.Join(2000); if (_publisherThread.IsAlive) _publisherThread.Abort(); } pubStatus.text = "Publisher: Stopped"; } public void StartSubscriber() { if (_subscriberRunning || (_subscriberThread != null && _subscriberThread.IsAlive)) return; _subscriberRunning = true; _subscriberThread = new Thread(SubscriberWork); _subscriberThread.Start(); subStatus.text = "Subscriber: Started"; } public void StopSubscriber() { _subscriberRunning = false; if (_subscriberThread != null && _subscriberThread.IsAlive) { _subscriberThread.Join(2000); if (_subscriberThread.IsAlive) _subscriberThread.Abort(); } subStatus.text = "Subscriber: Stopped"; } private void PublisherWork() { try { using (var pubSocket = new PublisherSocket()) { pubSocket.Bind("tcp://*:5557"); string filePath = "C:/Users/XXXXX/Desktop/0MQ Demo/Sending Folder/test.wav"; byte[] fileBytes = File.ReadAllBytes(filePath); while (_publisherRunning) { pubSocket.SendMoreFrame("File").SendFrame(fileBytes); Thread.Sleep(1000); } } } catch (ThreadAbortException) { } } private void SubscriberWork() { try { using (var subSocket = new SubscriberSocket()) { subSocket.Connect("tcp://localhost:5557"); subSocket.Subscribe("File"); while (_subscriberRunning) { string topic; byte[] fileBytes; if (subSocket.TryReceiveFrameString(TimeSpan.FromMilliseconds(500), out topic) && subSocket.TryReceiveFrameBytes(TimeSpan.FromMilliseconds(500), out fileBytes)) { SaveWavFile(fileBytes); } } } } catch (ThreadAbortException) { } } private void SaveWavFile(byte[] fileBytes) { string filePath = "C:/Users/XXXXX/Desktop/0MQ Demo/Receiving Folder/test.wav"; File.WriteAllBytes(filePath, fileBytes); } private void OnDestroy() { StopPublisher(); StopSubscriber(); NetMQConfig.Cleanup(); } }
内容的提问来源于stack exchange,提问作者AureliaSHABI
相关产品推荐
相关产品推荐

