Lord.Service 7.10.8

dotnet add package Lord.Service --version 7.10.8
                    
NuGet\Install-Package Lord.Service -Version 7.10.8
                    
This command is intended to be used within the Package Manager Console in Visual Studio, as it uses the NuGet module's version of Install-Package.
<PackageReference Include="Lord.Service" Version="7.10.8" />
                    
For projects that support PackageReference, copy this XML node into the project file to reference the package.
<PackageVersion Include="Lord.Service" Version="7.10.8" />
                    
Directory.Packages.props
<PackageReference Include="Lord.Service" />
                    
Project file
For projects that support Central Package Management (CPM), copy this XML node into the solution Directory.Packages.props file to version the package.
paket add Lord.Service --version 7.10.8
                    
#r "nuget: Lord.Service, 7.10.8"
                    
#r directive can be used in F# Interactive and Polyglot Notebooks. Copy this into the interactive tool or source code of the script to reference the package.
#:package Lord.Service@7.10.8
                    
#:package directive can be used in C# file-based apps starting in .NET 10 preview 4. Copy this into a .cs file before any lines of code to reference the package.
#addin nuget:?package=Lord.Service&version=7.10.8
                    
Install as a Cake Addin
#tool nuget:?package=Lord.Service&version=7.10.8
                    
Install as a Cake Tool

LordService 完整使用文档

企业级推送与消息队列集成库

支持钉钉/企业微信/飞书消息推送、RabbitMQ/RocketMQ 发布订阅与死信队列、Redis 分布式缓存与锁、AES-GCM 认证加密、多日志框架。

版本:v7.10.8 | 目标框架:net6.0 / net8.0 / net9.0 / net10.0


目录


安装

dotnet add package Lord.Service

完整 appsettings.json 配置参考

以下是所有模块的完整配置示例,实际使用时按需选取:

{
  "RabbitMQ": {
    "HostName": "localhost",
    "VirtualHost": "/",
    "UserName": "guest",
    "Password": "guest",
    "Prefix": "MyApp",
    "ServiceName": "订单服务"
  },

  "RocketMQ": {
    "NameServerAddress": "localhost:9876",
    "Group": "LordService",
    "Prefix": "MyApp",
    "ServiceName": "订单服务",
    "RequestTimeoutSeconds": 30,
    "ConsumeBatchSize": 1,
    "RetryCount": 5
  },

  "Redis": {
    "ConnectionString": "localhost:6379",
    "Prefix": "MyApp",
    "Database": 0,
    "ConnectTimeout": 5,
    "SyncTimeout": 5,
    "AllowAdmin": false,
    "Ssl": false,
    "Password": "",
    "ClientName": "LordService"
  },

  "Encryption": {
    "PublicKeyOrKey": "你的AES密钥Base64字符串",
    "PrivateKeyOrIV": "备用IV的Base64字符串",
    "EncryptType": "Aes"
  },

  "DingTalk": {
    "Alias": "系统通知群",
    "Url": "https://oapi.dingtalk.com",
    "Token": "your_dingtalk_access_token",
    "Secret": "your_dingtalk_secret"
  },

  "WeChatPush": {
    "Alias": "运维群",
    "Url": "https://qyapi.weixin.qq.com/cgi-bin/webhook/send",
    "Token": "your_wechat_key"
  },

  "LarkPush": {
    "Alias": "开发群",
    "Url": "https://open.feishu.cn/open-apis/bot/v2/hook",
    "Token": "your_lark_token",
    "Secret": "your_lark_secret"
  },

  "DingApp": {
    "AppKey": "your_ding_app_key",
    "AppSecret": "your_ding_app_secret",
    "AgentId": "your_agent_id"
  },

  "WeChatOfficial": {
    "AppId": "your_app_id",
    "AppSecret": "your_app_secret",
    "Token": "your_token",
    "EncodingAesKey": null
  },

  "WeChatMiniProgram": {
    "AppId": "your_mini_program_appid",
    "AppSecret": "your_mini_program_secret",
    "BaseUrl": "https://api.weixin.qq.com"
  },

  "Logging": {
    "LogLevel": {
      "Default": "Information",
      "Microsoft.Hosting.Lifetime": "Information"
    }
  }
}

多群组配置

钉钉/企业微信/飞书支持多群组,使用数组配置:

{
  "DingTalk": [
    {
      "Alias": "生产告警群",
      "Url": "https://oapi.dingtalk.com",
      "Token": "prod_alert_token",
      "Secret": "prod_alert_secret"
    },
    {
      "Alias": "运维群",
      "Url": "https://oapi.dingtalk.com",
      "Token": "ops_token",
      "Secret": "ops_secret"
    }
  ]
}

RabbitMQ 多连接池

不同业务可使用不同的 RabbitMQ 连接池:

{
  "RabbitMQ": {
    "HostName": "localhost",
    "VirtualHost": "/",
    "UserName": "guest",
    "Password": "guest",
    "Prefix": "MyApp"
  },
  "OrderMQ": {
    "HostName": "mq-order.internal",
    "VirtualHost": "/order",
    "UserName": "order_user",
    "Password": "order_pass",
    "Prefix": "OrderApp"
  },
  "OrderQueue": {
    "QueueName": "Order_Process_Queue",
    "ExchangeName": "Order_Process_Topic",
    "RouteKey": "Service.order.process",
    "MQPool": "OrderMQ"
  }
}

快速上手

最简配置(内存缓存 + NLog + AES)

// Program.cs
builder.Services.AddLordService(lord => lord.UseDefaults());

UseDefaults() 等价于:

builder.Services.AddLordService(lord => lord
    .UseHttp()                                          // RestSharp,超时 1 分钟
    .UseLogging(log => log.UseNLog())                  // NLog 日志
    .UseCache(cache => cache.UseMemoryCache())         // 内存缓存
    .UseEncryption(enc => enc.UseAES(a => a.FromConfiguration())) // AES-GCM
);

生产推荐配置

builder.Services.AddLordService(lord => lord
    .UseHttp(h => h.UseNetHttp().WithTimeout(TimeSpan.FromSeconds(30)))
    .UseLogging(log => log.UseNLog())
    .UseCache(cache => cache.UseRedis(r => r.FromConfiguration("Redis")))
    .UseEncryption(enc => enc.UseAES(a => a.FromConfiguration("Encryption")))
    .UseMQ(mq => mq
        .UseDeadLetter()
        .UseRabbitMQ(r => r.FromConfiguration("RabbitMQ")))
    .UsePush(push => push
        .AddDingTalk(d => d.FromConfiguration("DingTalk"))
        .AddWeChat(w => w.FromConfiguration("WeChatPush"))
        .AddLark(l => l.FromConfiguration("LarkPush")))
);

模块详解

1. 缓存(Redis / Memory)

注册方式
// 方式1:从配置加载 Redis
.UseCache(cache => cache.UseRedis(r => r.FromConfiguration("Redis")))

// 方式2:直接传连接字符串
.UseCache(cache => cache.UseRedis("localhost:6379"))

// 方式3:使用内存缓存
.UseCache(cache => cache.UseMemoryCache())

// 方式4:自定义缓存实现
.UseCache(cache => cache.UseCustom<MyRedisCache>(c => c.FromConfiguration("Redis").AsSingleton()))

// 方式5:直接传实例
.UseCache(cache => cache.UseCustom(new MyRedisCache()))
Redis 配置项说明
配置项 类型 默认值 说明
ConnectionString string localhost:6379 Redis 连接字符串
Prefix string "" Key 前缀,用于多系统共用 Redis 隔离
Database int 0 Redis 数据库编号
ConnectTimeout int 5 连接超时(秒)
SyncTimeout int 5 同步操作超时(秒)
AllowAdmin bool false 是否允许管理操作(如 FlushDatabase)
Ssl bool false 是否启用 SSL
Password string? null Redis 密码
ClientName string? LordService 客户端名称
业务使用
public class UserService
{
    private readonly ICache _cache;

    public UserService(ICache cache) => _cache = cache;

    // 基本读写
    public async Task<User?> GetUserAsync(int userId)
    {
        return await _cache.GetItemAsync<User>($"user:{userId}");
    }

    public async Task SetUserAsync(int userId, User user)
    {
        await _cache.SetItemAsync($"user:{userId}", user, TimeSpan.FromMinutes(30));
    }

    // 缓存穿透保护:AddOrGetCacheItem 保证同 key 只回源一次
    public async Task<User> GetOrLoadUserAsync(int userId)
    {
        return await _cache.AddOrGetCacheItemAsync(
            $"user:{userId}",
            async () => await LoadFromDbAsync(userId), // 只在缓存缺失时执行
            TimeSpan.FromMinutes(30),
            isSlidingExpiration: true); // 滑动过期:每次读取续期
    }

    // 批量操作(自动分批,每批 500 条)
    public async Task SetUsersBatchAsync(Dictionary<string, User> users)
    {
        await _cache.SetBatchAsync(users, TimeSpan.FromHours(1));
    }

    public async Task<Dictionary<string, User?>> GetUsersBatchAsync(IEnumerable<string> keys)
    {
        return await _cache.GetBatchAsync<User>(keys);
    }

