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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 00:25:26