TierQueue 使用指南
TierQueue = 内核 Ring + 消费组 + 死信 + 延迟 + 幂等 + Session 恰好一次的组合产品 (tc-tier-queue-spec)。数据真相源 = RingOfQueueKey;每组独立持久位点域(VersionedMetadata 版本链原子提交);投递簿记 pending 纯内存(崩溃丢失 = at-least-once 重放窗口)。
快速开始
var fs = TierFs.New("memory:"); // 或 TierFs.New("local:///path/to/volume")
await using var q = await new TierQueueBuilder(fs, new TierQueueOptions
{
QueueName = "orders",
PageSize = 64 << 10, // 64KB 页(消息粒度小,小页 = 细粒度环绕)
MemorySize = 8 << 20, // 8MB 内存(测试/dev;磁盘大积压按需调大)
}).StartAsync();
// 入队
var r = await q.EnqueueAsync(new byte[] { 1, 2, 3 }, default);
await q.WaitForDurableAsync(r.Address, default); // 等待/推进落盘水位覆盖该消息
// 消费($default 组)
var batch = await q.DequeueAsync(10, default);
foreach (var d in batch)
{
// 处理 d.Payload ...
await q.AckAsync([d.Address], default);
}
核心语义
确认契约(spec §6.1)
- AckAsync 返回 = 位点已持久(数据先于位点 flush → 原子提交 → 物化)。此后断电不重投。
- 未 ack 崩溃 = at-least-once 重投(pending 不持久,重放从 Cursor 起)。
- 乱序 ack = Skip 表(非前缀确认进 Skip 持久,重放跳过已确认空洞)。Skip 表上限 =
MaxInFlight,超出抛InvalidOperationException(消费侧背压)。
投递令牌与 fencing(spec §6.2)
var consumerA = q.CreateConsumerAsync("$default");
var consumerB = q.CreateConsumerAsync("$default");
var got = await consumerA.DequeueAsync(1, default);
var receipt = new DeliveryReceipt(got[0].Address, got[0].Token);
// B 抢占(minIdle=0 → 立即改派,epoch +1)
var claimed = await consumerB.ClaimAsync(new ClaimFilter(TimeSpan.Zero), default);
// A 迟交(旧 token)→ StaleDeliveryException(不双计)
await consumerA.AckAsync([receipt], default); // 抛
组级 GroupEpoch 单调递增(claim/抢占/复位 +1)。旧 token 的 Ack 被确定性拒绝。
消费组
var g = await q.CreateGroupAsync(new GroupOptions
{
Name = "payments",
StartAt = GroupStartAt.Earliest,
VisibilityTimeout = TimeSpan.FromSeconds(60),
MaxRedeliveries = 16,
RetryBackoff = n => TimeSpan.FromSeconds(Math.Pow(2, n)), // 指数退避
}, default);
var consumer = q.CreateConsumerAsync("payments");
- 每组独立位点(Cursor/Skip/Epoch/Fenced),持久于
{queue}.group.{name}。 - 组间天然并行;组内 Dequeue/Ack/Claim/Nack 串行(SemaphoreSlim)。
- 删组 → 截断下界抬升(最慢组不再钉住数据头)。
可见性超时与重试退避
VisibilityTimeout:在途未 ack 的重投等待。超时未 ack → 计数 +1 重投。
RetryBackoff(spec §7.1):可配置重投延迟。null/Zero = 立即重投(回队 ReadPoint 回卷);
否则保留在 pending,NextDeliverableTick 推到未来,到期由 ExpireOverdueAsync 回收重投。
// 指数退避:2s, 4s, 8s, ...
RetryBackoff = n => TimeSpan.FromSeconds(Math.Pow(2, n))
// 固定退避 5s
RetryBackoff = _ => TimeSpan.FromSeconds(5)
崩溃丢失 pending → 退避失效、立即重投(at-least-once 兜底,可接受)。
达 MaxRedeliveries → 死信路由(见下)。
死信队列(spec §7)
默认开启(独立 TierQueue 实例 {queue}.dlq)。达重投上限 → DLQ envelope 写入(durable 先行)
→ 原位终结(skip 标记)。
// DLQ 恰好一条 + 溯源字段
var dlqBatch = await q.DeadLetterQueue!.DequeueAsync(10, default);
DeadLetterEnvelope.TryUnwrap(dlqBatch[0].Payload, out var srcAddr, out var srcGroup,
out var finalCount, out var payload);
// 回放:死信还原原 payload 入目标队列
await q.DeadLetterQueue!.ReplayAsync(dlqBatch[0].Address, targetQueue, default);
// 或全部回放
await q.DeadLetterQueue!.ReplayAllAsync(targetQueue, default);
关闭档:DeadLetter = null → 达上限仅推进位点 + 告警,不进 DLQ。
延迟队列(spec §4)
Delayed 配置默认开启(BTree 就绪索引)。DueTime = UTC Ticks,晚于 MaxDelay 拒绝
(DelayedOptions.MaxDelay 缺省 24 小时)。
long due = DateTime.UtcNow.Ticks + TimeSpan.FromMinutes(5).Ticks;
await q.EnqueueAsync(new EnqueueOptions { DueTime = due }, payload, default);
await q.EnqueueBatchAsync(new EnqueueOptions { ProducerId = 42, Seq = 100 }, payloads, default);
// 批量幂等:Seq 按批内序号递增,键为 (42, 100), (42, 101)...
// 取消(按地址 / 按区间)
await q.CancelDelayedAsync(CancelDelayedFilter.ByAddress(r.Address), default);
await q.CancelDelayedAsync(CancelDelayedFilter.ByRange(from, to), default);
// 窥视
var pending = await q.ListDelayedAsync(default);
var page = await q.ListDelayedAsync(new DelayedQuery(from, to, Offset: 100, Limit: 50), default);
- 同 dueTime 按地址 FIFO(索引键 tiebreaker = 地址)。
- 未就绪延迟消息钉住截断下界(
MinFloor含延迟索引最小存活地址)。 - 崩溃窗口:索引未持久 → 恢复对账从 Ring 重建(零丢失)。
幂等生产(spec §6.4 档三)
Idempotency 配置开启后,(ProducerId, Seq) 判重(HashIndex)。同 key 重发返回首次地址,
Ring 恰一条。
await q.EnqueueAsync(new EnqueueOptions { ProducerId = 42, Seq = 1 }, payload, default);
var dup = await q.EnqueueAsync(new EnqueueOptions { ProducerId = 42, Seq = 1 }, payload, default);
// dup.Address == 首次地址
恰好一次(spec §6.3 档二)
组位点与业务效果经 Session 2PC 同生共死。业务存储需实现 ITransactionParticipant。
var session = SessionManager.Create(fs, "order-domain", HangingResolution.ForwardCommit,
("biz", bizStore), ("queue.orders", q.GetGroupParticipant("orders")));
session.Initialize();
await session.WaitForReadyAsync(default);
using var ts = session.OpenSession("consumer");
ts.Stage(() => bizStore.Process(order)); // 业务效果 staged
await consumer.AckInRoundAsync(ts, [receipt], default); // 确认挂进回合
await ts.CommitAsync(); // 2PC 原子:业务 + 位点同落地
Abort = 组位点与业务效果同回退,消息重投(at-least-once 兜底)。
积压治理(spec §8.2)
MaxBacklogBytes + Backlog 策略:
| 模式 | 行为 |
|---|---|
Block(默认) |
阻塞生产者至回收跟上,有界等待 BacklogBlockTimeout,超时抛 QueueFullException |
Reject |
超限即抛 QueueFullException |
EvictOldest |
淘汰最慢组(fence → DataLostFencedException),显式 ResetGroupAsync 才能续 |
new TierQueueOptions
{
MaxBacklogBytes = 1 << 30, // 1GB 积压上限
Backlog = BacklogPolicy.Block,
BacklogBlockTimeout = TimeSpan.FromSeconds(10),
}
溢出 payload
OverflowPolicy.Enabled + MinOverflowSize:超大消息分离到溢出引擎,入队/消费/死信往返一致。
运维与诊断
// 组管理
var groups = await q.ListGroupsAsync(default); // 组名清单
await q.DeleteGroupAsync("payments", default); // 删组(截断下界抬升)
await q.ResetGroupAsync("payments", GroupStartAt.Earliest, LogicalAddress.Empty, default);
// EvictOldest 被 fence 后复位续用
// 统计与积压观测
QueueStats s = await q.GetStatsAsync(default); // 头/尾/durable 尾/积压字节/各组位点
foreach (var g in s.Groups) Console.WriteLine($"{g.Name}: cursor={g.Cursor}");
// 消费者侧诊断
var pending = await consumer.PendingAsync(default); // PEL 在途表(地址/持有者/空闲时长/重投计数)
await consumer.NackAsync([d.Address], default); // 显式 nack——回队按可见性/退避重投
// 存储面
await q.FlushAsync(default); // Ring 全量刷盘(durable 尾推进到 tail)
await q.TruncateAsync(uptoExclusive, default); // 物理截断(回收治理)
配置参考
TierQueueOptions
| 参数 | 默认 | 说明 |
|---|---|---|
QueueName |
"tc.queue" |
引擎空间隔离名(数据 Ring = {name}.ring,组注册表 = {name}.groups) |
PageSize |
1MB | Ring 页大小(消息小→小页) |
MemorySize |
64MB | Ring 内存容量(测试勿抄大缺省——mem 卷物化实占) |
SegmentGrowthLimit |
64MB | 数据 Ring 段生长上限 |
MaxInFlight |
4096 | 组在途窗口上限(Skip 表 + pending 容量) |
BufferedAckBatchSize |
1 | AckBufferedAsync 自动 durable 阈值;1 = 不缓冲,普通 AckAsync 始终立即 durable |
BufferedAckCommitInterval |
Infinite | AckBufferedAsync 最长驻留时间;正值时无后续 Ack 也会后台 flush |
GroupRegistryPayloadSize |
1MB | 组注册表持久容量;超过时建组明确失败,禁止静默截断 |
ColdReadRatio |
0.25 | 冷页缓存占比(冷读回源 ClockCache) |
OverflowPolicy |
Disabled |
溢出策略 |
MinOverflowSize |
0 | 溢出阈值(Enabled 时生效) |
MaxBacklogBytes |
null | 积压上限(null = 不限) |
Backlog |
Block |
治理模式 |
BacklogBlockTimeout |
10s | Block 有界等待上限 |
DeadLetter |
new() |
死信配置(null = 关闭) |
Delayed |
new() |
延迟队列配置(null = 不装配索引) |
Idempotency |
null | 幂等生产配置(null = 关闭) |
即时、非溢出、无积压上限的 EnqueueBatchAsync 会使用 Ring 独占写窗口,将 tail 锁竞争降至按页领取;
启用延迟、幂等、溢出或积压治理时保留逐条完整协调路径。延迟查询和范围取消通过 BTree lower-bound
seek 起扫,不再从索引头扫描。AckBufferedAsync 是显式吞吐档:未达到阈值时调用返回不代表 durable,
需要以 FlushAcknowledgementsAsync 建立持久化边界;崩溃前缓冲项按 at-least-once 重投。
DelayedOptions 支持 NodeSize、MinFillPercent 与持久化策略;IdempotencyOptions
支持 HashTableCapacity、OverflowPoolCapacity。需要完整覆盖各结构的持久化类型、缓存、
文件打开 hints、段元组 flush 周期或引擎优化参数时,在 builder 注入统一变换器:
var builder = new TierQueueBuilder(fs, options)
.WithStorageOptionsFactory((component, defaults) => component switch
{
"ring" => defaults.WithHints(FileOpenHints.WriteThrough),
"delay" => defaults.WithMetaTupleFlushInterval(TimeSpan.FromSeconds(1)),
_ => defaults,
});
await using var queue = await builder.StartAsync();
该变换器应用于 Ring、组注册表、每个组状态域、延迟索引、幂等索引及 DLQ。MemoryFileSystem
上的默认 Ring 自动禁用文件预分配。若需替换整个结构实例,可分别使用
WithRingFactory、WithRegistryMetaFactory、WithGroupMetaFactory、
WithDelayedIndexFactory 与 WithIdempotencyIndexFactory。
GroupOptions
| 参数 | 默认 | 说明 |
|---|---|---|
Name |
required | 组名 |
StartAt |
Earliest |
创建位点(Earliest/Latest/Address) |
VisibilityTimeout |
60s | 在途未 ack 重投等待 |
MaxRedeliveries |
16 | 重投上限(达上限死信) |
RetryBackoff |
null | 重试退避策略(null/Zero = 立即) |
验证矩阵(spec §13)
18 项验收矩阵全覆盖(见 tests/TC.Tier.Products.Tests/Queue/):
| # | 项 | 测试文件 |
|---|---|---|
| 1 | FIFO 保序 | TierQueueTests |
| 2 | 组独立 | TierQueueGroupTests |
| 3 | 乱序 ack Skip 持久 | TierQueueTests |
| 4 | ack 后断电不重投 | TierQueueTests |
| 5 | 崩溃窗口重投有界 | TierQueueGroupTests |
| 6 | fencing 迟交拒绝 | TierQueueGroupTests |
| 7 | 档二恰好一次 | TierQueueSessionRoundTests |
| 8 | 幂等生产 | TierQueueDelayedTests |
| 9 | 延迟就绪序 | TierQueueDelayedTests |
| 10 | 延迟崩溃对账 | TierQueueDelayedTests |
| 11 | 延迟取消解锚 | TierQueueDelayedTests |
| 12 | MaxDelay 锚定治理 | TierQueueDelayedTests |
| 13 | 可见性超时 + 退避 | TierQueueGroupTests |
| 14 | 死信 + Replay | TierQueueDeadLetterTests |
| 15 | 治理三模式 | TierQueueGroupTests |
| 16 | 截断守卫(含延迟索引) | TierQueueDelayedTests |
| 17 | 溢出 payload 死信往返 | TierQueueDeadLetterTests |
| 18 | 组状态原子崩溃 | TierQueueDurabilityTests |