    // 分布式锁
    public async Task DoWithLockAsync(string resourceKey, Func<Task> action)
    {
        await using var handle = await _cache.DistributedLock.AcquireAsync(
            $"lock:{resourceKey}",
            TimeSpan.FromMinutes(1));

        if (handle.IsAcquired)
        {
            await action();
        }
    }

    // 自动续期锁(看门狗模式,适合长时间任务)
    public async Task DoWithAutoRenewalLockAsync(string resourceKey, Func<Task> action)
    {
        await using var handle = await _cache.DistributedLock.AcquireWithRenewalAsync(
            $"lock:{resourceKey}",
            TimeSpan.FromSeconds(30)); // 锁 30 秒,自动续期

        if (handle.IsAcquired)
        {
            await action(); // 即使任务超过 30 秒,锁也会自动续期
        }
    }
}

安全提示:同步 AddOrGetCacheItem 的回源等待有 30 秒超时保护,避免 cachePopulate 卡住时线程池饥饿。生产环境推荐使用异步 AddOrGetCacheItemAsync


2. 消息队列(RabbitMQ)

注册方式
.UseMQ(mq => mq
    .UseDeadLetter()  // 全局启用死信队列(所有队列自动生成 per-queue DLX)
    .UseRabbitMQ(r => r
        .FromConfiguration("RabbitMQ")           // 从配置加载
        // 或:.WithConnection("host", "/vhost", "user", "pass") // 直接配置
        .WithPrefix("MyApp")                      // 队列名前缀,隔离多系统
        .WithServiceName("订单服务"))              // 连接名称,便于运维识别
)
RabbitMQ 配置项说明
配置项 说明
HostName RabbitMQ 主机地址
VirtualHost 虚拟主机
UserName 用户名
Password 密码
Prefix 队列/交换机/路由键前缀
ServiceName 连接名称(便于 RabbitMQ 管理端识别)
IMQHub 简化 API(推荐)
// 定义消息(队列名自动从类型名推导)
public record OrderCreated(string OrderId, decimal Amount);

// 发布
await _hub.PublishAsync(new OrderCreated("ORD-001", 99.9m));

// 订阅
await _hub.SubscribeAsync<OrderCreated>(async (msg, sp) =>
{
    var logger = sp.GetRequiredService<ILogger<Program>>();
    logger.LogInformation("收到订单: {OrderId}, 金额: {Amount}", msg.OrderId, msg.Amount);
    return true; // true=Ack, false=Reject(触发重试或死信)
});

// 死信订阅
await _hub.SubscribeDeadLetterAsync<OrderCreated>(async (msg, sp) =>
{
    var logger = sp.GetRequiredService<ILogger<Program>>();
    logger.LogWarning("订单消息进入死信: {OrderId}", msg.OrderId);
    return true;
});
IMQHub 完整 API
// 发布
await hub.PublishAsync(message);                              // 类型推导
await hub.PublishAsync(message, cfg => cfg.Qos = 20);         // 自定义配置
await hub.PublishAsync("custom_name", message);               // 自定义队列名
await hub.PublishAsync("custom_name", message, cfg => { });   // 自定义名 + 配置
await hub.PublishAsync(rawJsonString);                        // 原始字符串

// 订阅(5 种重载)
await hub.SubscribeAsync<T>(msg => Task.FromResult(true));            // 最简
await hub.SubscribeAsync<T>(msg => true);                             // 同步
await hub.SubscribeAsync<T>(async (msg, sp) => true);                 // 注入 ServiceProvider
await hub.SubscribeAsync<T>(handler, cfg => { });                     // 自定义配置
await hub.SubscribeAsync<T>("name", handler);                        // 自定义队列名

// 死信订阅
await hub.SubscribeDeadLetterAsync<T>(async (msg, sp) => true);
await hub.SubscribeDeadLetterAsync<T>("name", handler, cfg => { });
MQ 拦截器
.UseMQ(mq => mq
    .AddFilter<MyQueueArgsFilter>()      // IMQQueueArgsFilter: 修改队列参数
    .AddFilter<MyDeadLetterFilter>()     // IMQDeadLetterFilter: 修改死信配置
    .AddFilter<MyCryptoFilter>()         // IMQCryptoFilter: 自定义加解密
    .AddFilter<MyJsonFilter>()           // IMQJsonFilter: 自定义序列化
    .AddFilter<MyExceptionFilter>()      // IMQExceptionFilter: 异常处理
    .UseRabbitMQ(r => r.FromConfiguration("RabbitMQ"))
)
消费端看门狗(RabbitMQ)

RabbitMQ 消费端内置看门狗机制:

  • 心跳检测:每 10 秒检查 consumer 存活状态
  • 自动恢复:consumer 关闭或 Channel 断开时自动重建
  • 指数退避:恢复失败时 1s→2s→4s→8s→16s→30s 退避重试
  • 低频持续:连续失败 20 次后进入低频模式(每 60 秒一次),永不永久停止
  • 超时保护:handler 执行超过 3 分钟自动 Reject 并取消后台任务

2.5 消息队列(RocketMQ)

说明:v7.11.0 新增 RocketMQ 基础支持,采用 NewLife.RocketMQ 客户端,适用于信创场景和国产化 MQ 需求。

注册方式
.UseMQ(mq => mq
    .UseRocketMQ(r => r
        .FromConfiguration("RocketMQ")           // 从配置加载
        // 或:.WithNameServer("127.0.0.1:9876")  // 直接配置 NameServer
        .WithGroup("LordService")                 // 默认生产/消费组
        .WithPrefix("MyApp")                      // Topic/Tag/ConsumerGroup 前缀
        .WithServiceName("订单服务"))              // 连接名称,便于运维识别
)
RocketMQ 配置项说明
配置项 说明
NameServerAddress NameServer 地址,如 127.0.0.1:9876
Group 默认生产/消费组名
Prefix Topic/Tag/ConsumerGroup/ProducerGroup 前缀
ServiceName 服务名称(便于识别)
RequestTimeoutSeconds 请求超时时间(秒)
ConsumeBatchSize 消费批量大小
RetryCount 最大重试次数
RocketMQ 与 RabbitMQ 字段映射
通用字段 RocketMQ 映射 说明
ExchangeName Topic 消息主题
RouteKey Tag 消息标签过滤
QueueName ConsumerGroup 消费者组
EnableDeadLetter EnableFailedMessage 失败消息能力
IMQHub 简化 API(推荐)
// 定义消息(Topic/Tag/ConsumerGroup 自动从类型名推导)
public record OrderCreated(string OrderId, decimal Amount);

// 发布(自动生成 Topic: OrderCreated_Topic, Tag: Service.ordercreated)
await _hub.PublishAsync(new OrderCreated("ORD-001", 99.9m));

// 订阅(自动生成 ConsumerGroup: OrderCreated_ConsumerGroup)
await _hub.SubscribeAsync<OrderCreated>(async (msg, sp) =>
{
    var logger = sp.GetRequiredService<ILogger<Program>>();
    logger.LogInformation("收到订单: {OrderId}, 金额: {Amount}", msg.OrderId, msg.Amount);
    return true; // true=消费成功, false=触发 RocketMQ 重试
});

// 失败消息订阅(通用接口,避免 RabbitMQ DLX 术语)
await _hub.SubscribeFailedAsync<OrderCreated>(async (msg, sp) =>
{
    logger.LogWarning("订单处理失败: {OrderId}", msg.OrderId);
    return true; // 确认已处理
});
失败消息处理

RocketMQ 使用 %DLQ%{ConsumerGroup} 约定 Topic 存储失败消息:

// 消费失败时返回 false,RocketMQ 自动重试
await _hub.SubscribeAsync<OrderCreated>(async msg =>
{
    try
    {
        // 业务处理
        return true; // 成功
    }
    catch
    {
        return false; // 失败,进入 RocketMQ 重试机制
    }
});

// 订阅失败消息
await _hub.SubscribeFailedAsync<OrderCreated>(async msg =>
{
    // 处理进入 %DLQ%{ConsumerGroup} 的失败消息
    return true;
});
高级功能

RocketMQ 支持以下高级功能(通过 IRocketMQAdvancedPushIRocketMQAdvancedReceive 接口):

1. 事务消息
// 获取高级发布接口
var push = await _factory.GetPushServiceAsync<OrderCreated>();
if (push is IRocketMQAdvancedPush advancedPush)
{
    // 发送事务消息(半消息)
    var sendResult = await advancedPush.PublishTransactionAsync(new OrderCreated("ORD-001", 99.9m));
    
    try
    {
        // 执行本地事务
        await _dbContext.SaveOrderAsync(new Order { OrderId = "ORD-001" });
        
        // 提交事务
        await advancedPush.EndTransactionAsync(sendResult, commit: true);
    }
    catch
    {
        // 回滚事务
        await advancedPush.EndTransactionAsync(sendResult, commit: false);
    }
}
2. 顺序消息
// 发送顺序消息(相同 orderKey 的消息进入同一队列)
await advancedPush.PublishOrderAsync(new OrderCreated("ORD-001", 99.9m), orderKey: "ORD-001");
3. 延迟消息
// 发送延迟消息(延迟级别 1-18)
await advancedPush.PublishDelayAsync(new OrderCreated("ORD-001", 99.9m), delayLevel: 3);
// 延迟级别对应:1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h
4. Request-Reply 模式
// 发送请求并等待响应
var response = await advancedPush.RequestAsync<OrderQuery, OrderResult>(
    new OrderQuery("ORD-001"), 
    timeout: 5000);
