设计类算法与工程算法模式
#算法 · #限流 · #雪花ID · #令牌桶 · #CRC · #负载均衡
有些算法不解决"排序"或"搜索"问题,而是处理工程中实际遇到的系统设计难题:限流、ID生成、负载分配、数据校验。这些算法简单却关键,是后端工程师的必备工具。
1. 限流算法
限流是保护系统免受过载的第一道防线。
1.1 固定窗口计数器
go
// 固定窗口:key → 当前窗口剩余配额
// 问题:窗口边界突发流量可能翻倍
type FixedWindow struct {
mu sync.Mutex
limit int
count int
resetAt time.Time
}优点:简单,内存占用小。 缺点:窗口边界的"临界问题"——两个窗口交界处可能放过 2× 限额的流量。
1.2 滑动窗口日志
go
// 滑动窗口:保留最近窗口内的所有请求时间戳
// 内存开销与请求量成正比
type SlidingLog struct {
mu sync.Mutex
limit int
window time.Duration
timestamps []time.Time
}优点:精确,无临界问题。 缺点:内存消耗大(需保存所有时间戳)。
1.3 滑动窗口计数器
go
// 滑动窗口计数器:结合固定窗口 + 上一个窗口的加权值
// 当前窗口计数 + 上一个窗口计数 × (窗口重叠比例)
type SlidingCounter struct {
mu sync.Mutex
limit int
window time.Duration
prevCount int
curCount int
curWindowStart time.Time
}
func (s *SlidingCounter) Allow() bool {
s.mu.Lock()
defer s.mu.Unlock()
now := time.Now()
if now.Sub(s.curWindowStart) >= s.window {
s.prevCount = s.curCount
s.curCount = 0
s.curWindowStart = now
}
// 加权计算当前有效请求数
elapsed := now.Sub(s.curWindowStart)
overlap := float64(s.window-elapsed) / float64(s.window)
current := float64(s.curCount) + float64(s.prevCount)*overlap
if current < float64(s.limit) {
s.curCount++
return true
}
return false
}优点:精确度和内存的折中。
1.4 令牌桶
go
// 令牌桶:以固定速率生成令牌,请求消费令牌
// 允许"突发":桶中最多可存储 burst 个令牌
type TokenBucket struct {
rate float64 // 令牌生成速率(个/秒)
burst float64 // 桶容量
tokens float64 // 当前令牌数
lastUpdate time.Time
mu sync.Mutex
}
func (tb *TokenBucket) Allow() bool {
tb.mu.Lock()
defer tb.mu.Unlock()
now := time.Now()
elapsed := now.Sub(tb.lastUpdate).Seconds()
tb.tokens += elapsed * tb.rate
if tb.tokens > tb.burst {
tb.tokens = tb.burst
}
tb.lastUpdate = now
if tb.tokens >= 1 {
tb.tokens--
return true
}
return false
}应用:Go rate.Limiter、Nginx limit_req、AWS API Gateway。
1.5 漏桶
与令牌桶对称:请求进入队列,以固定速率"漏出"。只能平滑,不能突发。
1.6 对比
| 算法 | 突发支持 | 内存 | 适用场景 |
|---|---|---|---|
| 固定窗口 | ❌ | 极小 | 简单计数限制 |
| 滑动窗口日志 | ❌ | 高 | 精确限流 |
| 滑动窗口计数器 | ❌ | 低 | 通用 Redis 限流 |
| 令牌桶 | ✅ | 低 | API 网关 |
| 漏桶 | ❌ | 中 | 流量整形 |
2. 分布式 ID 生成
2.1 Snowflake(雪花算法)
64 位 ID 结构:
┌──────────────────────────────────────────────────────────────────┐
│ 1b │ 41b │ 10b │ 12b │
│未用│ 毫秒时间戳 │ 机器 ID │ 序列号 │
└──────────────────────────────────────────────────────────────────┘go
type Snowflake struct {
mu sync.Mutex
timestamp int64 // 上次生成 ID 的时间戳
workerID int64 // 机器 ID (0-1023)
sequence int64 // 毫秒内序列 (0-4095)
}
const epoch = 1640995200000 // 2022-01-01 自定义起始时间
func (s *Snowflake) NextID() int64 {
s.mu.Lock()
defer s.mu.Unlock()
now := time.Now().UnixMilli()
if now == s.timestamp {
s.sequence = (s.sequence + 1) & 0xFFF // 12bit
if s.sequence == 0 {
// 当前毫秒序列号用完,等待下一毫秒
for now <= s.timestamp {
now = time.Now().UnixMilli()
}
}
} else {
s.sequence = 0
}
s.timestamp = now
return ((now - epoch) << 22) | (s.workerID << 12) | s.sequence
}优点:递增、不依赖数据库、高性能。 缺点:依赖机器时钟,时钟回拨会导致 ID 重复。
2.2 号段模式(Leaf)
从数据库批量取一段 ID 号段缓存在本地,用完再取。
2.3 UUID v7
基于时间戳排序的 UUID,天然递增,适合数据库主键。
3. 负载均衡算法
3.1 轮询(Round Robin)
go
type RoundRobin struct {
nodes []string
current int32
}
func (r *RoundRobin) Next() string {
idx := atomic.AddInt32(&r.current, 1)
return r.nodes[int(idx)%len(r.nodes)]
}3.2 加权轮询(Weighted Round Robin)
Nginx 的平滑加权轮询:每次选 current_weight 最大的节点,选中后 current_weight -= total_weight。
go
type WeightedNode struct {
addr string
weight int
currentWeight int
}
func smoothWRR(nodes []*WeightedNode) string {
total := 0
var best *WeightedNode
for _, n := range nodes {
n.currentWeight += n.weight
total += n.weight
if best == nil || n.currentWeight > best.currentWeight {
best = n
}
}
best.currentWeight -= total
return best.addr
}3.3 最少连接
每次选连接数最少的节点。适合长连接场景。
3.4 一致性哈希
见 一致性Hash与Raft。
4. 短链算法(URL Shortener)
将长 URL 映射为短链(如 t.cn/xYz9Pq),核心是唯一 ID ↔ 短码 的双向转换。
4.1 设计目标
| 要求 | 说明 |
|---|---|
| 唯一性 | 不同长 URL 必须生成不同短码 |
| 一致性 | 同一长 URL 应返回同一短码(幂等) |
| 短码长度 | 通常 6~8 字符,空间足够 |
| 不可预测 | 不能被人遍历所有短链 |
4.2 Base62 编码
最常见的方案:数字 ID → Base62(0-9a-zA-Z,共 62 个字符)。
7 位 Base62 → 62^7 ≈ 3.5 万亿 条短链(绰绰有余)
6 位 Base62 → 62^6 ≈ 568 亿 (够大多数场景)go
const base62Chars = "0123456789ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz"
// 数字 ID → 短码
func encode(id int64) string {
if id == 0 {
return string(base62Chars[0])
}
var result []byte
for id > 0 {
result = append([]byte{base62Chars[id%62]}, result...)
id /= 62
}
return string(result)
}
// 短码 → 数字 ID
func decode(shortCode string) int64 {
var id int64
for _, ch := range shortCode {
var val int64
switch {
case '0' <= ch && ch <= '9':
val = int64(ch - '0')
case 'A' <= ch && ch <= 'Z':
val = int64(ch-'A') + 10
case 'a' <= ch && ch <= 'z':
val = int64(ch-'a') + 36
}
id = id*62 + val
}
return id
}4.3 ID 生成方案
| 方案 | 原理 | 优点 | 缺点 |
|---|---|---|---|
| 哈希法 | MD5/SHA1(long_url) 取前 N 位 | 相同 URL 天然幂等 | 碰撞需处理(加盐/重试) |
| 自增 ID | Redis INCR / DB AUTO_INCREMENT | 简单不碰撞 | 依赖外部服务 |
| Snowflake | 雪花算法生成唯一 ID | 分布式高性能 | 短码可预测 |
| 随机生成 | rand(62^N),检查冲突后重试 | 简单 | 冲突概率随使用量增长 |
4.4 推荐方案:自增 ID + Base62 + 预生成
go
// 生产级架构:
// 1. 长 URL → 哈希判重(Redis/布隆过滤器)→ 若已存在直接返回
// 2. 不存在 → 从发号器获取唯一 ID(Redis INCR / MySQL 自增 / Snowflake)
// 3. ID → Base62 编码 → 短码
// 4. 存储映射: shortCode → longURL(KV 存储,如 Redis)
// longURL_hash → shortCode(用于幂等查询)
type Shortener struct {
storage map[string]string // shortCode → longURL
reverse map[string]string // longURL_hash → shortCode
counter int64
mu sync.Mutex
}
func (s *Shortener) Shorten(longURL string) string {
hash := sha256Hex(longURL)[:12] // 取前 12 位作为幂等 key
s.mu.Lock()
defer s.mu.Unlock()
// 幂等检查
if code, ok := s.reverse[hash]; ok {
return code
}
// 生成新 ID
s.counter++
code := encode(s.counter)
s.storage[code] = longURL
s.reverse[hash] = code
return code
}
func (s *Shortener) Expand(shortCode string) (string, bool) {
s.mu.Lock()
defer s.mu.Unlock()
url, ok := s.storage[shortCode]
return url, ok
}4.5 不可预测性增强
若需要短码不可被遍历猜测,可采用以下方式:
- ID 混淆:对自增 ID 做对称加密(如 XOR 一个固定密钥、Feistel 网络)
- 随机码:生成随机 6-8 位 Base62 串,冲突后重试
- 带盐哈希:
MD5(longURL + secret_salt),从中间截取(而非开头)
5. 数据校验算法
4.1 CRC 循环冗余校验
用于检测数据传输/存储中的意外错误,非加密用途。
4.2 布隆过滤器
详见 特殊结构。
4.3 HyperLogLog
基数估计:用极小的内存估算集合中不同元素的数量(误差约 2%)。
应用:Redis PFADD/PFCOUNT,统计 UV。
6. 缓存淘汰策略
详见 LRU 与 LFU 缓存淘汰。
| 策略 | 核心思想 | 适用场景 |
|---|---|---|
| LRU | 淘汰最久未使用的 | 通用缓存 |
| LFU | 淘汰使用频率最低的 | 热点数据缓存 |
| FIFO | 淘汰最早进入的 | 简单场景 |
| TTL | 淘汰过期的 | 临时数据 |
| W-TinyLFU | LRU+LFU 混合 | 现代高性能缓存(Caffeine) |
7. 分布式锁
7.1 基于 Redis 的分布式锁
go
// 单节点 Redis 分布式锁(SET NX + 过期时间防死锁)
func acquireLock(rdb *redis.Client, key, value string, ttl time.Duration) bool {
ok, _ := rdb.SetNX(ctx, key, value, ttl).Result()
return ok
}
// 释放锁 — 必须校验 value(防止误删别人的锁)
func releaseLock(rdb *redis.Client, key, value string) bool {
script := `
if redis.call("GET", KEYS[1]) == ARGV[1] then
return redis.call("DEL", KEYS[1])
else
return 0
end
`
result, _ := rdb.Eval(ctx, script, []string{key}, value).Int()
return result == 1
}
// 为什么释放锁要校验 value?
// - 锁可能已过期,被其他客户端获取
// - 如果直接 DEL,会误删别人持有的锁!
// - Lua 脚本保证 GET + DEL 的原子性7.2 Redlock — 多节点 Redis 分布式锁
单节点 Redis 存在 SPOF。Redlock 通过在多个独立 Redis 节点上获取锁来提高可靠性:
Redlock 算法(N=5 个节点,法定人数 = N/2+1 = 3):
1. 获取当前时间戳 t1
2. 依次对 N 个 Redis 节点执行 SET NX(每个有独立超时)
3. 计算获取锁成功的节点数 count
4. 计算总耗时 = t2 - t1
5. 如果 count ≥ 法定人数 且 总耗时 < 锁 TTL → 获取成功
6. 如果失败 → 向所有节点发送 DEL
锁的有效时间 = TTL - 总耗时(扣除获取锁的通信开销)Redlock 的争议(Martin Kleppmann vs Salvatore Sanfilippo):
| 论点 | 说明 |
|---|---|
| 反对 | 时钟跳跃可能破坏安全性;GC pause 可能导致锁过期;没有 fencing token |
| 支持 | 工程实用主义:大多数场景下"够用";fencing token 是额外机制,不应强耦合到锁 |
生产建议:
- 高性能场景 → 单节点 Redis + Redlock 客户端库(redsync)
- 强一致性场景 → etcd/zookeeper(基于 Raft/ZAB 共识,天然线性一致)
- 最佳实践 → 加 fencing token(单调递增序列号),防止过期锁的破坏
go
// 带 fencing token 的分布式锁用法
type LockResult struct {
Granted bool
Token int64 // 单调递增的 fencing token
}
// 使用锁的代码将 token 写入存储系统
// 存储系统拒绝过期 token 的写入(token < 当前已见最大 token → 拒绝)7.3 etcd 分布式锁
go
// etcd 分布式锁 — 基于租约(Lease)
import clientv3 "go.etcd.io/etcd/client/v3"
func etcdLock(cli *clientv3.Client, key string, ttl int64) (*clientv3.LeaseGrantResponse, error) {
// 1. 创建租约
lease, err := cli.Grant(ctx, ttl)
// 2. 用租约创建 key(事务保证原子性)
txn := cli.Txn(ctx).If(
clientv3.Compare(clientv3.CreateRevision(key), "=", 0),
).Then(
clientv3.OpPut(key, "locked", clientv3.WithLease(lease.ID)),
)
resp, err := txn.Commit()
if !resp.Succeeded {
return nil, ErrLockHeld // 锁被其他人持有
}
// 3. 续租(心跳)
keepAliveCh, _ := cli.KeepAlive(ctx, lease.ID)
go func() {
for range keepAliveCh {} // 自动续租
}()
return lease, nil
}8. 分布式 ID 生成 — 时钟回拨处理
雪花算法(Snowflake)依赖系统时钟单调递增,但时钟回拨(NTP 同步、虚拟机迁移)会导致 ID 冲突。
8.1 时钟回拨的三种处理策略
go
type Snowflake struct {
lastTimestamp int64
sequence int64
workerID int64
}
func (s *Snowflake) NextID() int64 {
timestamp := time.Now().UnixMilli()
if timestamp < s.lastTimestamp {
// 时钟回拨了!
offset := s.lastTimestamp - timestamp
switch strategy {
case "wait":
// 策略1: 等待时钟追上(适合短暂回拨)
time.Sleep(time.Duration(offset) * time.Millisecond)
timestamp = time.Now().UnixMilli()
case "throw":
// 策略2: 抛异常(严格模式,直接拒绝服务)
panic("clock moved backwards")
case "tolerant":
// 策略3: 容忍模式 — 使用上次时间戳 + 额外序列号
// 需要扩展序列号位宽来容纳"回拨补偿"
if offset < maxBackwardDrift {
timestamp = s.lastTimestamp
// 使用"回拨序列号段"避免冲突
s.sequence = s.backupSequence
}
}
}
// ... 正常生成逻辑 ...
}8.2 主流方案的回拨对策
| 方案 | 回拨对策 |
|---|---|
| Twitter Snowflake | 原始论文没处理,生产实现通常抛异常 |
| 百度 UidGenerator | DefaultUidGenerator: 抛异常; CachedUidGenerator: 容忍(ring buffer 缓存) |
| 美团 Leaf | Leaf-Segment: 不依赖时钟; Leaf-Snowflake: 等待 + 抛异常 |
| Sonyflake | 从时间戳改为"从起始时间的间隔",减少精度要求 |
| UUID v7 | 基于毫秒时间戳,不防回拨但概率极低 |
8.3 Leaf-Segment — 彻底避免时钟依赖
Leaf-Segment 方案(美团):
- 号段(segment)模式: 从 DB 一次取一段 ID(如 1000-2000)
- 用完后异步取下一段 → 内存分配,零时钟依赖
- 双 buffer: 当前号段用完前,预加载下一个号段
DB 表:
biz_tag | max_id | step | update_time
order | 2000 | 1000 | ...
每次更新: UPDATE id_alloc SET max_id = max_id + step WHERE biz_tag = 'order'
→ 原子递增,拿到号段范围9. 高并发计数系统设计
设计一个支持亿级 QPS 的计数系统,如微博点赞数、视频播放量。核心矛盾:写多读多 + 数据量大 + 实时性要求。
9.1 方案演进
mermaid
flowchart TB
subgraph "1. 单机计数器 (QPS < 1k)"
A["Redis INCR key"]
end
subgraph "2. 分片计数 (QPS < 10w)"
B["Redis INCR key:{shard}"]
end
subgraph "3. 本地缓存 + 异步聚合 (QPS < 100w)"
C["本地 Map 累加 → 定时 flush Redis"]
end
subgraph "4. 日志驱动 + 离线聚合 (QPS > 100w)"
D["本地写日志 → Kafka → Flink 聚合 → Redis"]
end方案 1: 单机 Redis INCR — 最简单
go
func addLike(postID string) error {
_, err := rdb.Incr(ctx, "like:"+postID).Result()
return err
}瓶颈:单 Redis key 热点 → CPU 单核 100%,QPS 上限约 10w。
方案 2: 分片计数 — 牺牲读实时性换写入扩展
go
const shardCount = 100
func addLikeSharded(postID string, userID string) error {
shard := hash(userID) % shardCount
key := fmt.Sprintf("like:%s:%d", postID, shard)
return rdb.Incr(ctx, key).Result()
}
func getLikeCount(postID string) (int64, error) {
var total int64
for i := 0; i < shardCount; i++ {
key := fmt.Sprintf("like:%s:%d", postID, i)
val, _ := rdb.Get(ctx, key).Int64()
total += val
}
return total, nil
}
// 写出: 分散到 100 个 key,无热点
// 读取: 需要 MGET 100 次或 pipeline 批量取方案 3: 本地 Buffer + 定时 flush — 最高写入吞吐
go
type LocalCounter struct {
mu sync.Mutex
buffers map[string]int64 // key → 本地累加值
rdb *redis.Client
}
func (c *LocalCounter) Incr(key string) {
c.mu.Lock()
c.buffers[key]++
c.mu.Unlock()
}
// 定时 flush (每秒或每 100 次)
func (c *LocalCounter) flushToRedis() {
c.mu.Lock()
snapshot := c.buffers
c.buffers = make(map[string]int64) // 快速交换
c.mu.Unlock()
pipe := c.rdb.Pipeline()
for key, delta := range snapshot {
pipe.IncrBy(ctx, key, delta)
}
pipe.Exec(ctx)
}
// 风险: 进程挂了会丢 buffer 中的数据 (可接受,计数不是金融交易)方案 4: 日志驱动 + 离线聚合 — 最实时 + 最准确
text
写入路径:
用户点赞 → 本地写日志(Kafka Producer [点赞 post_id=123]")
→ Kafka 分区 (按 post_id)
→ Flink 窗口聚合 (每 1s 滚动窗口)
→ Redis HSET post:123 likes 聚合后的值
读取路径:
实时: Redis GET post:123:likes
历史: ClickHouse SELECT SUM(like_count) FROM likes WHERE post_id = 1239.2 去重计数 — HyperLogLog
如果需求是"UV 数"而非"总次数",去重是关键。HyperLogLog 用 12KB 空间近似计数 2^64 个元素,误差约 0.81%。
go
// Redis HyperLogLog
for _, userID := range userIDs {
rdb.PFAdd(ctx, "uv:post:"+postID, userID)
}
count, _ := rdb.PFCount(ctx, "uv:post:"+postID).Result() // 近似 UV| 需求 | 方案 | 空间 | 精度 |
|---|---|---|---|
| 总计数 | INCR / 分片 INCR | O(1) per shard | 精确 |
| 去重计数 (UV) | HyperLogLog | 12KB | ±0.81% |
| 去重 + 精确 | Bloom Filter 判断存在 + Set 计数 | 取决于误判率 | 精确 |
| Top N 排行榜 | Sorted Set (ZINCRBY) | O(N) | 精确 |
9.3 写热点问题
text
问题的根源:
单个计数 key (如某个明星的微博) 成为写热点
→ Redis 单线程,该 key 的 INCR 成为瓶颈
→ 即使分片,多个线程/进程仍可能竞争同一个 key
解决:
1. 分片写入、批量读取 (方案 2)
2. 本地累加 + 异步 flush (方案 3)
3. 写 Kafka 异步消费 (方案 4)
4. 如果只是展示用且允许近似: Morris Counter (概率计数)go
// Morris Counter: 8bit 计数百万级,概率递增
type MorrisCounter struct {
value uint8
}
func (m *MorrisCounter) Incr() {
// 递增概率 = 1/(m.value * factor + 1)
p := 1.0 / (float64(m.value)*10.0 + 1.0)
if rand.Float64() < p {
m.value++
}
}
func (m *MorrisCounter) Estimate() float64 {
// 实际值 ≈ (e^factor * m.value - 1) / factor
// 见 Redis LFU 淘汰算法中的 Morris Counter 实现
return math.Pow(math.E, float64(m.value)*0.1) / 0.1
}
// Redis 的 LFU 淘汰策略使用了此技术10. 时间轮(Timing Wheel)
假设系统中有数百万个定时任务(如订单 30 分钟未支付自动取消),最小堆的插入和删除是 O(log n)。时间轮将定时任务的插入、删除、到期触发全部优化到 O(1)。
10.1 为什么最小堆不够用?
text
场景: 电商订单系统,每秒产生 1000 个新订单
每个订单需要在 30 分钟后检查是否支付
最小堆方案:
插入: O(log N) — 100万任务时约 20 次比较
删除/到期: O(log N)
内存: 每个定时器约 40B (堆节点 + 回调函数指针)
100万定时器的最小堆 → 插入约 20 次比较 → 看起来还好
但每秒钟 1000 次插入 → 每秒 20000 次比较 → 再加上取出、取消 → CPU 可观
关键瓶颈不是单次 O(log N),而是"海量定时器 + 高频操作"下的累积开销。10.2 单层时间轮原理
时间轮 = 环形数组(槽位数组) + 指针
假设精度 = 1 秒,一圈 = 60 秒 → 60 个槽位
槽位 0 1 2 3 ... 58 59
↓ ↓ ↓ ↓ ↓ ↓
[任务链表] [任务链表]
指针每秒前进一格,指向的槽位中的所有任务到期执行。
插入任务: 计算 (now + delay) % 60 → 直接放到对应槽位的链表中 → O(1)
删除任务: 从链表中摘除 → O(1)
到期触发: 指针走到槽位 → 遍历链表执行所有任务 → O(该槽位的任务数)go
type TimeWheel struct {
slots []*TaskList // 槽位数组
current int // 当前指针位置
tick time.Duration // 每格时间精度
slotNum int // 槽位总数
mu sync.Mutex
}
type Task struct {
key string
circle int // 还需转几圈
callback func()
prev, next *Task
}
func (tw *TimeWheel) Add(delay time.Duration, key string, callback func()) {
tw.mu.Lock()
defer tw.mu.Unlock()
// 计算目标槽位和圈数
ticks := int(delay / tw.tick)
slot := (tw.current + ticks) % tw.slotNum
circle := ticks / tw.slotNum
task := &Task{key: key, circle: circle, callback: callback}
tw.slots[slot].PushBack(task)
}
func (tw *TimeWheel) Tick() {
tw.mu.Lock()
defer tw.mu.Unlock()
// 取出当前槽所有任务
slot := tw.slots[tw.current]
for task := slot.Head(); task != nil; {
next := task.next
if task.circle > 0 {
task.circle-- // 圈数减 1,不执行
} else {
slot.Remove(task)
go task.callback() // 到期执行
}
task = next
}
tw.current = (tw.current + 1) % tw.slotNum
}10.3 单层时间轮的局限
text
问题: 如果定时范围很大(如 30 分钟),槽位数量会爆炸。
30 分钟 = 1800 秒
如果精度 = 1 秒 → 需要 1800 个槽位
如果精度 = 10ms → 需要 180000 个槽位
→ 内存过大,且大多数槽位是空的
既要精度高,又要范围大 → 分层时间轮10.4 分层时间轮(Hierarchical Timing Wheel)
核心思想:用多个不同粒度的时间轮,类似时钟的"时-分-秒"。
text
┌──────────────────┐
│ 第 3 层 (时轮) │ 1 格 = 1 小时,共 24 格
│ 覆盖 0-23 小时 │
└────────┬─────────┘
│ 当时针走一格 → 把该格的任务降级到分轮
▼
┌──────────────────┐
│ 第 2 层 (分轮) │ 1 格 = 1 分钟,共 60 格
│ 覆盖 0-59 分钟 │
└────────┬─────────┘
│ 当分针走一格 → 把该格的任务降级到秒轮
▼
┌──────────────────┐
│ 第 1 层 (秒轮) │ 1 格 = 1 秒,共 60 格
│ 覆盖 0-59 秒 │
└────────┬─────────┘
│ 每秒走一格 → 到期执行
▼
[执行回调]mermaid
flowchart TB
subgraph Layer3["第三层: 时轮 (1格=1h, 24格)"]
H0["0h"] --> H1["1h"] --> H2["2h"] --> H23["23h"]
end
subgraph Layer2["第二层: 分轮 (1格=1min, 60格)"]
M0["0min"] --> M1["1min"] --> M59["59min"]
end
subgraph Layer1["第一层: 秒轮 (1格=1s, 60格)"]
S0["0s"] --> S1["1s"] --> S59["59s"]
end
H1 -->|"时针走1格<br/>将H1的任务降级"| Layer2
M1 -->|"分针走1格<br/>将M1的任务降级"| Layer1
S0 -->|"每秒走1格<br/>到期执行"| EXEC["执行回调"]数学分析:
text
三级时间轮: 时(24格) + 分(60格) + 秒(60格)
总槽位数: 24 + 60 + 60 = 144 格
覆盖范围: 24 × 60 × 60 = 86400 秒 = 24 小时
精度: 1 秒
如果只用单层时间轮达到相同效果:
需要 86400 个槽位
→ 分层设计节省了 86400 / 144 ≈ 600 倍内存
通用公式:
N 级时间轮,每级 Wi 格
总槽位 = Σ Wi
覆盖范围 = Π Wi × 基础精度10.5 完整实现
go
type HierarchicalTimeWheel struct {
levels []*TimeWheel // 多层时间轮,从粗到细
tick time.Duration // 最细粒度
}
func NewHierarchical(levelConfigs []LevelConfig) *HierarchicalTimeWheel {
htw := &HierarchicalTimeWheel{tick: levelConfigs[0].Tick}
for i, cfg := range levelConfigs {
tw := &TimeWheel{
slots: make([]*TaskList, cfg.Slots),
slotNum: cfg.Slots,
tick: cfg.Tick,
}
// 上层的一格 = 下层的一整圈
if i > 0 {
tw.tick = levelConfigs[i-1].Tick *
time.Duration(levelConfigs[i-1].Slots)
}
htw.levels = append(htw.levels, tw)
}
return htw
}
func (htw *HierarchicalTimeWheel) Add(delay time.Duration, key string, cb func()) {
// 找到合适的时间轮层级
for i := len(htw.levels) - 1; i >= 0; i-- {
tw := htw.levels[i]
if delay >= tw.tick {
tw.Add(delay, key, cb)
return
}
}
// 小于最小精度,放入最底层
htw.levels[0].Add(delay, key, cb)
}
func (htw *HierarchicalTimeWheel) Tick() {
// 只驱动最底层(最细粒度)
htw.levels[0].Tick()
// 检查是否要降级上层任务
// 当底层转完一圈 → 降级上一层对应槽位的任务
for i := 1; i < len(htw.levels); i++ {
if htw.levels[i-1].current == 0 {
htw.cascadeDown(i)
}
}
}
// 降级:将上层当前槽位的任务重新分配到下层
func (htw *HierarchicalTimeWheel) cascadeDown(level int) {
upper := htw.levels[level]
lower := htw.levels[level-1]
slot := upper.slots[upper.current]
for task := slot.Head(); task != nil; {
next := task.next
slot.Remove(task)
// 重新计算在下一层的位置
lower.Add(task.remainDelay, task.key, task.callback)
task = next
}
}10.6 Kafka 中的时间轮
Kafka 大量使用时间轮来管理延迟操作:
| 场景 | 延迟类型 | 实现 |
|---|---|---|
| Producer ACK 等待 | acks=all 等待 follower 同步 | 时间轮管理超时 |
| Consumer 拉取等待 | fetch.max.wait.ms | 时间轮触发返回 |
| 延迟创建主题/分区 | 管理后台延迟任务 | SystemTimer (基于时间轮) |
| 事务超时 | transaction.timeout.ms | 时间轮到期中止事务 |
Kafka 的实现(org.apache.kafka.common.utils.Timer):
java
// Kafka 时间轮的核心数据结构
class TimingWheel {
private final long tickMs; // 每格毫秒数
private final int wheelSize; // 槽位数
private final long interval; // 一圈覆盖时间 = tickMs * wheelSize
private final TimerTaskList[] buckets; // 槽位数组
private long currentTime; // 当前时间(向下取整到 tickMs)
private volatile TimingWheel overflowWheel; // 溢出层(上层时间轮)
}mermaid
flowchart TB
subgraph Kafka["Kafka 分层时间轮"]
L0["层0: tickMs=1ms, wheelSize=20<br/>覆盖 20ms"]
L1["层1: tickMs=20ms, wheelSize=20<br/>覆盖 400ms"]
L2["层2: tickMs=400ms, wheelSize=20<br/>覆盖 8s"]
L3["层3: tickMs=8s, wheelSize=15<br/>覆盖 120s"]
end
L0 -->|"超出一圈 → 溢出到"| L1
L1 -->|"超出一圈 → 溢出到"| L2
L2 -->|"超出一圈 → 溢出到"| L310.7 Netty 中的时间轮
Netty 的 HashedWheelTimer 使用了单层时间轮 + 圈数:
java
// Netty HashedWheelTimer
new HashedWheelTimer(
100, TimeUnit.MILLISECONDS, // 每格 100ms
512 // 512 个槽位
);
// 覆盖范围: 100ms × 512 = 51.2 秒
// Netty 没有使用多层,而是用单层 + circle 计数器
// 设计出发点是 I/O 超时通常在秒级,51.2 秒够用10.8 对比总结
| 维度 | 最小堆 | 单层时间轮 | 分层时间轮 |
|---|---|---|---|
| 插入 | O(log N) | O(1) | O(1) |
| 删除 | O(log N) | O(1) | O(1) |
| 到期触发 | O(log N) | O(1) | O(1) |
| 内存 | O(N) | O(M+N) | O(ΣWi + N) |
| 精度 | 任意 | tick 的整数倍 | tick 的整数倍 |
| 适用场景 | 定时器少 | 定时器多,范围小 | 定时器多,范围大 |
| 代表实现 | Linux Timerfd | Netty | Kafka |
11. 海量 Token/长连接生命周期管理
假设系统有 1000 万个在线长连接,每个连接有自定义有效期(几秒到几天不等)。传统方案(数据库轮询、Redis 过期)都无法同时满足"精准到期"和"低内存"。
11.1 为什么 Redis 过期机制不够用
text
Redis 过期删除策略:惰性删除 + 定期随机抽样
惰性删除:访问 key 时检查是否过期 → 过期则删除
→ 如果 key 从不被访问 → 永远不会被惰性删除 → 占着内存
定期删除:每秒 10 次,每次随机抽 20 个 key
→ 如果过期 key 占比 > 25% → 重复抽取
→ 每次最多 16 轮 → 每轮 20 个 = 最多 320 个/次
→ 每秒最多删 320 × 10 = 3200 个过期 key
1000 万个 token,100 万个同时过期:
3200 个/秒 → 需要 312 秒 ≈ 5 分钟才能删完
→ 这 5 分钟内,100 万过期 token 的内存无法释放!
对于自定义有效期(持续新产生 token),过期间隔分布在整个时间线上
→ 如果生成速率 1000/秒,1 小时后有 360 万 token → 但有的 1 分钟过期,有的 1 天过期
→ Redis 随机抽样无法保证"恰好到期的"能被及时清理11.2 方案设计:延迟队列 + 分层时间轮
mermaid
flowchart TB
subgraph Incoming["Token 创建"]
A["新 Token<br/>TTL = 30 分钟"] --> B["计算到期时间<br/>expireAt = now + 30min"]
B --> C["存入 Hash:<br/>token:detail:{tokenID}"]
B --> D["投递到延迟队列:<br/>expireAt → tokenID"]
end
subgraph DelayQueue["延迟队列(基于时间轮)"]
D --> E["分层时间轮<br/>毫秒→秒→分→时"]
E --> F{"当前 tick 到期?"}
F -->|"是"| G["取出到期 tokenID 列表"]
F -->|"否"| E
end
subgraph Handler["到期处理"]
G --> H["查询 token:detail:{tokenID}"]
H --> I{"Token 已被续期?"}
I -->|"否"| J["执行过期逻辑:<br/>关闭连接、回收资源"]
I -->|"是"| K["跳过(已续期)"]
end
subgraph Renewal["续期处理"]
L["客户端心跳/业务操作"] --> M["延长 TTL"]
M --> N["更新 token:detail 的 expireAt"]
M --> O["重新投递到延迟队列<br/>新的 expireAt"]
Note right of O: "旧的延迟任务到期时<br/>会检查是否已续期"
end11.3 核心实现
go
type TokenManager struct {
details map[string]*TokenDetail // tokenID → 详情
timeWheel *HierarchicalTimeWheel // 分层时间轮
mu sync.RWMutex
}
type TokenDetail struct {
ID string
UserID string
ExpireAt time.Time
Conn net.Conn // 关联的长连接
RenewedAt time.Time
}
// 创建 Token
func (tm *TokenManager) Create(userID string, ttl time.Duration) string {
tokenID := generateTokenID()
expireAt := time.Now().Add(ttl)
detail := &TokenDetail{
ID: tokenID, UserID: userID,
ExpireAt: expireAt, Conn: conn,
}
tm.mu.Lock()
tm.details[tokenID] = detail
tm.mu.Unlock()
// 投递到时间轮
tm.timeWheel.Add(ttl, tokenID, func() {
tm.handleExpiry(tokenID, expireAt)
})
return tokenID
}
// 到期处理
func (tm *TokenManager) handleExpiry(tokenID string, expectedExpireAt time.Time) {
tm.mu.Lock()
detail, ok := tm.details[tokenID]
if !ok {
tm.mu.Unlock()
return // 已删除
}
// 关键检查:是否已被续期?
if detail.ExpireAt.After(expectedExpireAt) {
// Token 已被续期 → 旧延迟任务作废
tm.mu.Unlock()
return
}
// 真正到期 → 执行清理
delete(tm.details, tokenID)
tm.mu.Unlock()
// 关闭连接、记录日志、回收资源
detail.Conn.Close()
log.Info("token expired", "tokenID", tokenID, "userID", detail.UserID)
}
// 续期
func (tm *TokenManager) Renew(tokenID string, newTTL time.Duration) error {
tm.mu.Lock()
detail, ok := tm.details[tokenID]
if !ok {
tm.mu.Unlock()
return ErrTokenNotFound
}
newExpireAt := time.Now().Add(newTTL)
detail.ExpireAt = newExpireAt
detail.RenewedAt = time.Now()
tm.mu.Unlock()
// 重新投递到时间轮(新 TTL)
tm.timeWheel.Add(newTTL, tokenID, func() {
tm.handleExpiry(tokenID, newExpireAt)
})
return nil
}11.4 内存压榨 — 冷数据下沉
1000 万个长连接全部存在内存中仍然压力大。对于 TTL 很长的 token(如 7 天),可以采用冷热分离:
text
热数据(TTL < 1 小时):
→ 全部在内存中 + 时间轮管理
→ 快速到期检测
温数据(1 小时 < TTL < 1 天):
→ 时间轮只保留前 1 小时的精度
→ 1 小时后降级为分钟级轮询(每分钟扫一次 DB/Redis)
冷数据(TTL > 1 天):
→ 不在内存中维护时间轮
→ 定时任务(每 5 分钟)扫一次 DB 中到期的 token
→ 或依赖 Redis TTL 兜底go
func (tm *TokenManager) CreateOrSelectiveMemory(token *TokenDetail, ttl time.Duration) {
if ttl < time.Hour {
// 热:全内存管理
tm.details[token.ID] = token
tm.timeWheel.Add(ttl, token.ID, tm.handleExpiry)
} else if ttl < 24*time.Hour {
// 温:前 1 小时精确管理,之后降级
tm.details[token.ID] = token
tm.timeWheel.Add(time.Hour, token.ID, func() {
tm.degradeToWarm(token.ID) // 降级到轮询队列
})
} else {
// 冷:仅 DB/Redis,定时轮询
tm.coldTokens[token.ID] = token.ExpireAt
}
}11.5 对比总结
| 方案 | 到期精度 | 内存占用 | 1000 万 token | 适用 |
|---|---|---|---|---|
| Redis 自带过期 | 秒级(随机抽样) | 中 | ∼1GB(每个 key ~100B) | 小规模、对延迟不敏感 |
| 纯时间轮 | 毫秒级 | 高(所有 token 在内存) | 需大量堆外内存 | 实时性要求高 |
| 时间轮 + 冷热分离 | 热: 毫秒 / 冷: 分钟 | 低(热数据占用少) | 可承载 | 大规模、混合 TTL |
| DB 轮询 | 取决于扫表间隔 | 最低 | 扫 1000 万行的开销 | TTL 很长、不频繁的操作 |
11.6 Redis 结合时间轮的混合方案
text
生产环境最实用的方案: Redis 存储 token 详情 + 本地时间轮管理热数据
结构:
Redis Hash: token:{tokenID} → {userID, expireAt, ...}
Redis SortedSet: token_expiry → {score=expireAt_timestamp, member=tokenID}
本地时间轮管理最近 N 分钟到期的 token(热数据)
流程:
1. 新 token → 写入 Redis Hash + 写入 SortedSet (score=expireAt)
2. 本地时间轮每隔 1 秒从 SortedSet 的头部(最早到期的)拉取下一批
3. 时间轮触发到期 → Redis Hash 查询确认 → 执行清理 → 从 SortedSet 删除
4. SortedSet 的 score 是精确的 → 毫秒级精度
5. 服务重启 → 重新从 SortedSet 拉取 → 无需担心本地状态丢失go
func (tm *TokenManager) prefetchLoop() {
ticker := time.NewTicker(1 * time.Second)
for range ticker.C {
// 从 Redis SortedSet 拉取下一分钟到期的 token
now := time.Now().UnixMilli()
tokenIDs, _ := rdb.ZRangeByScore(ctx, "token_expiry", &redis.ZRangeBy{
Min: "0",
Max: strconv.FormatInt(now+60000, 10), // 未来 60 秒
Offset: 0,
Count: 1000,
}).Result()
for _, tokenID := range tokenIDs {
// 如果本地还没管理 → 加入时间轮
if !tm.isLocalManaged(tokenID) {
detail := tm.loadFromRedis(tokenID)
tm.addToLocalWheel(detail)
}
}
}
}
登录后即可发表评论 👇