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() 保存为 _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 / UnConsumeAsyncTopic 为 null/空白时抛 ArgumentExceptionProduceAsyncContent 为 null 时抛 ArgumentNullException)。OnAsync / OffAsync 失败上报使用 null 安全回退消息(实例启动失败,未返回错误信息 / 实例停止失败,未返回错误信息);ConsumeAsync / UnConsumeAsync 上报 订阅失败,未返回错误信息 / 取消订阅失败,未返回错误信息

自动加载流程

  1. MqOperate 构造 → MonitorAsync() 保存为 _initializationTask
  2. 扫描 lib/mq/ 中匹配 Snet.*.dll(或 DllWatcherFormat)的文件 → 通过 PluginOperate 加载
  3. 扫描 config/mq/ 已有 *.Mq.Config.json → 创建实例
  4. FileSystemWatcher:
    • DLL 创建/删除 → 加载插件 / 释放匹配实例并卸载插件
    • Config 创建/修改 → 创建实例或释放后重新加载实例
    • Config 删除 → 释放并移除活跃实例与重试队列实例
  5. 启动失败的实例 → 进入 OnFailIoc → 每 1000 ms 后台重试
  6. 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");

参见