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
<PackageReference Include="mgzhenhong.ASW2.MessageQueue.RedisStorage" Version="2.1.11" />
<PackageVersion Include="mgzhenhong.ASW2.MessageQueue.RedisStorage" Version="2.1.11" />
<PackageReference Include="mgzhenhong.ASW2.MessageQueue.RedisStorage" />
paket add mgzhenhong.ASW2.MessageQueue.RedisStorage --version 2.1.11
#r "nuget: mgzhenhong.ASW2.MessageQueue.RedisStorage, 2.1.11"
#:package mgzhenhong.ASW2.MessageQueue.RedisStorage@2.1.11
#addin nuget:?package=mgzhenhong.ASW2.MessageQueue.RedisStorage&version=2.1.11
#tool nuget:?package=mgzhenhong.ASW2.MessageQueue.RedisStorage&version=2.1.11
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 | Versions 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. |
-
net10.0
- FreeRedis (>= 1.5.5)
- mgzhenhong.ASW2.MessageQueue (>= 2.1.11)
NuGet packages
This package is not used by any NuGet packages.
GitHub repositories
This package is not used by any popular GitHub repositories.