Lord.Service 7.10.6 License Info

Lord.Service 7.10.6

LordService 完整使用文档

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

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

版本:v7.10.5 | 目标框架: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)

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 的事务消息实际是普通消息)。

注册方式
.UseMQ(mq => mq
    // 方式一:直接指定引擎(同时注册 IMQProvider)
    .UseMQProvider(MQProviderType.RocketMQ)
    // 或:.UseMQProvider(MQProviderType.RabbitMQ)

    // 方式二:先配置引擎,再注册 Provider(可自定义配置)
    .UseRocketMQ(r => r.FromConfiguration("RocketMQ"))
    .UseMQProvider(MQProviderType.RocketMQ)
)
一站式获取所有服务
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. 注册时一行切换,业务代码零改动

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.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(或启用自动创建)
    • 高级功能(事务/顺序/延迟消息)需要手动转换接口类型