📨 MqOperate — 消息队列编排器
命名空间: Snet.Core.mq | 继承: CoreUnify<MqOperate, MqData> | 源码: Snet.Core/mq/
MqOperate 是 Snet 的消息中间件统一编排引擎。它通过 FileSystemWatcher 监控 plugin/config 目录变化,自动发现并管理所有 IMq 实例的生命周期。
架构
插件 DLL 目录 (lib/mq)
↓ FileSystemWatcher
MqOperate
├── ConcurrentDictionary<string, IMq> InstanceIoc -- 活跃实例
├── ConcurrentDictionary<string, IMq> OnFailIoc -- 启动失败,重试中
├── ConcurrentDictionary<string, PluginModel> PluginModelIoc -- 插件元数据
├── Channel<QueueData> 消息管道 -- 容量 65535
└── 后台重试 + 消费者任务池 -- TaskNumber 控制
配置 (MqData)
| 属性 | 类型 | 默认值 | 说明 |
|---|---|---|---|
LibFolder |
string |
BaseDirectory/lib/mq |
插件 DLL 目录 |
LibConfigFolder |
string |
BaseDirectory/config/mq |
插件配置目录 |
DllWatcherFormat |
string |
"Snet.*.dll" |
DLL 文件匹配模式 |
ConfigWatcherFormat |
string |
"*.Mq.Config.json" |
配置文件匹配模式 |
InterfaceFullName |
string |
"Snet.Model.interface.IMq" |
目标接口全名 |
TaskNumber |
int |
5 |
后台消费者任务数 |
API
生命周期
// 构造时将 MonitorAsync() 保存为 _initializationTask — 开始监听文件变化并加载已有插件;
// Dispose/DisposeAsync 会等待该任务(上限 5 秒)
public MqOperate(MqData basics)
// 启动指定/所有实例
public async Task<OperateResult> OnAsync(List<string>? ISns = null)
// 停止指定/所有实例
public async Task<OperateResult> OffAsync(List<string>? ISns = null)
// 移除并释放指定/所有实例
public async Task<OperateResult> RemoveAsync(List<string>? ISns = null)
// 移除并释放单个实例
public async Task<OperateResult> DisposeAsync(string ISn)
生产/消费
// 生产 — 向指定/所有 IMq 实例发布消息
public async Task<OperateResult> ProduceAsync(string Topic, string Content, List<string>? ISns = null)
public async Task<OperateResult> ProduceAsync(string Topic, byte[] Content, List<string>? ISns = null)
// 生产 — 同步版本(阻塞调用线程直至消息入队;
// 自 26.250.1 起移除 Task.Run 包装 — 直接 GetAwaiter().GetResult())
public OperateResult Produce(string Topic, string Content, List<string>? ISns = null)
public OperateResult Produce(string Topic, byte[] Content, List<string>? ISns = null)
// 消费 — 在指定/所有 IMq 实例上订阅主题
public async Task<OperateResult> ConsumeAsync(string Topic, List<string>? ISns = null)
public async Task<OperateResult> UnConsumeAsync(string Topic, List<string>? ISns = null)
自 26.250.1 起:ProduceAsync / ConsumeAsync / UnConsumeAsync 在 Topic 为 null/空白时抛 ArgumentException(ProduceAsync 在 Content 为 null 时抛 ArgumentNullException)。OnAsync / OffAsync 失败上报使用 null 安全回退消息(实例启动失败,未返回错误信息 / 实例停止失败,未返回错误信息);ConsumeAsync / UnConsumeAsync 上报 订阅失败,未返回错误信息 / 取消订阅失败,未返回错误信息。
自动加载流程
- MqOperate 构造 → MonitorAsync() 保存为 _initializationTask
- 扫描 lib/mq/ 中匹配
Snet.*.dll(或DllWatcherFormat)的文件 → 通过 PluginOperate 加载 - 扫描 config/mq/ 已有 *.Mq.Config.json → 创建实例
- FileSystemWatcher:
- DLL 创建/删除 → 加载插件 / 释放匹配实例并卸载插件
- Config 创建/修改 → 创建实例或释放后重新加载实例
- Config 删除 → 释放并移除活跃实例与重试队列实例
- 启动失败的实例 → 进入 OnFailIoc → 每 1000 ms 后台重试
- Produce 操作 → 有界 Channel 管道(容量 65535)→
TaskNumber个消费者分发
使用示例
using Snet.Core.mq;
// 1. 创建编排器 — 自动开始监控
var mqOrchestrator = await MqOperate.InstanceAsync(new MqData {
LibFolder = "./plugins/mq",
LibConfigFolder = "./config/mq",
TaskNumber = 5
});
// 2. 向所有已加载的 MQ 实例发布消息
await mqOrchestrator.ProduceAsync("sensors/temp", "25.3°C");
// 3. 在所有实例上订阅
await mqOrchestrator.ConsumeAsync("sensors/temp");
// 4. 停止所有实例
await mqOrchestrator.OffAsync();
// 5. 释放单个实例
await mqOrchestrator.DisposeAsync("mq-instance-001");