5. Tag/SQL92 过滤
// 获取高级消费接口
var receive = await _factory.GetReceiveServiceAsync<OrderCreated>();
if (receive is IRocketMQAdvancedReceive advancedReceive)
{
    // Tag 过滤
    await advancedReceive.SubscribeWithFilterAsync<OrderCreated>(
        async msg => { /* 处理消息 */ return true; },
        RocketMQExpressionType.Tag,
        "TagA || TagB");
    
    // SQL92 过滤
    await advancedReceive.SubscribeWithFilterAsync<OrderCreated>(
        async msg => { /* 处理消息 */ return true; },
        RocketMQExpressionType.SQL92,
        "amount > 100 AND status = 'ACTIVE'");
}
6. 多 Topic 订阅
// 一个 Consumer 同时消费多个 Topic
await advancedReceive.SubscribeMultipleTopicsAsync<OrderCreated>(
    async msg => { /* 处理消息 */ return true; },
    topics: "Topic1;Topic2;Topic3");
7. Pop 消费模式(RocketMQ 5.0+)
// 轻量消费模式,无需客户端 Rebalance
await advancedReceive.PopConsumeAsync<OrderCreated>(
    async msg => { /* 处理消息 */ return true; },
    maxNums: 32,
    invisibleTime: 30000);
云厂商支持

RocketMQ 支持主流云厂商的托管服务:

{
  "RocketMQ": {
    "NameServerAddress": "MQ_INST_xxx.aliyuncs.com:80",
    "CloudProvider": "Aliyun",
    "AccessKey": "your_access_key",
    "SecretKey": "your_secret_key",
    "InstanceId": "MQ_INST_xxx"
  }
}
云厂商 CloudProvider 必需配置
阿里云 Aliyun AccessKey, SecretKey, InstanceId
华为云 Huawei AccessKey, SecretKey, InstanceId, EnableSsl
腾讯云 Tencent AccessKey, SecretKey, Namespace
Apache ACL Apache AccessKey, SecretKey
自建 RocketMQ 密码配置(ACL)

v7.10.6 新增:Builder 链式配置,无需手写 appsettings。

自建 RocketMQ 开启 ACL 认证(broker.conf 配置 aclEnable=true + aclAccessKey/aclSecretKey):

.UseMQ(mq => mq.UseRocketMQ(r => r
    .WithNameServer("127.0.0.1:9876")
    .WithAcl("your_access_key", "your_secret_key")))   // ✅ Apache ACL 认证

或者 appsettings.json 配置

{
  "RocketMQ": {
    "NameServerAddress": "127.0.0.1:9876",
    "AccessKey": "your_access_key",
    "SecretKey": "your_secret_key"
  }
}

云厂商托管服务链式配置

.UseMQ(mq => mq.UseRocketMQ(r => r
    .WithNameServer("MQ_INST_xxx.aliyuncs.com:80")
    .WithCloudProvider(RocketMQCloudProvider.Aliyun, "ak", "sk", "MQ_INST_xxx")
    .WithSsl(true)))   // 华为云等需要 SSL

完整 Builder 方法

方法 说明
WithNameServer(addr) NameServer 地址
WithGroup(group) 默认生产/消费组
WithPrefix(prefix) Topic/Tag/Group 前缀
WithAcl(ak, sk) 自建 RocketMQ ACL 认证(等价 CloudProvider=Apache)
WithCloudProvider(type, ak, sk, extra) 云厂商托管(Aliyun/Huawei/Tencent)
WithSsl(enable) SSL/TLS 开关(华为云等)
WithServiceName(name) 服务名称
消费端看门狗(RocketMQ)

RocketMQ 消费端内置看门狗机制(RocketMQSubscriberRocketMQReceive 均支持):

  • 心跳检测:每 10 秒检查消费活动状态
  • 停滞检测:超过 5 分钟无消费活动,判定 consumer 可能已死,触发恢复
  • 自动恢复:Dispose 旧 Consumer + 创建新 Consumer + Start
  • 指数退避:恢复失败时 1s->2s->4s->8s->16s->30s 退避重试
  • 低频持续:连续失败 20 次后进入低频模式(每 60 秒一次),永不永久停止

看门狗默认启用,无需额外配置。详见 MQ 看门狗使用指南

注意事项
  1. 第一版限制:事务消息、顺序消息、延迟消息、ACL 等高级特性通过扩展接口提供,需要手动转换接口类型
  2. Topic 预创建:生产环境建议由运维预先创建 Topic,不依赖自动创建
  3. 多 Provider 并存:同一容器不能同时启用 RabbitMQ 和 RocketMQ
  4. 失败消息:推荐使用 IMQFailedMessageHub 通用接口,避免依赖特定 MQ 术语

2.6 IMQProvider 统一入口(推荐)

v7.10.4 新增IMQProvider 是 MQ 的统一提供者接口,一站式获取所有 MQ 服务。 支持细粒度能力检查,避免"静默降级"问题(如 RabbitMQ 的事务消息实际是普通消息)。

注册方式

v7.10.6 更新UseRocketMQ / UseRabbitMQ 已自动包含 IMQProvider 注册,无需单独调用 UseMQProvider

.UseMQ(mq => mq
    // ✅ 推荐:直接配置引擎,IMQProvider 自动注册
    .UseRocketMQ(r => r
        .WithNameServer("127.0.0.1:9876")
        .WithAcl("ak", "sk")))

UseMQProvider 保留为快捷别名(不关心连接参数、只需切换引擎类型时使用):

.UseMQ(mq => mq.UseMQProvider(MQProviderType.RocketMQ))  // 等价于 UseRocketMQ()
.UseMQ(mq => mq.UseMQProvider(MQProviderType.RabbitMQ))  // 等价于 UseRabbitMQ()
一站式获取所有服务
public class OrderService
{
    private readonly IMQProvider _provider;

    public OrderService(IMQProvider provider) => _provider = provider;

    // 场景 1:IMQHub(日常发布订阅,类型/Topic 驱动)
    public async Task PublishWithHub(OrderEvent order)
    {
        var hub = _provider.GetHub();
        await hub.PublishAsync(order);  // 类型自动推导队列名
    }

    // 场景 2:IMQPublisher(基础发布)
    public async Task PublishWithPublisher(OrderEvent order)
    {
        var publisher = _provider.GetPublisher();
        await publisher.PublishAsync("OrderTopic", order);
    }

    // 场景 3:IMQSubscriber(基础订阅)
    public async Task SubscribeOrders()
    {
        var subscriber = _provider.GetSubscriber();
        await subscriber.SubscribeAsync<OrderEvent>("OrderTopic", "OrderGroup", async order => {
            return true;
        });
    }

    // 场景 4:IMQFactory(按配置节获取)
    public async Task PublishWithFactory(OrderEvent order)
    {
        var factory = _provider.GetFactory();
        var push = await factory.GetPushServiceAsync("OrderMQ");
        await push.PublishAsync(order);
    }
}
细粒度能力检查

IMQProvider 提供 5 个能力标志,业务代码可以提前检查,避免使用不支持的功能:

能力标志 RabbitMQ RocketMQ 说明
SupportsTransactionMessages 事务消息(半消息 + 二次确认)
SupportsDelayedMessages 延迟消息(RocketMQ 18 级)
SupportsRequestReply Request-Reply 模式
SupportsSql92Filtering SQL92 过滤表达式
SupportsOrderedMessages 严格顺序消息
public class PaymentService
{
    private readonly IMQProvider _provider;

    public PaymentService(IMQProvider provider) => _provider = provider;

    // ✅ 事务消息 + 能力检查
    public async Task SendTransactionMessage(PaymentEvent payment)
    {
        if (!_provider.SupportsTransactionMessages)
        {
            throw new NotSupportedException(
                $"当前 MQ 引擎 ({_provider.ProviderType}) 不支持事务消息,请切换到 RocketMQ");
        }

        var publisher = _provider.GetAdvancedPublisher()
            ?? throw new InvalidOperationException("无法获取高级发布者");

        var result = await publisher.PublishTransactionAsync("PaymentTopic", payment);
        try
        {
            await _db.SavePaymentAsync(payment);       // 本地事务
            await publisher.EndTransactionAsync(result, commit: true);
        }
        catch
        {
            await publisher.EndTransactionAsync(result, commit: false);
            throw;
        }
    }

