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 的异步变体 |
