MQ 编排器 - Snet Docs

📨 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");

参见