    // ✅ 延迟消息 + 能力检查(RabbitMQ 降级为调度框架)
    public async Task SendDelayMessage(PaymentEvent payment, int delayMinutes)
    {
        if (!_provider.SupportsDelayedMessages)
        {
            // RabbitMQ:使用调度框架替代
            await _scheduler.ScheduleAsync(
                () => _provider.GetPublisher().PublishAsync("PaymentTopic", payment),
                TimeSpan.FromMinutes(delayMinutes));
            return;
        }

        // RocketMQ:原生延迟消息(1-18 级)
        var publisher = _provider.GetAdvancedPublisher()!;
        await publisher.PublishDelayAsync("PaymentTopic", payment, delayLevel: 3);
    }
}
性能说明
  • 能力检查是常量属性,零运行时开销
  • GetHub() / GetFactory() 内部懒加载缓存,只解析一次 DI
  • GetPublisher() / GetSubscriber() 每次返回新实例(Transient 语义),内部连接池(Producer/Consumer)复用,实例创建成本极低
  • Provider 本身是 Singleton,全应用共享一个实例
未来扩展(MQTT 等)

新增 MQ 引擎只需:

  1. MQProviderType 枚举中新增成员(如 MQTT = 2
  2. 实现 IMQProvider(如 MQTTProvider
  3. 实现对应的 Hub/Publisher/Subscriber 适配器
  4. 注册时一行切换,业务代码零改动

2.7 PushMessage 统一消息模型(推荐)

v7.10.6 新增PushMessage渠道无关的统一推送消息模型。 业务代码只依赖 PushMessage 构建消息,推送时由各渠道适配器自动转换为钉钉/飞书/企业微信的 webhook JSON。 切换渠道业务代码零改动(只换注册)。

为什么需要 PushMessage?
❌ 旧方式:业务代码直接使用渠道专属消息类
DingTalkMessage.Markdown("告警", "CPU过高")  → 切飞书要改 LarkMessage.Post(...),方法名/属性全不同

✅ 新方式:PushMessage 统一模型(渠道无关)
PushMessage.Markdown("告警", "CPU过高")      → 钉钉/飞书/企业微信通用
五种消息类型
using Lord.Service.ApiPush.Messages;

// 1. 纯文本
PushMessage.Text("hello world", "13800000000");          // 可选 @用户

// 2. Markdown 富文本
PushMessage.Markdown("告警", "CPU 使用率过高", "138xxxx");

// 3. 链接卡片
PushMessage.Link("标题", "描述", "https://example.com", "picUrl");

// 4. 按钮卡片(整体跳转 / 多个按钮)
PushMessage.ActionCard("确认", "是否继续?", "确认", "https://a.com");
PushMessage.ActionCard("标题", "内容", new[] { ("确认", "https://a.com"), ("取消", "https://b.com") });

// 5. 信息流卡片
PushMessage.FeedCard(new[] { ("标题1", "https://a.com", "pic1.jpg") });
链式构建(IMessageBuilder)

PushMessage / DingTalkMessage / LarkMessage / WeChatMessage 均实现统一的 IMessageBuilder<T> 接口, 支持完全一致的链式构建体验:

var msg = PushMessage.Markdown("告警", "CPU 使用率过高")
    .Add("主机", "192.168.1.1")                              // 单条键值对
    .Add(new Dictionary<string, string> {                     // 批量字典
        ["CPU"] = "95%",
        ["内存"] = "80%"
    })
    .Add(new { 负载 = "2.5", 磁盘 = "70%" })                  // 实体属性自动转键值对
    .At("13800000000");                                       // @用户

// DingTalkMessage / LarkMessage / WeChatMessage 用法完全一致
var ding = new DingTalkMessage("告警").Add("主机", "192.168.1.1").At("138xxxx");
var lark = new LarkMessage("告警").Add("主机", "192.168.1.1");
var wechat = new WeChatMessage("告警").Add("主机", "192.168.1.1");
推送(渠道自动转换)
// 注入推送工厂(钉钉/飞书/企业微信)
public class AlertService(IApiPush push)
{
    public async Task SendAlertAsync(string title, string content)
    {
        var msg = PushMessage.Markdown(title, content).Add("主机", "192.168.1.1");
        await push.PushAsync(msg);   // 自动转换为当前渠道 JSON
    }
}

// 切换渠道 = 只改注册,业务代码零改动!
// 钉钉:    .UsePush(p => p.AddDingTalk(c => c.FromConfiguration("DingPush")))
// 飞书:    .UsePush(p => p.AddLark(c => c.FromConfiguration("LarkPush")))
// 企业微信: .UsePush(p => p.AddWeChat(c => c.FromConfiguration("WeChatPush")))
MQ 场景(生产端发 PushMessage,消费端任意渠道推送)
// 生产端:MQ 里传 PushMessage(渠道无关 JSON)
var msg = PushMessage.Markdown("异常通知", null)
    .Add("异常信息", ex.Message)
    .Add(new { 服务器 = ip, 时间 = DateTime.Now });
await hub.PublishAsync(msg);

// 消费端:订阅 PushMessage,用任意渠道推送
await _hub.SubscribeAsync<PushMessage>(async msg =>
{
    return await push.PushAsync(msg);   // 自动转当前渠道
});

注意:MQ 场景下生产端和消费端的消息类型必须一致(都用 PushMessage 或都用 DingTalkMessage)。

渠道适配器(内部原理)
PushMessage.Markdown("告警", "CPU过高").Add("主机", "192.168.1.1")
    ├─ 钉钉适配器     → {"msgtype":"markdown","markdown":{"title":"告警",...}}
    ├─ 飞书适配器     → {"msg_type":"post","content":{"zh_cn":{"title":"告警",...}}}
    └─ 企业微信适配器 → {"msgtype":"markdown","markdown":{"content":"#### 告警..."}}

各推送实现构造时自动注册自己的适配器,PushAsync(PushMessage) 按实例类型自动转换。

与渠道专属消息类的选择
场景 推荐
需要切换渠道 / MQ 传输 PushMessage(渠道无关)
固定使用钉钉 + 高级定制 DingTalkMessage
固定使用飞书 + 高级定制 LarkMessage
固定使用企业微信 + 高级定制 WeChatMessage

2.8 MQ 消息加密开关

v7.10.6 新增:消息级加密控制。无敏感信息的消息可明文传输(提升性能、避免加解密配置不匹配问题)。

三级控制(优先级从高到低)
IMQSettings(消息级) > UseMQ 全局配置 > appsettings.json(引擎配置)
1. 消息级(IMQSettings.EnableEncryption / EncryptionPrefixes)
// 单个消息关闭加密(无敏感信息)
await hub.PublishAsync(msg, s => { s.EnableEncryption = false; });

// 前缀白名单:只加密指定前缀的队列(其余明文)
await hub.SubscribeAsync<T>(handler, s => { s.EncryptionPrefixes = "secret_;private_"; });
2. UseMQ 全局配置(引擎无关:RabbitMQ/RocketMQ/未来 Kafka)
.UseMQ(mq => mq
    // 全局关闭加密(所有队列明文)
    .EnableEncryption(false)

    // 或:前缀白名单(只加密 "secret_" 和 "private_" 开头的队列)
    .EncryptionPrefixes("secret_", "private_")

    .UseRocketMQ(r => r.FromConfiguration("RocketMQ")))
3. appsettings.json(引擎配置级)
{
  "RocketMQ": {
    "NameServerAddress": "localhost:9876",
    "EnableEncryption": false,
    "EncryptionPrefixes": "secret_;private_"
  }
}
生产端与消费端一致性

生产端和消费端使用相同的队列名时自动保持一致的加密策略(都按队列名前缀判断)。


2.9 MQ 事务功能详解

事务消息是 RocketMQ 的原生能力(半消息 + 二次确认),RabbitMQ 不支持(调用时降级为普通消息并输出警告)。

为什么需要事务消息?

场景:订单系统先落库、再发 MQ 通知下游。两者要么都成功、要么都失败:

❌ 普通消息:先发消息后落库 → 落库失败但消息已发出(下游处理了不存在的订单)
❌ 先落库后发消息 → 发送失败但订单已存在(下游不知道有新订单)

✅ 事务消息:
1. 发送"半消息"(Broker 暂存,不投递)
2. 执行本地事务(落库)
3. 提交 → Broker 投递消息;回滚 → Broker 丢弃消息
4. 若第 2 步崩溃 → Broker 回查事务状态,决定投递还是丢弃
使用方式(RocketMQ)
// 1. 获取高级发布者(IMQProvider 已随 UseRocketMQ 自动注册)
public class OrderService(IMQProvider provider)
{
    public async Task CreateOrderAsync(Order order)
    {
        // 检查引擎能力(RabbitMQ 会返回 false)
        if (!provider.SupportsTransactionMessages)
        {
            throw new NotSupportedException("当前 MQ 引擎不支持事务消息,请使用 RocketMQ");
        }

        var publisher = provider.GetAdvancedPublisher()
            ?? throw new InvalidOperationException("无法获取高级发布者");

        // 2. 发送半消息(Broker 暂存,不投递),返回事务对象
        // ✅ 方式 A:Topic 由类型自动推导(OrderEvent → OrderEvent_Topic,推荐)
        var tx = await publisher.PublishTransactionAsync(order);

        // 方式 B:显式指定 Topic
        // var tx = await publisher.PublishTransactionAsync("OrderTopic", order);

        try
        {
            // 3. 执行本地事务
            await _db.SaveOrderAsync(order);

            // 4. 提交事务 → Broker 投递消息(✅ 事务对象直接提交,直观方便)
            await tx.CommitAsync();
        }
        catch
        {
            // 5. 回滚事务 → Broker 丢弃消息(✅ 事务对象直接回滚)
            await tx.RollbackAsync();
            throw;
        }
    }
}
事务消息的可靠性保障
┌─────────────┐   半消息     ┌──────────────┐   提交     ┌──────────┐
│   业务代码   │ ──────────→ │  RocketMQ    │ ────────→ │  消费者   │
│             │              │  (Broker)    │            │          │
│  1. 发半消息  │              │  暂存不投递   │            │          │
│  2. 本地事务  │              │              │            │          │
│  3. 提交/回滚 │ ←────────── │  回查(可选)  │            │          │
└─────────────┘   回查状态    └──────────────┘            └──────────┘
  • 半消息:Broker 收到但暂不投递,消费者不可见
  • 二次确认:业务提交/回滚后 Broker 才投递或丢弃
  • 事务回查:业务进程崩溃时,Broker 主动回查本地事务状态(需业务实现回查接口,本库提供基础支持)
事务对象 API(推荐)

PublishTransactionAsync 返回的 MQTransactionResult 对象自带提交/回滚方法。

两种重载

// 方式 A:自动推导 Topic(消息类型 → {TypeName}_Topic,泛型安全)
var tx = await publisher.PublishTransactionAsync(order);
//   例如 OrderEvent → OrderEvent_Topic;List<int> → List_Int32_Topic

// 方式 B:显式指定 Topic
var tx = await publisher.PublishTransactionAsync("OrderTopic", order);
var tx = await publisher.PublishTransactionAsync("OrderTopic", order);
await tx.CommitAsync();    // 提交事务(投递消息)
await tx.RollbackAsync();  // 回滚事务(丢弃消息)
  • 事务对象只能结束一次:重复 Commit/Rollback 抛出 InvalidOperationException
  • 引擎未绑定(如 RabbitMQ 降级)时 Commit/Rollback 会明确告警,不再静默失败
  • 旧 API EndTransactionAsync(result, commit) 仍然可用(兼容)
RabbitMQ 的降级行为(警告)
// RabbitMQ 上调用事务 API:
var result = await publisher.PublishTransactionAsync("OrderTopic", order);
// ⚠️ 输出警告:RabbitMQ 不支持事务消息,降级为普通消息
// ⚠️ 消息立即投递,EndTransactionAsync(commit:false) 无法回滚!

// 生产环境必须提前检查能力:
if (!provider.SupportsTransactionMessages)
{
    // 方案 A:改用 RocketMQ
    // 方案 B:本地事务表 + 补偿机制
}
顺序消息(RocketMQ 原生)
// 相同 orderKey 的消息路由到同一队列,保证消费顺序
await publisher.PublishOrderAsync("OrderTopic", order, orderKey: order.OrderId);
延迟消息(RocketMQ 原生 18 级)
// 延迟级别 1-18 对应:1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h
await publisher.PublishDelayAsync("OrderTopic", order, delayLevel: 3);
Request-Reply(RocketMQ 原生)
// 发送请求并等待响应(同步 RPC 模式)
var response = await publisher.RequestAsync<OrderQuery, OrderResult>(
    new OrderQuery("ORD-001"), timeout: 5000);
能力检查总览
provider.SupportsTransactionMessages   // 事务消息(RabbitMQ: false, RocketMQ: true)
provider.SupportsDelayedMessages       // 延迟消息(RabbitMQ: false, RocketMQ: true)
provider.SupportsRequestReply          // Request-Reply(RabbitMQ: false, RocketMQ: true)
provider.SupportsSql92Filtering        // SQL92 过滤(RabbitMQ: false, RocketMQ: true)
provider.SupportsOrderedMessages       // 顺序消息(RabbitMQ: false, RocketMQ: true)

3. 消息推送(钉钉/企业微信/飞书)

注册方式
.UsePush(push => push
    // 钉钉机器人
    .AddDingTalk(d => d
        .FromConfiguration("DingTalk")           // 从配置加载
        // 或:.AddGroup("告警群", "token", "secret") // 直接添加群组
    )
    // 企业微信机器人
    .AddWeChat(w => w
        .FromConfiguration("WeChatPush")
        // 或:.AddGroup("运维群", "key")
    )
    // 飞书机器人
    .AddLark(l => l
        .FromConfiguration("LarkPush")
        // 或:.AddGroup("开发群", "token", "secret")
    )
    // 钉钉应用推送(工作通知)
    .AddDingApp(a => a
        .FromConfiguration("DingApp")
        // 或:.WithCredentials("appKey", "appSecret", "agentId")
    )
)
业务使用
public class NotificationService
{
    private readonly IDingTalkApiFactory _dingFactory;
    private readonly IWeChatApiFactory _weChatFactory;
    private readonly ILarkApiFactory _larkFactory;

    public NotificationService(
        IDingTalkApiFactory dingFactory,
        IWeChatApiFactory weChatFactory,
        ILarkApiFactory larkFactory)
    {
        _dingFactory = dingFactory;
        _weChatFactory = weChatFactory;
        _larkFactory = larkFactory;
    }

    // 推送到默认群
    public async Task<bool> NotifyDingTalkAsync(string content)
    {
        var push = _dingFactory.GetPushService();
        return await push.PushAsync(content);
    }

    // 推送到指定群(按 Alias)
    public async Task<bool> NotifyDingTalkGroupAsync(string alias, string content)
    {
        var push = _dingFactory.GetPushService(alias);
        return await push.PushAsync(content);
    }

    // 使用消息格式化器
    public async Task<bool> NotifyWithFormatAsync()
    {
        var push = _dingFactory.GetPushService();
        return await push.PushAsync(format => format.Text("服务器 CPU 超过 90%"));
    }

    // 获取所有推送服务
    public async Task BroadcastAsync(string content)
    {
        foreach (var push in _dingFactory.GetAllPushService())
        {
            await push.PushAsync(content);
        }
    }
}

安全提示:推送失败日志已自动脱敏,不会将完整消息内容写入日志。


4. 加密(AES-GCM / RSA)

注册方式
// AES-GCM(推荐,默认)
.UseEncryption(enc => enc.UseAES(a => a
    .FromConfiguration("Encryption")        // 从配置加载
    // 或:.WithKeys("base64_key", "base64_iv")  // 直接配置
))

// RSA
.UseEncryption(enc => enc.UseRSA(r => r
    .FromConfiguration("Encryption")
    // 或:.WithKeys("public_key_xml", "private_key_xml")
))

// DES(已标记 Obsolete,仅用于历史数据解密,不推荐生产新用)
// .UseEncryption(enc => enc.UseDES(d => d.FromConfiguration("Encryption")))
AES 密钥要求
项目 要求
密钥格式 Base64 编码字符串
解码后长度 必须为 16 / 24 / 32 字节
加密算法 AES-GCM(认证加密)
密文格式 [LRD 0x02][nonce(12)][tag(16)][ciphertext]
压缩 加密前先 Deflate 压缩

重要:v7.11.0 起,AES 新加密统一使用 AES-GCM。密钥长度不合法时会抛 InvalidOperationException,不再静默截断/补零。

生成 AES 密钥
using System.Security.Cryptography;

var key = RandomNumberGenerator.GetBytes(32); // 256 位
var iv = RandomNumberGenerator.GetBytes(16);  // 备用 IV(仅 legacy 解密使用)
var keyBase64 = Convert.ToBase64String(key);
var ivBase64 = Convert.ToBase64String(iv);

// 写入 appsettings.json
// "Encryption": { "PublicKeyOrKey": "...", "PrivateKeyOrIV": "...", "EncryptType": "Aes" }
业务使用
public class SecureService
{
    private readonly IEncryptProvider _encryptor;

    public SecureService(IEncryptProvider encryptor) => _encryptor = encryptor;

    public byte[] Encrypt(string plainText)
    {
        var bytes = Encoding.UTF8.GetBytes(plainText);
        return _encryptor.Encryption(bytes);
    }

    public string Decrypt(byte[] cipher)
    {
        var bytes = _encryptor.Decryption(cipher);
        return Encoding.UTF8.GetString(bytes);
    }
}
加密配置示例
{
  "Encryption": {
    "PublicKeyOrKey": "AAECAwQFBgcICQoLDA0ODxAREhMUFRYXGBkaGxwdHh8=",
    "PrivateKeyOrIV": "AAECAwQFBgcICQoLDA0ODxAREhMUFRYXGBkaGxwdHh8=",
    "EncryptType": "Aes"
  }
}

5. HTTP 请求

注册方式
// RestSharp(默认)
.UseHttp(h => h.UseRestSharp().WithTimeout(TimeSpan.FromSeconds(30)))

// HttpClient
.UseHttp(h => h.UseNetHttp().WithTimeout(TimeSpan.FromSeconds(30)))
业务使用
public class ApiService
{
    private readonly IHttpFactory _httpFactory;

    public ApiService(IHttpFactory httpFactory) => _httpFactory = httpFactory;

    public async Task<MyResult?> GetDataAsync()
    {
        var factory = _httpFactory.CreateFactory("https://api.example.com");
        var response = await factory.CreateRequest("/api/v1/data")
            .AddHeader("Authorization", "Bearer token")
            .AddQueryParameter("page", "1")
            .GetAsync<MyResult>();
        return response.Content;
    }

    public async Task PostDataAsync(object payload)
    {
        var factory = _httpFactory.CreateFactory("https://api.example.com");
        await factory.CreateRequest("/api/v1/submit")
            .AddJsonBody(payload)
            .PostAsync();
    }
}
HTTP 重试策略
方法 默认重试 说明
GET / HEAD / OPTIONS / PUT / DELETE 3 次 幂等方法自动重试
POST / PATCH 不重试 防止重复提交
4xx(除 429) 不重试 客户端错误不可恢复
429 / 5xx 重试 服务端临时错误
超时/连接失败 重试 网络问题可恢复

如需对 POST 启用重试:

var request = factory.CreateRequest("/api/submit")
    .AddJsonBody(payload)
    .SetRetryCount(3)
    .AllowRetryNonIdempotent()  // 显式允许 POST 重试
    .PostAsync<MyResult>();

安全提示:HttpClient 缓存 key 已从 url.Host 改为 scheme://host:port,避免同 host 不同端口复用错误客户端。响应消息在读取内容后立即释放,避免连接池耗尽。


6. 日志(NLog / Log4Net)

注册方式
// NLog(默认)
.UseLogging(log => log.UseNLog())

// Log4Net
.UseLogging(log => log.UseLog4Net())

// 自定义日志配置
.UseLogging(log => log.UseNLog(builder =>
{
    builder.SetMinimumLevel(LogLevel.Information);
    builder.AddFilter<Microsoft.Hosting.Lifetime>("Microsoft", LogLevel.Warning);
}))

日志配置文件 nlog.config / log4net.config 优先从应用根目录加载,不存在时使用内置默认配置。


7. 微信公众号

注册方式
.UseWeChatOfficial(wx =>
{
    wx.FromConfiguration("WeChatOfficial");
    // 或:wx.WithSettings("appId", "appSecret", "token", "encodingAesKey?");
    // 可选:wx.UseMessageHandler<MyCustomHandler>();
})
微信公众号配置
{
  "WeChatOfficial": {
    "AppId": "your_app_id",
    "AppSecret": "your_app_secret",
    "Token": "your_token",
    "EncodingAesKey": null
  }
}
Controller 接入(安全 POST 入口)
[ApiController]
[Route("api/wechat")]
public class WeChatController : ControllerBase
{
    private readonly IWeChatSecureMessageHandler _handler;

    public WeChatController(IWeChatSecureMessageHandler handler)
        => _handler = handler;

    // GET:微信 URL 验证
    [HttpGet]
    public string Get(
        [FromQuery] string signature,
        [FromQuery] string timestamp,
        [FromQuery] string nonce,
        [FromQuery] string echostr)
    {
        var service = _handler.MsgOfficialService;
        return service.ValidateSignature(signature, timestamp, nonce, echostr, out var response)
            ? response : "fail";
    }

    // POST:微信消息回调(自动验签)
    [HttpPost]
    public async Task<string> Post(
        [FromQuery] string signature,
        [FromQuery] string timestamp,
        [FromQuery] string nonce)
    {
        using var reader = new StreamReader(Request.Body);
        var xml = await reader.ReadToEndAsync();
        // 先验签,再处理 XML
        return await _handler.HandleMessageAsync(xml, signature, timestamp, nonce);
    }
}

安全提示IWeChatSecureMessageHandler 在处理 XML 前会先校验微信签名。签名校验使用固定时间比较(CryptographicOperations.FixedTimeEquals),防止时序攻击。

事件订阅
public class WeChatEventService
{
    private readonly IWeChatEventHandler _events;

    public WeChatEventService(IWeChatEventHandler events)
    {
        _events = events;

        // 用户关注
        _events.OnUserSubscribed += async args =>
        {
            Console.WriteLine($"用户关注: {args.FromUser}");
            await Task.CompletedTask;
        };

        // 文本消息
        _events.OnTextMessageReceived += async args =>
        {
            Console.WriteLine($"收到文本: {args.Content}");
            await Task.CompletedTask;
        };
    }
}

8. 微信小程序

注册方式
.UseWeChatMiniProgram(wx =>
{
    wx.FromConfiguration("WeChatMiniProgram");
    // 或:wx.WithSettings("appId", "appSecret");
    // 多小程序:wx.AddProgram("alias", "appId", "appSecret");
})
配置
{
  "WeChatMiniProgram": {
    "AppId": "your_appid",
    "AppSecret": "your_secret",
    "BaseUrl": "https://api.weixin.qq.com"
  }
}

9. JSON 序列化

// Newtonsoft.Json(默认)
.UseJson(j => j.UseNewtonsoftJson(settings =>
{
    settings.DateFormatString = "yyyy-MM-dd HH:mm:ss";
    settings.NullValueHandling = NullValueHandling.Ignore;
}))

// System.Text.Json
.UseJson(j => j.UseTextJson(opts =>
{
    opts.Serialize = o => o.PropertyNamingPolicy = JsonNamingPolicy.CamelCase;
}))

安全最佳实践

1. 加密

  • 生产环境必须使用 AES-GCM,不要使用 DES
  • AES 密钥长度必须为 16/24/32 字节(推荐 32 字节 = AES-256)
  • 密钥不要硬编码在代码中,使用 appsettings.json 或环境变量
  • 定期轮换密钥

2. 日志脱敏

  • 推送失败日志已自动脱敏,不会记录完整消息内容
  • 微信 XML 日志已自动脱敏 ContentRecognitionFromUserName 等敏感字段
  • 如需自定义脱敏,使用 LogSanitizer 工具类

3. 微信公众号

  • POST 回调必须使用 IWeChatSecureMessageHandler 进行签名校验
  • 签名校验使用固定时间比较,防止时序攻击
  • EncodingAesKey 配置后可支持加密模式

4. HTTP 请求

  • POST/PATCH 默认不重试,防止重复提交
  • 响应消息在读取内容后立即释放,避免连接池耗尽
  • HttpClient 缓存按 scheme://host:port 隔离

5. RabbitMQ

  • 连接创建有 30 秒总超时,不会无限阻塞
  • 消费端 handler 超时 3 分钟后自动取消,避免资源泄漏
  • 看门狗永不永久停止,失败后进入低频持续重试

6. Redis

  • 同步 AddOrGetCacheItem 有 30 秒超时保护
  • 分布式锁自动续期有重入保护,防止续期任务堆积
  • 批量操作自动分批(每批 500 条),避免大命令阻塞 Redis

v7.10.8 变更与迁移

新功能

功能 说明
死信前缀白名单 DeadLetterPrefixes EncryptionPrefixes 同模式:队列名匹配前缀才创建/转发死信,不匹配的队列消费失败后直接丢弃("不重要的扔了就扔了")。builder(.DeadLetterPrefixes("important_", ...))、配置级、settings 级三级配置,RabbitMQ / RocketMQ 双引擎生效
RocketMQ 死信主题探针 死信订阅启动前自动发送 LORD_SERVICE_DLQ_PROBE 触发 %DLQ%{group} 主题创建,修复"Consumer 提前订阅不存在的 %DLQ% 主题后无法自动恢复"的问题(探针不会泄漏到主订阅,已实测)
RocketMQSubscriber 死信前缀过滤 SubscribeFailedAsync 与失败转发路径对齐 RocketMQReceive:消费组不匹配白名单时禁止订阅死信 / 失败直接丢弃(未配置前缀时行为与旧版完全一致)

重要修复

问题 修复
RocketMQFactory 死信前缀继承失效 括号嵌套错误(死信继承误嵌加密前缀 if 块内),配置 EncryptionPrefixesDeadLetterPrefixes 不再从 config 继承
RocketMQSubscriber 无视前缀白名单 消费失败超限后无条件转发 %DLQ%,现按前缀白名单丢弃
MQDeadLetterContext.Enabled 默认值 7.10.7 回归:默认 false 导致 EnableDeadLetter / UseDeadLetter() 静默不产生死信配置,超限消息被 broker 直接丢弃;恢复默认 true(8c81bee 审计修复)

迁移指南

  1. 死信前缀:默认(不配置)行为与旧版完全一致——全部队列启用死信。需要"不重要队列不建死信"时:

    .UseMQ(mq => mq.UseDeadLetter().DeadLetterPrefixes("important_", "critical_").UseRabbitMQ())
    

    或配置级:

    "RabbitMQ": { "DeadLetterPrefixes": "important_;critical_" }
    

    队列名以任一前缀开头(忽略大小写)才创建死信;不匹配的队列消费失败超限后直接丢弃消息。


v7.10.6 变更与迁移

新功能

功能 说明
PushMessage 统一消息模型 渠道无关推送,切换钉钉/飞书/企业微信零改动
IMessageBuilder 一致性接口 4 个消息类统一 Add 字典/实体/At 用户
消息加密开关 EnableEncryption / EncryptionPrefixes(消息级 + UseMQ 全局 + 配置级)
RocketMQ ACL 链式配置 WithAcl / WithCloudProvider / WithSsl
IMQProvider 自动注册 UseRocketMQ / UseRabbitMQ 已包含,无需单独调用
HTTP 幂等保护 POST 默认不重试 + AllowRetryNonIdempotent()
IHttpStreamResponse 释放 实现 IDisposable

重要修复

问题 修复
钉钉推送 @所有人逻辑写反 未指定 users 时不再误 @全员
飞书富文本消息必然抛异常 style 对象缺属性名 → 省略可选字段
Redis 信号量超时破坏 Wait 返回值检查 + 防 SemaphoreFullException
Redis ClearAsync 误清整库 无前缀无 pattern 时拒绝执行
Redis DateTime 时区偏移 保持 UTC Kind
DeepCopy 循环引用 StackOverflow AsyncLocal 深度计数,抛明确异常
NetHttp SendAsync 无限递归 修复调用自身
探针消息被消费无限重试 移除探针机制
看门狗停滞误判 连接活性检查 + 卡死检测

迁移指南

  1. IMQProvider 注册:原来 UseMQProvider() 单独调用可以删除,UseRocketMQ / UseRabbitMQ 已自动注册。

  2. 消息推送:新项目推荐直接使用 PushMessage(渠道无关);存量 DingTalkMessage 代码不受影响(API 兼容)。

  3. 加密:默认仍为 AES-GCM(安全)。无敏感信息的消息可用 EnableEncryption(false) 关闭加密。

  4. HTTP:POST 默认不重试(防重复提交)。需要重试时显式调用 .AllowRetryNonIdempotent()


v7.11.0 变更与迁移

安全增强

变更项 旧行为 新行为
默认加密 DES AES-GCM
AES 模式 CBC 无认证 GCM 认证加密
AES 密钥 静默截断/补零 严格校验 16/24/32 字节
推送日志 完整记录 content 自动脱敏
微信日志 完整记录 XML 自动脱敏敏感字段
微信 POST 无签名校验入口 新增 IWeChatSecureMessageHandler
签名比较 字符串相等 固定时间比较

稳定性增强

变更项 旧行为 新行为
MQ 连接创建 CancellationToken.None 30 秒内部超时
MQ 看门狗 20 次失败后永久停止 低频持续重试(60 秒一次)
MQ handler 超时 后台任务继续运行 取消后台任务
HTTP 响应 未释放 using var 确保释放
HTTP POST 默认重试 3 次 默认不重试
HTTP 缓存 key url.Host scheme://host:port

性能增强

变更项 旧行为 新行为
Redis RemoveBatch 一次性删除 分批 500 条
Redis SetBatch 一次性写入 分批 500 条
Redis 锁续期 可能重入堆积 重入保护
Redis singleflight 无超时 30 秒超时

架构改进

变更项 旧行为 新行为
Builder 注册 8 处 BuildServiceProvider 全部移除,改用 BindConfiguration 或懒加载

迁移指南

  1. DES → AES:如果旧数据是用 DES 加密的,解密时仍可使用 UseDES。新数据必须使用 UseAES

  2. 推送配置FromConfiguration() 的配置在首次调用 GetPushService(sectionName) 时懒加载,行为与之前一致。

  3. HTTP POST 重试:如果业务依赖 POST 自动重试,调用 .AllowRetryNonIdempotent() 显式开启。

  4. 微信 POST:将 Controller 中的 HandleMessageAsync(xml) 替换为 HandleMessageAsync(xml, signature, timestamp, nonce)

  5. RabbitMQ → RocketMQ 迁移

    无需修改的代码(使用通用接口):

    // ✅ 以下代码从 RabbitMQ 切换到 RocketMQ 无需修改
    await hub.PublishAsync(new OrderCreated("ORD-001", 99.9m));
    await hub.SubscribeAsync<OrderCreated>(async msg => { return true; });
    await hub.SubscribeFailedAsync<OrderCreated>(async msg => { return true; });
    

    需要修改的代码(仅注册部分):

    // ❌ 旧代码(RabbitMQ)
    .UseMQ(mq => mq.UseRabbitMQ(r => r.FromConfiguration("RabbitMQ")))
    
    // ✅ 新代码(RocketMQ)
    .UseMQ(mq => mq.UseRocketMQ(r => r.FromConfiguration("RocketMQ")))
    

    字段映射关系: | RabbitMQ 概念 | RocketMQ 概念 | 自动映射 | |--------------|--------------|---------| | Exchange | Topic | ✅ 自动 | | Queue | ConsumerGroup | ✅ 自动 | | RouteKey | Tag | ✅ 自动 | | DLX/DLQ | %DLQ%{ConsumerGroup} | ✅ 自动 |

    注意事项

    • 同一容器不能同时启用 RabbitMQ 和 RocketMQ
    • RocketMQ 需要在 broker 上预先创建 Topic(或启用自动创建)
    • 高级功能(事务/顺序/延迟消息)需要手动转换接口类型
Product Compatible and additional computed target framework versions.
.NET net6.0 is compatible.  net6.0-android was computed.  net6.0-ios was computed.  net6.0-maccatalyst was computed.  net6.0-macos was computed.  net6.0-tvos was computed.  net6.0-windows was computed.  net7.0 was computed.  net7.0-android was computed.  net7.0-ios was computed.  net7.0-maccatalyst was computed.  net7.0-macos was computed.  net7.0-tvos was computed.  net7.0-windows was computed.  net8.0 is compatible.  net8.0-android was computed.  net8.0-browser was computed.  net8.0-ios was computed.  net8.0-maccatalyst was computed.  net8.0-macos was computed.  net8.0-tvos was computed.  net8.0-windows was computed.  net9.0 is compatible.  net9.0-android was computed.  net9.0-browser was computed.  net9.0-ios was computed.  net9.0-maccatalyst was computed.  net9.0-macos was computed.  net9.0-tvos was computed.  net9.0-windows was computed.  net10.0 is compatible.  net10.0-android was computed.  net10.0-browser was computed.  net10.0-ios was computed.  net10.0-maccatalyst was computed.  net10.0-macos was computed.  net10.0-tvos was computed.  net10.0-windows was computed. 
Compatible target framework(s)
Included target framework(s) (in package)
Learn more about Target Frameworks and .NET Standard.

NuGet packages

This package is not used by any NuGet packages.

GitHub repositories

This package is not used by any popular GitHub repositories.

Version Downloads Last Updated
7.10.8 40 8/11/2026
7.10.7 94 8/6/2026
7.10.6 95 8/5/2026
7.10.5 92 8/4/2026
7.10.1 103 8/1/2026
7.10.0 99 7/30/2026
7.0.8 114 7/29/2026
7.0.7 103 7/29/2026
7.0.6 103 7/16/2026
7.0.5 118 4/24/2026
7.0.3 145 2/10/2026
7.0.2 136 2/6/2026
7.0.1 138 2/4/2026
7.0.0 150 1/26/2026
6.2.0 146 1/23/2026
6.0.9 142 1/12/2026
6.0.8 313 11/30/2025
Loading failed

7.10.7:
- 国密算法支持(SM4 对称 + SM2 非对称),基于 BouncyCastle
- 多套加密并行注册(UseAES+UseSM4+UseSM2),IEncryptProviderFactory 按名称解析
- 修复多套并行时 EncryptOption 单例覆盖(每个 Provider 独立配置实例)
- 注册风格与 MQ 一致:UseEncryption(s => s.UseSM4(m => m.FromConfiguration()))
- SM4: CBC+随机IV+Magic头+DoS防护;SM2: 单次上限检查+双密钥格式
- 新增 11 个国密单元测试(加解密往返/随机IV/大数据/非法密钥/超限)

7.10.6:
- 【新功能】统一消息模型 PushMessage + 渠道适配器(钉钉/飞书/企业微信)
 - 业务代码用 PushMessage 构建消息(渠道无关),推送时自动转换为各渠道 webhook JSON
 - 切换渠道(钉钉→飞书→企业微信)业务代码零改动,只换 AddDingTalk/AddLark/AddWeChat 注册
 - IApiPush 新增 Push(PushMessage)/PushAsync(PushMessage) 统一入口(默认接口实现)
 - 支持 Markdown/Text/Link/ActionCard/FeedCard 五种消息类型
- 【新功能】消息构建器统一接口 IMessageBuilder{T}(一致性)
 - PushMessage/DingTalkMessage/LarkMessage/WeChatMessage 均支持 Add 字典、Add 实体、At 用户
 - Add(IDictionary) 批量键值对、Add(TEntity) 属性自动转键值对
 - 修复泛型重载决议问题(传入 Dictionary 时运行时分流,避免被当实体反射)
- 【根本性改进】移除 RocketMQ 探针消息机制(LORD_SERVICE_DLQ_PROBE)
 - 根因:探针用于创建 %DLQ% 主题,但 RocketMQ Consumer 会自动订阅 %RETRY%/%DLQ% 队列,探针被业务消费 → 解密失败 → 无限重试
 - 验证:Consumer 订阅不存在的 Topic 时客户端会重试等待,失败消息转发时 Producer 发布(autoCreateTopic)自动创建 Topic,探针完全不需要
 - 保留存量探针跳过检查(兼容清理旧版本积压)
- 【看门狗】修复消费停滞检测误判:生产端正常无消息时被误判"停滞"导致每 5 分钟反复重建 consumer
 - 主检查改为 consumer 连接活性(Active/Disposed),停滞检测改为 handler 卡死检测(处理中标记 _processingTick,30 分钟阈值)
- 【看门狗】修复 Dispose 竞态:等待恢复任务结束再释放锁(避免 ObjectDisposedException)
- 【语义修复】Consumer.Tags 恢复 null(RocketMQ 协议缺省=不过滤),移除 Array.Empty 赋值(空数组可能导致过滤语义变化)
- 【新功能】消息级加密开关(IMQSettings.EnableEncryption + EncryptionPrefixes 前缀白名单)
 - 无需敏感信息的消息可关闭加密(提升性能、避免加解密配置不匹配)
 - 支持按队列名前缀白名单只加密部分队列(如 "secret_;private_")
 - 生产端与消费端使用相同队列名时自动保持一致的加密策略
 - RocketMQConfig/RabbitMQConfig 提供全局默认值(Topic 驱动接口场景)
- 【重构】提取 RocketMQMessageHelper(探针识别/重试退避)与 MQEncryptionHelper(加密策略),消除重复代码
- 【新功能】UseMQ 全局加密配置(引擎无关:RabbitMQ/RocketMQ/未来 Kafka)
 - mq.EnableEncryption(false) 全局关闭加密;mq.EncryptionPrefixes("secret_", "private_") 白名单前缀
 - 通过 PostConfigureAll 应用到所有已注册引擎,Factory 创建 settings 时自动继承
 - RabbitMQ 生产/消费链路已接入加密开关(RabbitMQPush.ApplyEncryption / RabbitMQReceive.GetMessage)
 - 优先级:IMQSettings(消息级) > UseMQ 全局配置 > appsettings 配置
- 【新功能】RocketMQ Builder 新增 ACL/云厂商/SSL 链式配置
 - WithAcl(accessKey, secretKey) Apache ACL 认证
 - WithCloudProvider(Aliyun/Huawei/Tencent, ak, sk, extra) 云厂商托管
 - WithSsl(true) SSL/TLS 开关
- 【测试】新增 MQBuilder 全局加密配置测试 + RocketMQ Builder ACL 测试 + 全部 MQ 相关测试通过(42 个)
- 【全面审查修复】Redis/加密/推送格式/工具类 30+ 个 bug
 - 修复 DingTextFormat @所有人逻辑写反(未指定 users 时不再误 @全员)
 - 修复 LarkTextFormat 富文本消息必然抛异常(style 对象缺属性名)
 - 修复 Redis AddOrGetCacheItem 信号量超时破坏(Wait 返回值检查,防 SemaphoreFullException)
 - 修复 Redis ClearAsync 误清整库(无前缀无 pattern 时拒绝清理,与同步版一致)
 - 修复 Redis DateTime 时区偏移(保持 UTC Kind,不再 ToLocalTime)
 - 修复 object 类型缓存读写不对称(读侧先还原原始值再走 JSON)
 - 修复 RedisDistributedLock 前缀双冒号问题(TrimEnd(':'))
 - AES-GCM 解密增加认证前长度上限(防内存 DoS)+ 密钥解码缓存(性能)
 - RSA 加密增加单次上限检查(190 字节明确报错,提示用 AES)
 - DES 启用时输出安全警告(56 位密钥可破解,仅建议兼容旧数据)
 - 修复 DeepCopy 循环引用防护失效(AsyncLocal 深度计数,循环引用抛明确异常而非 StackOverflow)
 - 修复 MemoryCache.Increment 类型判断(int/uint 等数值类型也能正确递增)
 - 修复 TypeConvertUtil 数组/object 分支畸形 JSON 抛异常(统一返回 default)
 - 修复 LarkTextFormat 卡片 config/header 序列化为字符串(改为对象结构)
- 【HTTP 类库审查修复】性能/稳定性/易用性全面优化
 - 修复 NetHttp SendAsync 无限递归 bug(调用自身导致死循环,单元测试捕获)
 - 修复 NetHttp 流式响应生命周期 bug(response 提前 Dispose 导致流关闭)
 - RestHttp 与 NetHttp 重试策略对齐:POST 默认不重试(防重复提交)+ 4xx 不重试(429 限流除外)
 - IHttpRequest 新增 AllowRetryNonIdempotent(显式允许非幂等重试)
 - IHttpStreamResponse 实现 IDisposable(流使用后可释放)
 - AddJsonBody(string) 语义统一:原始 JSON 原样发送(修复二次转义问题)

7.10.5:
- IMQProvider 性能优化(懒加载缓存 Hub/Factory/IJsonFormat,避免重复 DI 解析)
- 修复 RocketMQSettings 默认值点号问题(RocketMQ Topic/Group 不允许点号,改为下划线)
- 文档新增 IMQProvider 使用指南(一站式入口 + 能力检查 + MQTT 扩展指引)

7.10.4:
- IMQProvider 新增细粒度能力检查(SupportsTransactionMessages/SupportsDelayedMessages/SupportsRequestReply/SupportsSql92Filtering/SupportsOrderedMessages)
- IMQProvider 新增便捷方法(GetPublisher/GetSubscriber/GetHub/GetFactory),一站式获取所有 MQ 服务
- 修复 RabbitMQ 事务消息"静默降级"问题(现在可通过 SupportsTransactionMessages 提前检查)
- 修复 RabbitMQ 延迟消息"静默降级"问题(现在可通过 SupportsDelayedMessages 提前检查)
- 新增 RocketMQCloudHelper 共享帮助类,减少 160+ 行重复代码

7.10.3:
- 修复 RocketMQPublisher 缺少 Producer 活性检查问题(添加 IsProducerAlive 双检锁,与 RocketMQPush 对齐)
- 修复 EnsureFailedTopicAsync 中 Producer 资源泄漏(改用 finally 确保释放)
- 抽取 ConfigureCloudProvider 到共享帮助类 RocketMQCloudHelper(减少 160 行重复代码)
- 优化 CleanupExpiredRetryEntries 避免 LINQ 分配(改用 foreach 直接遍历)
- 文档版本同步更新至 v7.10.3

7.10.2:
- 修复 RocketMQ 模块 116 个 nullable reference 警告(提升类库质量)
- 修复 _retryCounts 字典无限增长问题(添加基于时间的自动清理机制)
- 修复 _failedProducer 资源泄漏问题(Dispose 时正确释放)
- 修复 _failedProducer 并发初始化竞态条件(双检锁保护)
- 优化 typeNameCache 缓存淘汰策略(达到上限时清空重建,确保新类型能被缓存)
- 修复 Tags 属性 null 赋值问题(改用 Array.Empty 兼容 C# 10)
- 修复 Type.FullName 可能为 null 的问题(添加回退值)
- 修复 message.Tags/Keys 可能为 null 的问题(添加默认值)

7.10.0:
- AES 加密改用随机 IV(每次加密生成新 IV,向后兼容旧数据)
- 新增 Polly 8 弹性策略(LordResilience),MQ 连接重试改为指数退避+抖动
- HTTP 重试改为指数退避+随机抖动(RestSharp + HttpClient 双实现)
- 移除 RabbitMQFactory/RabbitMQService 终结器(消除 GC 线程死锁风险)
- _timedOutTags 添加基于时间的自动淘汰(心跳中清理超过 5 分钟的泄漏条目)
- 修复 LarkPush.PushAsync 空路径 bug(异步飞书推送发送到错误端点)
- 修复 DingTalkPush/LarkPush 并发竞态条件(GetAuthentication 修改共享对象)
- 修复 WeChatPush/LarkPush 死代码和 null 检查不一致

7.0.8:
- IMQHub 新增显式名称参数重载(PublishAsync/SubscribeAsync/SubscribeDeadLetterAsync)
- 新增 MQNameAttribute,支持为消息类型指定自定义名称前缀
- IMQHub 中文 XML 注释
- RabbitMQ Prefix 前缀隔离(多系统共用 broker 时区分队列名)
- RabbitMQ 消费端看门狗重写(心跳检测+无限重试+指数退避+Channel 重建)
- 生产端 RabbitMQPush 新增 Channel 自动恢复机制
- 死信队列(DLX)纳入看门狗保护
- RSA 改为标准保密模式(加密用公钥,解密用私钥)
- 默认加密算法从 DES 改为 AES
- 修复缓存 AddOrGetCacheItem 默认值误判问题
- 修复服务工厂缓存 token 污染问题
- 修复消费端异常消息无限 requeue 问题
- 移除 FreeRedis 支持(聚焦 StackExchange.Redis)