mgzhenhong.ASW2.MessageQueue.RedisStorage 2.1.11

dotnet add package mgzhenhong.ASW2.MessageQueue.RedisStorage --version 2.1.11
                    
NuGet\Install-Package mgzhenhong.ASW2.MessageQueue.RedisStorage -Version 2.1.11
                    
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="mgzhenhong.ASW2.MessageQueue.RedisStorage" Version="2.1.11" />
                    
For projects that support PackageReference, copy this XML node into the project file to reference the package.
<PackageVersion Include="mgzhenhong.ASW2.MessageQueue.RedisStorage" Version="2.1.11" />
                    
Directory.Packages.props
<PackageReference Include="mgzhenhong.ASW2.MessageQueue.RedisStorage" />
                    
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 mgzhenhong.ASW2.MessageQueue.RedisStorage --version 2.1.11
                    
#r "nuget: mgzhenhong.ASW2.MessageQueue.RedisStorage, 2.1.11"
                    
#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 mgzhenhong.ASW2.MessageQueue.RedisStorage@2.1.11
                    
#: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=mgzhenhong.ASW2.MessageQueue.RedisStorage&version=2.1.11
                    
Install as a Cake Addin
#tool nuget:?package=mgzhenhong.ASW2.MessageQueue.RedisStorage&version=2.1.11
                    
Install as a Cake Tool

ASW.MessageQueue.RedisStorage

基于 FreeRedis (Redis 协议客户端) 的 IStorage 第二实现, 与 Sqlite 实现行为等价。

与 Sqlite 实现的语义差异

1. 持久性

维度 Sqlite Redis
默认持久化 写入即持久化到磁盘(WAL) 内存为主, 默认不持久化
进程崩溃 数据不丢 (WAL 已 fsync) 可能丢数据, 除非开启 AOF/RDB
部署要求 单文件 .db 单独 Redis 服务, 需要运维

强烈建议:生产环境必须开启 Redis AOF (appendonly yes, appendfsync everysec), 否则 server 进程崩溃可能丢失最近未落盘的消息。

2. Dispose 语义

维度 Sqlite Redis
Dispose storage 关闭 db 连接, 释放内存 仅停止分发, 不删除数据
Dispose factory 释放所有 storage 实例 关闭 ConnectionMultiplexer, 数据保留
清理数据 需手动 File.Delete(db) 需手动 DEL 相关 key

重要: Redis storage 的 Dispose 不会清理数据, 与 Sqlite 行为不同。 如需清理测试数据, 用 KEYS aswmq:* + DEL, 或调用 StorageFail 触发自定义清理逻辑。

3. 跨 key 原子性

维度 Sqlite Redis
原子性保障 BEGIN IMMEDIATE 事务 (跨进程也安全) Lua 脚本 (Redis 单线程执行)
进程内并发 AutoResetEvent 串行化 (快速路径) 依赖 Lua 原子性, 无需进程内锁
跨进程并发 Sqlite 文件锁 + 事务 (需小心 busy_timeout) Lua 天然原子

测试建议: 阶段二 Sqlite 实现的并发 fetch 测试在单实例内 被 AutoResetEvent 保护, 跨进程场景未覆盖; Redis 实现的 Test_Redis_ConcurrentFetch_DoesNot_Duplicate 用多 storage 实例共享 同一组 key, 真正验证 Lua 原子性, 是更严苛的并发验证。

4. 计数复杂度

方法 Sqlite Redis
GetStorageMessageCount SELECT COUNT(*) (O(N)) HLEN (O(1))
GetQueueMessageCount SELECT COUNT(*) (O(N)) HLEN (O(1))
GetQueueMessageCount(q, state) SELECT COUNT(*) WHERE state=? (O(N)) ZCARD (O(1))

不维护独立计数 Hash — 所有计数用 HLEN / ZCARD 实时计算, 保证与实际状态不漂移, 这是阶段三的关键设计决定。

5. Cluster 部署

所有键带 hash tag {aswmq:<storage>}, 同一 storage 的 key 落在同一 slot, 保证 Lua 脚本能跨 key 操作而不被路由到不同节点。

键空间

