Table of Contents

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 支持 NodeSizeMinFillPercent 与持久化策略;IdempotencyOptions 支持 HashTableCapacityOverflowPoolCapacity。需要完整覆盖各结构的持久化类型、缓存、 文件打开 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 自动禁用文件预分配。若需替换整个结构实例,可分别使用 WithRingFactoryWithRegistryMetaFactoryWithGroupMetaFactoryWithDelayedIndexFactoryWithIdempotencyIndexFactory

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