NetMQ - Snet Docs

NetMQ 中间件(ZeroMQ)

包名: Snet.NetMQ | 类: NetMQOperate | 基类: MqAbstract<NetMQOperate, NetMQData.Basics>

使用 ZeroMQ(NetMQ -- ZeroMQ 的原生 .NET 移植版本)提供轻量级、高性能的消息传递。支持发布/订阅消息模式,无需代理服务器。

概述

NetMQOperate 封装了 NetMQ 套接字,用于无代理的高速消息传递。它继承自 MqAbstract<NetMQOperate, NetMQData.Basics> 并实现了 IMq,在所有 Snet 中间件中提供一致的生产者/消费者语义。

ZeroMQ 适用于低延迟要求较高的场景,在这些场景中代理会带来不必要的额外开销。

安装

dotnet add package Snet.NetMQ

快速开始

using Snet.NetMQ;

// NetMQ 单实例只能发布或只能订阅——需要创建两个实例。
// ZMQ slow-joiner:先启动发布端,否则订阅端会错过其订阅建立之前发送的消息。

// --- 发布端(UModel = PubModel,默认) ---
var pub = await NetMQOperate.InstanceAsync(new NetMQData.Basics
{
    Address = "tcp://127.0.0.1:8866"
});

await pub.OnAsync();

// 发布消息
await pub.ProduceAsync("temperature", "25.4", System.Text.Encoding.UTF8);

// --- 订阅端(UModel = SubModel) ---
var sub = await NetMQOperate.InstanceAsync(new NetMQData.Basics
{
    Address = "tcp://127.0.0.1:8866",
    UModel = UseModel.SubModel
});

await sub.OnAsync();

// 绑定事件再消费
sub.OnDataEventAsync += async (sender, e) =>
{
    if (e.Status)
        Console.WriteLine($"接收到数据: {e.ResultData}");
    else
        Console.WriteLine($"消费失败: {e.Message}");
};

await sub.ConsumeAsync("temperature");

// 后续:取消订阅
// await sub.UnConsumeAsync("temperature");
// await pub.DisposeAsync();
// await sub.DisposeAsync();

生产与消费

发布与订阅是 IMq 的一等操作——先绑定数据事件,再消费;发布可用字符串或原始字节:

// 1) 先绑定事件再消费——收到的消息在此到达
mq.OnDataEventAsync += async (sender, e) =>
{
    if (e.Status)
        Console.WriteLine($"收到: {e.ResultData}");
    else
        Console.WriteLine($"消费失败: {e.Message}");
};

// 2) 消费:订阅主题
await mq.ConsumeAsync("temperature");

// 3) 生产:发布字符串消息(UTF-8)
await mq.ProduceAsync("temperature", "25.4", System.Text.Encoding.UTF8);

// 4) 生产:发布原始字节
var bytes = System.Text.Encoding.UTF8.GetBytes("payload");
await mq.ProduceAsync("temperature", bytes);

// 后续:停止消费
// await mq.UnConsumeAsync("temperature");

ProduceAsync(topic, string, Encoding?)ProduceAsync(topic, byte[]) 均可发布;ConsumeAsync/UnConsumeAsync 管理主题订阅。事件见:事件

配置

参数 类型 默认值 描述
SN string 序列号 / 设备标识符
Address string "tcp://127.0.0.1:8866" 套接字地址(如 tcp://127.0.0.1:8866
UModel enum PubModel 使用模式:SubModel, PubModel
TimeOut int 1000 操作超时(毫秒)
ResponseType enum Content 响应类型

支持的模式

NetMQ 支持 2 种模式:PubModel, SubModel。

模式 描述
PubModel 发布者模式 -- 一对多消息扇出到所有连接的订阅者
SubModel 订阅者模式 -- 接收来自发布者的匹配主题前缀过滤器的消息

架构

ZeroMQ 无需中央代理即可运行。发布者和订阅者直接通信,具有以下特点:

  • 超低延迟:无代理中转,直接套接字到套接字的消息传递
  • 零基础设施:无需安装、配置或维护服务器
  • 无自动重连:网络中断后需重新调用 OnAsync 恢复连接

支持的操作

操作 方法 描述
连接 / 绑定 OnAsync() 绑定或连接 NetMQ 套接字
断开 OffAsync() 关闭 NetMQ 套接字
发布 ProduceAsync(topic, string, Encoding?) 发布带主题的字符串消息
发布 ProduceAsync(topic, byte[]) 发布带主题的原始字节
订阅 ConsumeAsync(topic) 订阅主题模式
取消订阅 UnConsumeAsync(topic) 移除主题订阅

事件

事件 签名 描述
OnDataEvent EventHandler<EventDataResult> 收到消息时触发
OnDataEventAsync EventHandlerAsync<EventDataResult> OnDataEvent 的异步变体
OnInfoEvent EventHandler<EventInfoResult> 收到状态和对等事件时触发
OnInfoEventAsync EventHandlerAsync<EventInfoResult> OnInfoEvent 的异步变体
OnLanguageEvent EventHandler<EventLanguageResult> 语言变更时触发
OnLanguageEventAsync EventHandlerAsync<EventLanguageResult> OnLanguageEvent 的异步变体

另请参阅