{aswmq:<s>}:meta             Hash   expirationPolicy (JSON) 等元数据
{aswmq:<s>}:msgs             Hash   全局存储消息体, field=Uid
{aswmq:<s>}:msgs:counter     String 存储消息 Uid 计数器 (INCR)
{aswmq:<s>}:msgs:bytime      ZSet   按 PublishTime 排序, member=Uid
{aswmq:<s>}:queues           Hash   队列元数据, field=queueName
{aswmq:<s>}:queues:watermark  Hash   队列分发水位, field=queueName, value=已入队最大 StorageMessageUid
                               (enqueue.lua 内原子更新, 供实例重建后跳过已入队消息)
{aswmq:<s>}:q:<q>:msgs       Hash   队列消息体, field=queueUid
{aswmq:<s>}:q:<q>:idle       ZSet   空闲消息, score=EnqueueTime.Ticks
{aswmq:<s>}:q:<q>:consuming  ZSet   消费中, score=下发时间.Ticks
{aswmq:<s>}:q:<q>:confirmed  ZSet   已确认, score=确认时间.Ticks
{aswmq:<s>}:q:<q>:dropped    ZSet   已丢弃, score=丢弃时间.Ticks
{aswmq:<s>}:q:<q>:counter    String 队列内 Uid 计数器 (INCR)
aswmq:storages               Set    全部 storage 名称 (供 GetStorages 枚举)

Lua 脚本原子性

脚本 操作 关键性
store INCR counter 分配全局 Uid + HSET msgs + ZADD bytime 存储原子性 (uid 分配 + 写入同脚本, 减少往返)
fetch ZRANGE idle + HGET msgs + 改写 consumer + ZREM idle + ZADD target 解决并发 fetch 时同一条消息被多次下发
respond ZSCORE consuming 校验 + ZREM + ZADD + HSET 消息体 (RejectedCount++ 等) 解决并发响应时状态被覆盖
enqueue HEXISTS 队列检查 + HGET watermark 水位检查 + INCR counter + HSET msgs + ZADD idle + HSET watermark 入队原子性 + 重启不重复 (水位检查/更新同脚本, 减少往返)
reset ZRANGE consuming + ZREM + ZADD idle + HSET 清空 consumer 解决重置时遗漏 Uid
clean_bytime ZRANGEBYSCORE + HDEL + ZREM 清理过期消息的原子性

性能优化 (阶段四补审)

单条消息路径 (发布 → 入队) 从 6-7 次命令往返优化到 2 次 (store.lua + enqueue.lua):

  • store.lua 合并 uid 分配 (INCR) 与存储
  • enqueue.lua 合并队列存在检查、水位检查、入队、水位持久化
  • 分发免 Redis 查询: StoreMessage 直接使用内存中的消息构造分发 (原每次分发全扫 bytime + HGET)

端到端实测 (5000 条, 测试环境): Sqlite ~400 条/秒, Redis (Garnet) ~1500 条/秒 (优化前 29 条/秒, 52 倍提升; Garnet 的 Lua 事务模式有锁开销, 真实 Redis 会更高)。

阶段三契约测试

RedisStorageContractTests 继承 StorageContractTestsBase, 自动获得 27 个契约测试。 无本地 Redis 时 Assert.Inconclusive 跳过所有测试, 不影响 CI 流程。

关键测试: Test_Redis_ConcurrentFetch_DoesNot_Duplicate - 5 个 RedisStorage 实例 共享同一 storage 名称, 并发 fetch 100 条消息, 验证合计 100 条且 Uid 不重复。 这是阶段二 Sqlite 单实例测试无法覆盖的跨实例并发场景。

使用方式

// 启动本地 Redis (Docker)
docker run -d --name asw-redis-test -p 6379:6379 redis:7-alpine

// 在代码中使用
var factory = new RedisStorageFactory("127.0.0.1:6379");
var result = factory.DeclearStorage("my-storage");
var storage = result.Data!;
storage.ExpirationPolicy = new MessageExpirationPolicy {
    ExpirationType = MessageExpirationType.AfterSeconds,
    ExpirationSeconds = 86400  // 1 天
};
storage.StartDistribute(filter, callback);
// ... 业务逻辑
storage.Dispose();
factory.Dispose();
Product Compatible and additional computed target framework versions.
.NET 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
2.1.11 92 8/29/2026
2.1.10 100 8/23/2026
2.1.9 92 8/19/2026
2.1.8 103 8/17/2026
2.1.7 110 7/31/2026