🐰 RabbitMQ 中间件
包名: Snet.RabbitMQ | 类: RabbitMQOperate | 基类: MqAbstract<RabbitMQOperate, RabbitMQData.Basics>
通过 RabbitMQ 提供 AMQP 0-9-1 消息传递。支持交换机、队列、路由键以及标准的发布/订阅模式,用于工业消息代理。
概述
RabbitMQOperate 连接到 RabbitMQ 服务器,并实现了 IMq 接口(IProducer + IConsumer),通过 AMQP 进行消息生产和消费。它继承自 MqAbstract<RabbitMQOperate, RabbitMQData.Basics>。
安装
dotnet add package Snet.RabbitMQ
快速开始
using Snet.RabbitMQ;
var rabbit = await RabbitMQOperate.InstanceAsync(new RabbitMQData.Basics
{
IpAddress = "127.0.0.1",
Port = 6688
});
await rabbit.OnAsync();
// 绑定事件再消费
rabbit.OnDataEventAsync += async (sender, e) =>
{
if (e.Status)
Console.WriteLine($"接收到数据: {e.ResultData}");
else
Console.WriteLine($"消费失败: {e.Message}");
};
await rabbit.ConsumeAsync("snet.route");
// 生产消息(可选)
await rabbit.ProduceAsync("snet.route", "hello rabbitmq", System.Text.Encoding.UTF8);
// 后续:取消消费
// await rabbit.UnConsumeAsync("snet.route");
// await rabbit.DisposeAsync();
生产与消费
发布与订阅是 IMq 的一等操作——先绑定数据事件,再消费;发布可用字符串或原始字节:
// 1) 先绑定事件再消费——收到的消息在此到达
rabbit.OnDataEventAsync += async (sender, e) =>
{
if (e.Status)
Console.WriteLine($"收到: {e.ResultData}");
else
Console.WriteLine($"消费失败: {e.Message}");
};
// 2) 消费:订阅主题
await rabbit.ConsumeAsync("snet.route");
// 3) 生产:发布字符串消息(UTF-8)
await rabbit.ProduceAsync("snet.route", "hello rabbitmq", System.Text.Encoding.UTF8);
// 4) 生产:发布原始字节
var bytes = System.Text.Encoding.UTF8.GetBytes("payload");
await rabbit.ProduceAsync("snet.route", bytes);
// 后续:停止消费
// await rabbit.UnConsumeAsync("snet.route");
ProduceAsync(topic, string, Encoding?)与ProduceAsync(topic, byte[])均可发布;ConsumeAsync/UnConsumeAsync管理主题订阅。事件见:事件。
配置
| 参数 | 类型 | 默认值 | 描述 |
|---|---|---|---|
SN |
string | — | 序列号 / 设备标识符 |
IpAddress |
string | "127.0.0.1" |
RabbitMQ 服务器 IP 地址 |
Port |
int | 6688 | AMQP 端口 |
UserName |
string | "shunnet" |
AMQP 登录用户名 |
Password |
string | "shunnet" |
AMQP 登录密码 |
ExChangeName |
string | "exchang" |
发布消息的交换机名称 |
MessageExpirationTime |
int | 86400000 | 消息过期时间(毫秒) |
ResponseType |
enum | Content | 响应类型 |
交换机类型
| 类型 | 描述 |
|---|---|
| Direct | 将消息路由到绑定键与路由键完全匹配的队列 |
| Topic | 将消息路由到绑定键模式与路由键匹配的队列(支持 * 和 # 通配符) |
| Fanout | 将消息路由到所有绑定的队列,忽略路由键 |
| Headers | 基于消息头而非路由键进行消息路由 |
支持的操作
| 操作 | 方法 | 描述 |
|---|---|---|
| 连接 | OnAsync() |
连接到 RabbitMQ 服务器 |
| 断开 | OffAsync() |
断开与服务器的连接 |
| 发布 | ProduceAsync(topic, string, Encoding?) |
向交换机发布带有路由键的字符串消息 |
| 发布 | ProduceAsync(topic, byte[]) |
发布原始字节 |
| 订阅 | ConsumeAsync(topic) |
绑定队列并开始消费 |
| 取消订阅 | UnConsumeAsync(topic) |
取消消费(BasicCancel)——队列与未消费消息保留;主题无消费登记时返回失败「此主题未在消费」。如需删除队列及遗留消息,请通过 Broker 管理 API 处理 |
| 发布(全参数) | Publish(content, Topic, Queue, RoutingKey, Type, Durable, ...) / PublishAsync(...) |
以显式交换机类型、队列、路由键与持久化设置发布(内容为字符串或 byte[]) |
| 订阅(全参数) | Consume(Topic, Queue, RoutingKey, Type, AutoAck, Durable, ...) / ConsumeAsync(...) |
以显式队列、交换机类型、持久化与确认模式消费 |
全参数重载默认
Type = "topic"、AutoAck = false(手动确认语义)。手动确认模式下,处理抛异常的消息会被拒绝(BasicReject,requeue: true)并重新投递——不会被静默丢弃。
事件
| 事件 | 签名 | 描述 |
|---|---|---|
OnDataEvent |
EventHandler<EventDataResult> |
从队列消费到消息时触发 |
OnDataEventAsync |
EventHandlerAsync<EventDataResult> |
OnDataEvent 的异步变体 |
OnInfoEvent |
EventHandler<EventInfoResult> |
收到状态和代理信息时触发 |
OnInfoEventAsync |
EventHandlerAsync<EventInfoResult> |
OnInfoEvent 的异步变体 |
OnLanguageEvent |
EventHandler<EventLanguageResult> |
语言变更时触发 |
OnLanguageEventAsync |
EventHandlerAsync<EventLanguageResult> |
OnLanguageEvent 的异步变体 |
