📨 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 — 开始监听文件变化并加载已有插件
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)
// 消费 — 在指定/所有 IMq 实例上订阅主题
public async Task<OperateResult> ConsumeAsync(string Topic, List<string>? ISns = null)
public async Task<OperateResult> UnConsumeAsync(string Topic, List<string>? ISns = null)
自动加载流程
1. MqOperate 构造 → MonitorAsync() 启动
2. 扫描 lib/mq/ 已有 *.dll → 通过 PluginOperate 加载
3. 扫描 config/mq/ 已有 *.Mq.Config.json → 创建实例
4. FileSystemWatcher:
- DLL 创建 → 加载插件
- DLL 删除 → 卸载部分实例
- Config 创建/修改 → 创建/重新配置实例
- Config 删除 → 移除对应实例
5. 启动失败的实例 → 进入 OnFailIoc → 后台定时重试
6. Consume/Produce 操作 → 通过 Channel 管道 → 并发消费者分发
使用示例
using Snet.Core.mq;
// 1. 创建编排器 — 自动开始监控
var mqOrchestrator = new MqOperate(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");
