StackSaga Cassandra 恢复目录缓存与容量计算指南 (Directory Caching & Calculations Guide)
1. 核心架构概述 (Executive Overview)
在 StackSaga 的 Cassandra 响应式架构中,待恢复的事务是通过五级分层活跃目录 (5-Tier Hierarchical Active Directory) 进行索引发现的:
es_days_by_year ← 第 1 层 (Tier 1): 哪一日历天存在待处理工作?
└── es_recovery_windows_by_day ← 第 2 层 (Tier 2): 该日内的哪一个 UTC 分钟窗口?
└── es_instances_by_recovery_window ← 第 3 层 (Tier 3): 哪个 Pod 实例在该窗口内写入了数据?
└── es_buckets_by_instance ← 第 4 层 (Tier 4): 该 Pod 创建了哪些分桶 (Bucket) 分区?
└── es_recovery_transactions_by_instance ← 第 5 层 (Tier 5): 实际的事务指针记录。
1.1. 5 倍写入放大风险 (The 5x Write Amplification Hazard)
在未经优化的朴素实现中,每一次事务写入都需要依次执行 5 条数据库写语句(每个层级一条):
Tx #1 ──> 写 Tier 1, 写 Tier 2, 写 Tier 3, 写 Tier 4, 写 Tier 5 (5 次写入)
Tx #2 ──> 写 Tier 1, 写 Tier 2, 写 Tier 3, 写 Tier 4, 写 Tier 5 (5 次写入)
Tx #3 ──> 写 Tier 1, 写 Tier 2, 写 Tier 3, 写 Tier 4, 写 Tier 5 (5 次写入)
...
Tx #50000 ─> 写 Tier 1, 写 Tier 2, 写 Tier 3, 写 Tier 4, 写 Tier 5 (5 次写入)
在 100,000 笔事务的高并发负载下,数据库将承受 500,000 次写入——其中 400,000 次完全是多余的幂等覆写 (Redundant Upserts),反复写入那些早已存在完全相同目录键的元数据表中。 这会导致严重的性能衰退与集群压力:
-
CommitLog 和 Memtable 膨胀:被大量重复的元数据写入无端填充。
-
SSTable 压缩压力激增:在目录表上触发不必要的 Compaction 压缩风暴。
-
5 倍网络往返开销:在实时业务流量路径上产生昂贵的响应式 Publisher 调度消耗。
1.2. 解决方案:基于延迟确认的乐观幂等写入 (Optimistic Idempotent Upsert with Deferred Confirmation)
父目录层级(第 1 到 4 层)的作用纯粹是指示牌 (Signposts),用于指引重试 Worker 快速定位数据。 它们绝不存储事务业务载荷,也不存储单独的事务 ID。
为了同时规避写入放大并消除不必要的线程阻塞等待:
-
零线程阻塞 (Zero Thread Waiting, 纯非阻塞):
-
并发线程之间绝不相互阻塞等待。
-
如果某个分桶在 Cassandra 中尚未确认,到达的并发线程直接并行执行第 1–4 层的写入。
-
因为 Cassandra 的写操作具备天然幂等性 (Idempotent),针对第 1–4 层的多个并发 upsert 覆盖写入完全相同的键值,既不会报错也不会产生脏数据。
-
-
延迟确认机制 (Deferred Confirmation):
-
一旦任何线程成功完成向 Cassandra 写入第 1–4 层,即在内存中调用
markBucketConfirmed。 -
从那一刻起,后续所有到达的线程都会检测到
requiresParentDirectoryInsert == false,从而仅写入第 5 层 (Tier 5 ONLY)(快速直达路径 Fast Path)。
-
-
自动化容灾与自愈能力 (Self-Healing):
-
如果某个线程在写入第 1–4 层时遭遇数据库抖动或超时,它不会确认该分桶。
-
并发或随后的其他业务线程将自然接续重试写入,直到其中一个成功确认,从而绝对保证目录路径不会遗漏孤立。
-
[分桶 0 初始状态 (未确认)]
│
├─► Tx #1: 已确认? 否 ──► 写入 Tier 1–4 + Tier 5 ──► DB 成功! ──► markConfirmed(0) ✔️
│ │
├─► Tx #2: 已确认? 否 ──► 写入 Tier 1–4 + Tier 5 ──► DB 成功! │ (此时已变为已确认状态!)
│ (并发执行,零阻塞等待) │
│ ▼
│ [已激活确认状态 (Confirmed)]
│ │
├─► Tx #3: 已确认? 是 ────────────────────────────────────► 仅写入 Tier 5 (快速直达路径)
├─► Tx #4: 已确认? 是 ────────────────────────────────────► 仅写入 Tier 5 (快速直达路径)
└─► Tx #50,000: 已确认? 是 ───────────────────────────────► 仅写入 Tier 5 (快速直达路径)
对于一个容纳 50,000 笔事务的分桶,仅在最开始的数毫秒内有极少数并发线程(通常仅 2 到 5 个请求)执行了第 1–4 层的写入。 其余的 49,995+ 笔事务 全部零阻塞直接落入第 5 层,立刻实现了降低 ~80% 数据库写入流量的飞跃性性能提升。
2. 数学计算引擎 (Mathematical Calculation Engine)
缓存与跟踪管理器 (RecoveryDirectoryTracker) 在内存中对每笔事务执行三个独立的数学计算,耗时仅约 5 纳秒,全程无需任何分布式锁。
2.1. 计算 1:目标时间窗口索引与 UTC 午夜跨天回绕 (UTC Midnight Rollover)
事务调度被划分为粒度为 windowIntervalMinutes 的离散时间窗口(默认:1 分钟,即 24 小时一天共 1440 个时间窗口)。
每个 Saga 领域实体均可配置独立的延迟跨度:
-
重试 (
RETRY): 瞬态下游故障,目标指向紧邻的未来时间窗口W + delayWindows(默认:W + 1)。 -
崩溃恢复守护 (
RESTORE): 类似看门狗的看门事务记录,目标指向遥远的未来时间窗口W + delayWindows(默认:W + 600,即提前 10 小时)。
2.1.1. 公式与伪代码
totalWindowsPerDay = 1440 / windowIntervalMinutes
minutesSinceMidnight = (utcHour * 60) + utcMinute
currentWindowIndex = minutesSinceMidnight / windowIntervalMinutes
targetWindowIndex = currentWindowIndex + delayWindows
// 计算日历天偏移量 (精准处理跨越午夜零点的场景)
daysOffset = targetWindowIndex / totalWindowsPerDay
targetWindowInDay = targetWindowIndex % totalWindowsPerDay
targetDate = utcDate + daysOffset 天
minuteOfDay = targetWindowInDay * windowIntervalMinutes
2.1.2. 步骤解析
-
每日总窗口数 (
totalWindowsPerDay):-
若
windowIntervalMinutes = 1:totalWindowsPerDay = 1440 / 1 = 1440。 -
若
windowIntervalMinutes = 5:totalWindowsPerDay = 1440 / 5 = 288。
-
-
当前窗口索引 (
currentWindowIndex):-
计算自 UTC 午夜起流逝的分钟数 (
minutesSinceMidnight) 并除以windowIntervalMinutes。
-
-
目标窗口索引 (
targetWindowIndex):-
将领域配置的延迟跨度
delayWindows累加到currentWindowIndex。
-
-
UTC 午夜跨天 (
daysOffset与targetWindowInDay):-
若
targetWindowIndex < totalWindowsPerDay:写操作仍落在当天的日历日期内(daysOffset = 0)。 -
若
targetWindowIndex >= totalWindowsPerDay:写操作跨越了 UTC 午夜零点: -
daysOffset依据targetWindowIndex / totalWindowsPerDay自动递进日历天。 -
targetWindowInDay通过取模运算 (%) 平滑回绕到次日的起始分钟。
-
-
Cassandra 集群列键 (
minuteOfDay):-
映射回当天的自然分钟刻度 (0 到 1439),精确契合表结构的集群列定义 (
minute_of_day ASC)。
-
2.2. 计算 2:单调槽位分配与奇偶分桶隔离 (Odd/Even Bucket Segregation)
每个目标窗口在本地维护一个 JVM 级别的 AtomicLong 计数器。
当事务分配到该窗口时,它原子地获取一个单调递增的序列槽位 (1, 2, 3…)。
2.2.1. 公式与伪代码
slot = counter.incrementAndGet() // 单调递增: 1, 2, 3...
bucketSequence = (slot - 1) / bucketSize
// RETRY 使用偶数分桶索引: 0, 2, 4, 6...
retryBucketIndex = bucketSequence * 2
// RESTORE 使用奇数分桶索引: 1, 3, 5, 7...
restoreBucketIndex = (bucketSequence * 2) + 1
2.2.2. 为什么必须对奇数桶和偶数桶进行物理隔离?
| 核心维度 | 偶数分桶 (RETRY) |
奇数分桶 (RESTORE) |
|---|---|---|
分桶索引序列 |
0, 2, 4, 6… |
1, 3, 5, 7… |
触发生成条件 |
瞬态步骤调用失败(极低概率,如 HTTP 503) |
事务初始化启动(每一笔正常 Saga 事务均会生成) |
数据删除模式 |
分桶处理完毕后通过 O(1) 分区墓碑 (Partition Tombstone) 一次性整桶丢弃 |
事务正常成功后通过单行删除语句 (Individual Row Deletes) 移除 |
墓碑扫描状态 |
0% 单元格墓碑(Worker 拥有 100% 极速纯净扫描路径) |
承担高并发正常完成所产生的大量单元格墓碑 (Cell Tombstones) |
Cassandra 安全性保障 |
彻底杜绝扫描阶段的 |
将墓碑频繁删除的扰动完全隔离在正常重试扫描路径之外 |
2.3. 计算 3:内存边界跟踪判断 (requiresParentDirectoryInsert)
为了精准判定当前事务是否需要写入第 1–4 层,每个时间窗口状态使用并发线程安全集合记录已确认的分桶索引:
Set<Long> confirmedBuckets = ConcurrentHashMap.newKeySet();
当槽位分配计算得到 bucketIndex 时,立即执行检查:
boolean requiresParentDirectoryInsert = !confirmedBuckets.contains(bucketIndex);
-
若
confirmedBuckets.contains(bucketIndex)为false: -
该分桶尚未在 Cassandra 中确认持久化。
-
requiresParentDirectoryInsert判定为true。 -
当前线程执行第 1–4 层的幂等写入,在数据库确认成功后回调:
tracker.markBucketConfirmed(domain, type, allocation); -
若
confirmedBuckets.contains(bucketIndex)为true: -
该分桶的目录层级已被确认存在于 Cassandra 中。
-
requiresParentDirectoryInsert判定为false。 -
当前线程完全绕过第 1–4 层,直接仅写入第 5 层 (Tier 5 ONLY)。
3. 完整数字演算流程示例 (Numerical Walkthrough)
假定系统配置如下:
* bucketSize = 50,000
* windowIntervalMinutes = 1
* 当前 UTC 时间:10:00:00 UTC (minutesSinceMidnight = 600, currentWindow = 600)
* 操作类型:RETRY(延迟 = 1 个窗口 → 目标时间窗口为 601)
| 步骤 / 事务 | 槽位 (Slot) | 分桶序列 | 分桶索引 | requiresParentInsert |
Cassandra 数据库具体操作 |
|---|---|---|---|---|---|
Tx #1 |
1 |
(1 - 1) / 50000 = 0 |
0 * 2 = 0 |
|
5 层全量初始化写入: |
Tx #2 |
2 |
(2 - 1) / 50000 = 0 |
0 * 2 = 0 |
|
并发幂等写入: |
Tx #3 .. #50,000 |
3 .. 50000 |
0 |
0 |
|
极速直达路径写入: |
Tx #50,001 |
50001 |
(50001 - 1) / 50000 = 1 |
1 * 2 = 2 |
|
分桶滚动写入: |
Tx #50,002 |
50002 |
1 |
2 |
|
极速直达路径写入: |
4. 跟踪管理核心类的职责架构
所有恢复目录跟踪类均封装于 org.stacksaga.cassandra.recovery 包下:
| 类名 (Class Name) | 架构定位与核心职责 (Architectural Role) |
|---|---|
|
枚举类,区分 |
|
不可变 Record,封装 |
|
承载结果的 Record,包含目标 |
|
JVM 本地线程安全状态载体,持有特定时间窗口的 |
|
计算组件,负责将时间戳与领域配置自动转换为目标 |
|
核心跟踪服务,统一管理活跃 |
5. 反应式编程集成代码示例 (Project Reactor)
以下展示反应式服务层如何结合使用 RecoveryDirectoryTracker 与延迟确认机制:
package org.stacksaga.cassandra.recovery;
import lombok.RequiredArgsConstructor;
import org.springframework.stereotype.Service;
import reactor.core.publisher.Mono;
@Service
@RequiredArgsConstructor
public class ReactiveRecoveryService {
private final RecoveryDirectoryTracker tracker;
private final CassandraRecoveryDao recoveryDao;
/**
* 注册瞬态失败的重试事务
*/
public Mono<Void> registerRetry(String domain, String transactionId, byte[] payload) {
// 步骤 1: 非阻塞方式极速分配槽位 (~5 纳秒)
return tracker.allocateReactive(domain, RecoveryType.RETRY)
.flatMap(allocation -> {
// 步骤 2: 基于延迟确认标志位进行分支判断
if (allocation.requiresParentDirectoryInsert()) {
// 未确认状态: 并发写入第 1–4 层,DB 成功后置为已确认,随后追加第 5 层
return Mono.when(
recoveryDao.insertTier1DaysByYear(allocation.targetDate()),
recoveryDao.insertTier2WindowsByDay(allocation.targetDate(), allocation.targetMinuteOfDay()),
recoveryDao.insertTier3InstanceMarker(allocation.targetDate(), allocation.targetMinuteOfDay()),
recoveryDao.insertTier4BucketIndex(allocation.targetDate(), allocation.targetMinuteOfDay(), allocation.bucketIndex())
)
// 一旦 Cassandra 返回成功,立即在内存中翻转确认标志!
.doOnSuccess(unused -> tracker.markBucketConfirmed(domain, RecoveryType.RETRY, allocation))
.then(recoveryDao.insertTier5RetryTransaction(allocation, transactionId, payload));
} else {
// 快速通道 (已确认状态): 零阻塞等待,仅写入第 5 层
return recoveryDao.insertTier5RetryTransaction(allocation, transactionId, payload);
}
});
}
/**
* 在事务启动时注册崩溃守护看门狗记录
*/
public Mono<String> registerRestoreWatchdog(String domain, String transactionId) {
return tracker.allocateReactive(domain, RecoveryType.RESTORE)
.flatMap(allocation -> {
if (allocation.requiresParentDirectoryInsert()) {
return Mono.when(
recoveryDao.insertTier1DaysByYear(allocation.targetDate()),
recoveryDao.insertTier2WindowsByDay(allocation.targetDate(), allocation.targetMinuteOfDay()),
recoveryDao.insertTier3InstanceMarker(allocation.targetDate(), allocation.targetMinuteOfDay()),
recoveryDao.insertTier4BucketIndex(allocation.targetDate(), allocation.targetMinuteOfDay(), allocation.bucketIndex())
)
.doOnSuccess(unused -> tracker.markBucketConfirmed(domain, RecoveryType.RESTORE, allocation))
.then(recoveryDao.insertTier5RestoreWatchdog(allocation, transactionId))
.thenReturn(buildWatchdogPath(allocation, transactionId));
} else {
return recoveryDao.insertTier5RestoreWatchdog(allocation, transactionId)
.thenReturn(buildWatchdogPath(allocation, transactionId));
}
});
}
private String buildWatchdogPath(RecoveryAllocationResult alloc, String txId) {
return alloc.targetDate() + "/" + alloc.targetMinuteOfDay() + "/" + alloc.bucketIndex() + "/" + txId;
}
}
6. 内存防膨胀机制:过期时间窗口清理 (Stale Window Pruning)
由于容器化生产环境中的微服务 Pod 长期连续运行,JVM 内存中的跟踪状态绝不能无限膨胀。
标准节点 (Standard-Nodes) 永远只向未来视界写入数据:
-
Retry 写入 W + 1(或领域配置的延迟 W + D)。
-
Restore 写入 W + 600(或领域配置的延迟 W + D)。
一旦真实物理时间超过时间窗口 W,将永远没有任何 Pod 会再次向该窗口写入记录。
RecoveryDirectoryTracker 提供了定时修剪方法 pruneExpiredWindows:
@Scheduled(fixedDelay = 300_000) // 每 5 分钟执行一次
public void pruneStaleRecoveryWindows() {
LocalDateTime nowUtc = LocalDateTime.now(ZoneOffset.UTC);
int currentMinuteOfDay = nowUtc.getHour() * 60 + nowUtc.getMinute();
// 清理所有已经成为过去式的时间窗口状态
tracker.pruneExpiredWindows(nowUtc.toLocalDate(), currentMinuteOfDay);
}
这提供了坚实的不变量保证:
-
RecoveryDirectoryTracker在 JVM 堆内存中的活跃占用被严格约束在 50 KB 以内。 -
无论微服务连续平稳运行多少天或多少个月,均彻底实现零内存泄露。