学习目标
能力目标
完成本章学习后,你将能够:
- 设计直播弹幕系统架构,理解 WebSocket 长连接、消息广播的读扩散模型、消息队列削峰与弹幕限流策略。
- 掌握朋友圈 Feed 流的三种模式(推模式、拉模式、推拉结合),能根据用户体量和活跃度选择合适的方案并实现时间线排序与可见性控制。
- 制定数据统计页面的加速方案,综合运用预计算、多级缓存、CDN、异步加载与物化视图等手段解决高并发读取问题。
- 识别并解决缓存雪崩问题,理解热 key 过期引发的多米诺效应,掌握互斥锁重建、热点 key 永不过期、多级缓存与布隆过滤器等防护手段。
- 实现长短链转换与扫码登录,掌握发号器、Base62 编码、301/302 重定向选型,以及二维码扫码登录的完整时序设计。
前置知识
- 了解 WebSocket 协议与 HTTP 长连接的区别
- 熟悉 Redis 基本操作(String/Hash/ZSet/Pub-Sub)
- 了解 Go 语言并发编程(goroutine、channel、sync 包)
动手做三件事
- 用 Go + gorilla/websocket 实现一个支持多房间的弹幕广播服务,包含连接管理和消息限流。
- 用 Redis ZSet 实现一个简易的朋友圈时间线,模拟推模式和拉模式两种读取方式。
- 实现一个长短链转换服务,包含发号器、Base62 编码、Redis 映射和 HTTP 302 重定向。
一、直播弹幕系统设计
1.1 用生活类比先建立直觉
弹幕系统就像体育场广播: 想象一个十万人的体育场,每个观众都能往广播站递纸条(发弹幕),广播站需要把纸条内容念给所有人听(广播)。问题来了:
- 如果十万人同时递纸条,广播站会被淹没(消息队列削峰)。
- 如果每张纸条都念一遍,观众听不过来(限流)。
- 如果广播站只有一个人念,念不过来(分布式扩容)。
- 如果某个区域信号不好听不到,观众会投诉(长连接保活)。
graph TB
U1[用户A
发弹幕] --> MQ[消息队列
Kafka削峰]
U2[用户B
发弹幕] --> MQ
U3[用户C
发弹幕] --> MQ
MQ --> W1[弹幕处理Worker1
限流+排序]
MQ --> W2[弹幕处理Worker2
限流+排序]
W1 --> R[Redis Pub/Sub
按房间分发]
W2 --> R
R --> WS1[WebSocket节点1
推送给房间观众]
R --> WS2[WebSocket节点2
推送给房间观众]桥接: 体育场广播类比映射到弹幕系统——纸条递送对应消息队列削峰,广播站限速对应弹幕限流(每秒最多展示 N 条),多广播员对应分布式 Worker,区域信号对应 WebSocket 长连接保活。
1.2 工程要点
弹幕系统的核心设计点:
| 模块 | 技术选型 | 设计要点 |
|---|---|---|
| 长连接 | WebSocket | 双向通信,服务端可主动推送;心跳保活 30 秒 |
| 消息削峰 | Kafka / NSQ | 弹幕写入先入队列,消费端按速率处理 |
| 消息分发 | Redis Pub/Sub | 按房间订阅频道,WebSocket 节点订阅对应频道 |
| 弹幕限流 | 令牌桶 / 滑动窗口 | 每秒最多展示 N 条弹幕,超出丢弃或排队 |
| 弹幕排序 | 服务端时间戳 | 保证同一房间内弹幕按时间有序展示 |
| 水平扩展 | 一致性哈希 | WebSocket 连接按 roomId 哈希到不同节点 |
Go 实现 WebSocket 弹幕广播核心代码:
package main
import (
"log"
"net/http"
"sync"
"github.com/gorilla/websocket"
)
// Room 表示一个直播房间
type Room struct {
mu sync.RWMutex
clients map[*websocket.Conn]bool
}
// 步骤1:创建房间管理器
type RoomManager struct {
mu sync.RWMutex
rooms map[string]*Room
}
func NewRoomManager() *RoomManager {
return &RoomManager{rooms: make(map[string]*Room)}
}
// 步骤2:用户加入房间
func (rm *RoomManager) Join(roomId string, conn *websocket.Conn) {
rm.mu.Lock()
room, ok := rm.rooms[roomId]
if !ok {
room = &Room{clients: make(map[*websocket.Conn]bool)}
rm.rooms[roomId] = room
}
rm.mu.Unlock()
room.mu.Lock()
room.clients[conn] = true
room.mu.Unlock()
log.Printf("用户加入房间 %s, 当前在线: %d", roomId, len(room.clients))
}
// 步骤3:用户离开房间
func (rm *RoomManager) Leave(roomId string, conn *websocket.Conn) {
rm.mu.RLock()
room, ok := rm.rooms[roomId]
rm.mu.RUnlock()
if !ok {
return
}
room.mu.Lock()
delete(room.clients, conn)
room.mu.Unlock()
}
// 步骤4:广播弹幕(读扩散:一条消息推给房间所有人)
func (rm *RoomManager) Broadcast(roomId string, message []byte) {
rm.mu.RLock()
room, ok := rm.rooms[roomId]
rm.mu.RUnlock()
if !ok {
return
}
room.mu.RLock()
defer room.mu.RUnlock()
for conn := range room.clients {
// 步骤5:非阻塞发送,发送失败则关闭连接
err := conn.WriteMessage(websocket.TextMessage, message)
if err != nil {
log.Printf("发送失败: %v", err)
conn.Close()
delete(room.clients, conn)
}
}
}
var upgrader = websocket.Upgrader{
CheckOrigin: func(r *http.Request) bool { return true },
}
func main() {
rm := NewRoomManager()
http.HandleFunc("/ws", func(w http.ResponseWriter, r *http.Request) {
roomId := r.URL.Query().Get("room")
conn, err := upgrader.Upgrade(w, r, nil)
if err != nil {
log.Println("WebSocket升级失败:", err)
return
}
defer conn.Close()
rm.Join(roomId, conn)
defer rm.Leave(roomId, conn)
for {
// 步骤6:读取用户发来的弹幕
_, msg, err := conn.ReadMessage()
if err != nil {
break
}
// 步骤7:广播给房间内所有用户
rm.Broadcast(roomId, msg)
}
})
log.Println("弹幕服务启动 :8080/ws?room=xxx")
http.ListenAndServe(":8080", nil)
}
⚠️ 新手必踩的坑:
conn.WriteMessage不是并发安全的。如果同一个连接同时被多个 goroutine 调用 WriteMessage 会引发 panic。解决方案:每个连接用一个独立的 channel 做发送队列,由单个 goroutine 负责写。
二、朋友圈与 Feed 流
2.1 用生活类比先建立直觉
朋友圈 Feed 流就像报社送报纸:
- 推模式(写扩散):你写完文章后,报社帮你复印 N 份,分别送到每个粉丝家门口的邮箱里(写入粉丝的 inbox)。粉丝看报时直接从自己邮箱拿(读取快),但你写文章时很累——如果有 100 万粉丝就要复制 100 万份(写入慢)。
- 拉模式(读扩散):你写完文章只放在报社公告栏(自己的 outbox),粉丝想看报时自己去所有关注人的公告栏收集(读取慢),但你写文章很轻松(写入快)。
- 推拉结合:大 V(粉丝多)只放公告栏(拉模式),普通用户(粉丝少)送到粉丝邮箱(推模式)。粉丝看报时:先拿邮箱里的(普通用户推来的),再去大 V 公告栏拉(大 V 的),合并排序。
graph TB
subgraph 推模式-写扩散
P1[发布者写动态] -->|写入自己的outbox| O1[发布者outbox]
P1 -->|推送到粉丝inbox| I1[粉丝A inbox]
P1 -->|推送到粉丝inbox| I2[粉丝B inbox]
end
subgraph 拉模式-读扩散
P2[发布者写动态] -->|只写入自己的outbox| O2[发布者outbox]
F1[粉丝读取] -->|主动拉取关注人outbox| O2
end
subgraph 推拉结合
BV[大V发动态] -->|只写outbox| O3[大Voutbox]
NU[普通用户发动态] -->|推送到粉丝inbox| I3[粉丝inbox]
F2[粉丝读取] -->|合并inbox+拉取大V| I3
F2 -->|拉取大V最新动态| O3
end桥接: 报社类比映射到 Feed 流——推模式是写时扩散(发动态时推到所有粉丝收件箱),拉模式是读时聚合(看动态时从所有关注人收件箱拉取),推拉结合是大 V 用拉、普通用户用推,兼顾写入效率和读取体验。
2.2 工程要点
三种 Feed 流模式对比
| 模式 | 写入开销 | 读取开销 | 适用场景 | 缺点 |
|---|---|---|---|---|
| 推模式 | O(粉丝数) | O(1) | 粉丝数少(小于1000) | 大 V 写入爆炸 |
| 拉模式 | O(1) | O(关注数) | 关注数少(小于500) | 读取延迟高 |
| 推拉结合 | 普通 O(粉丝数),大 V O(1) | O(关注的大 V 数) | 混合场景 | 逻辑复杂 |
Feed 流补充要点
- 推模式:发布动态时遍历粉丝列表,将动态 ID 写入每个粉丝的 Redis ZSet(score=时间戳)。大 V 发动态时不推,只写入自己的 outbox。
- 拉模式:读取时从关注列表中取出所有关注人的 outbox,合并后按时间排序。关注人太多时只取最近活跃的前 N 个。
- 推拉结合:粉丝数小于阈值的用户用推模式,粉丝数大于等于阈值的大 V 用拉模式。读取时先读自己的 inbox(推来的),再补充拉取关注的大 V 的 outbox,合并排序。
Go 实现 Feed 流核心逻辑:
package main
import (
"context"
"fmt"
"time"
"github.com/redis/go-redis/v9"
)
type FeedService struct {
rdb *redis.Client
bigVThreshold int64 // 大V粉丝数阈值
}
// 步骤1:发布动态
func (s *FeedService) Publish(ctx context.Context, userId string, postId string, timestamp int64) error {
// 步骤2:写入发布者outbox(所有模式都需要)
outboxKey := fmt.Sprintf("outbox:%s", userId)
s.rdb.ZAdd(ctx, outboxKey, redis.Z{
Score: float64(timestamp),
Member: postId,
})
// 步骤3:查询粉丝数,决定推还是拉
fansKey := fmt.Sprintf("fans:%s", userId)
fanCount, _ := s.rdb.SCard(ctx, fansKey).Result()
if fanCount < s.bigVThreshold {
// 步骤4:普通用户-推模式:写入每个粉丝的inbox
fans, _ := s.rdb.SMembers(ctx, fansKey).Result()
for _, fanId := range fans {
inboxKey := fmt.Sprintf("inbox:%s", fanId)
s.rdb.ZAdd(ctx, inboxKey, redis.Z{
Score: float64(timestamp),
Member: postId,
})
// 步骤5:inbox只保留最近1000条
s.rdb.ZRemRangeByRank(ctx, inboxKey, 0, -1001)
}
}
// 大V不推,粉丝读取时主动拉取outbox
return nil
}
// 步骤6:读取Feed流(推拉结合)
func (s *FeedService) GetFeed(ctx context.Context, userId string) ([]string, error) {
// 步骤7:先读inbox(推模式写入的)
inboxKey := fmt.Sprintf("inbox:%s", userId)
posts, _ := s.rdb.ZRevRange(ctx, inboxKey, 0, 19).Result()
// 步骤8:补充拉取关注的大V的outbox
followingKey := fmt.Sprintf("following:%s", userId)
followingIds, _ := s.rdb.SMembers(ctx, followingKey).Result()
for _, followId := range followingIds {
fansKey := fmt.Sprintf("fans:%s", followId)
fanCount, _ := s.rdb.SCard(ctx, fansKey).Result()
if fanCount >= s.bigVThreshold {
// 步骤9:大V的动态从outbox拉取
outboxKey := fmt.Sprintf("outbox:%s", followId)
bigVPosts, _ := s.rdb.ZRevRange(ctx, outboxKey, 0, 19).Result()
posts = append(posts, bigVPosts...)
}
}
// 步骤10:合并后重新按时间排序,取前20条
// 实际项目中用Redis ZUNIONSTORE或服务端排序
return posts, nil
}
func main() {
rdb := redis.NewClient(&redis.Options{Addr: "localhost:6379"})
svc := &FeedService{rdb: rdb, bigVThreshold: 1000}
ctx := context.Background()
// 模拟发布动态
svc.Publish(ctx, "user:001", "post:1001", time.Now().Unix())
// 模拟读取Feed流
feed, _ := svc.GetFeed(ctx, "user:002")
fmt.Println("Feed流:", feed)
}
⚠️ 新手必踩的坑: 推模式的 inbox 必须设置上限(如最近 1000 条),否则长期不登录的用户的 inbox 会无限膨胀,浪费 Redis 内存。同时要设置过期时间,长期不活跃用户的 inbox 自动清理。
三、数据统计页面加速
3.1 用生活类比先建立直觉
数据统计页面加速就像餐厅出菜加速: 想象一家餐厅的"今日菜品销量排行榜"展示屏:
- 不优化:每次有顾客问"今天什么卖得好",厨师都重新去数一遍今天卖了多少菜(实时查 DB 聚合),顾客等半天。
- 预计算:每小时自动统计一次销量排行榜,写在黑板上(定时任务聚合 + 缓存),顾客看黑板就行。
- CDN:把黑板照片贴到门口展示屏上,路过的人不用进店就能看(CDN 静态化)。
- 异步加载:排行榜先展示"加载中…",后台慢慢查,查到了再填充(前端分批请求)。
- 物化视图:直接在数据库里建一张"每日销量汇总表"(物化视图),查询时直接查汇总表,不用实时聚合。
桥接: 餐厅类比映射到数据统计加速——预计算是定时聚合写入缓存,CDN 是静态资源边缘缓存,异步加载是前端分批渲染避免白屏,物化视图是数据库预聚合表。
3.2 工程要点
| 加速方案 | 原理 | 优点 | 缺点 | 适用场景 |
|---|---|---|---|---|
| 预计算 | 定时任务聚合结果写入 Redis | 查询极快(O(1)) | 数据有延迟 | 实时性要求不高的统计 |
| Redis 缓存 | 查询结果缓存,命中直接返回 | 减少 DB 压力 | 缓存一致性 | 读多写少的统计 |
| CDN | 静态页面推到边缘节点 | 就近访问,延迟极低 | 只适合静态内容 | 全国分布的展示页 |
| 异步加载 | 前端分批请求,先渲染骨架屏 | 首屏快,体验好 | 总加载时间不变 | 多模块仪表盘 |
| 物化视图 | 数据库预聚合表 | 查询直接读汇总表 | 占额外存储 | 复杂 SQL 聚合 |
| 读写分离 | 读请求走从库 | 分散主库压力 | 主从延迟 | 高并发读 |
Go 定时预计算示例:
package main
import (
"context"
"fmt"
"time"
"github.com/redis/go-redis/v9"
)
type StatsService struct {
rdb *redis.Client
}
// 步骤1:定时预计算热门内容排行榜
func (s *StatsService) PrecomputeHotRanking(ctx context.Context) {
ticker := time.NewTicker(5 * time.Minute)
defer ticker.Stop()
for range ticker.C {
// 步骤2:从数据库聚合最近24小时的热门内容
// SQL: SELECT content_id, COUNT(*) as cnt FROM views
// WHERE created_at > NOW() - INTERVAL 24 HOUR
// GROUP BY content_id ORDER BY cnt DESC LIMIT 100
// 步骤3:将结果写入Redis ZSet(模拟数据)
hotKey := "stats:hot:24h"
items := []redis.Z{
{Score: 9821, Member: "content:1001"},
{Score: 8732, Member: "content:1002"},
{Score: 7621, Member: "content:1003"},
}
s.rdb.ZAdd(ctx, hotKey, items...)
// 步骤4:设置过期时间,防止定时任务挂了数据不更新
s.rdb.Expire(ctx, hotKey, 10*time.Minute)
fmt.Println("热门排行榜已更新")
}
}
// 步骤5:查询时直接读缓存
func (s *StatsService) GetHotRanking(ctx context.Context, limit int64) ([]string, error) {
// 步骤6:先查Redis缓存
hotKey := "stats:hot:24h"
result, err := s.rdb.ZRevRange(ctx, hotKey, 0, limit-1).Result()
if err == nil && len(result) > 0 {
return result, nil
}
// 步骤7:缓存未命中,回源查DB(实际项目中这里查物化视图)
// 注意加互斥锁防止缓存击穿(见第四节)
return []string{"content:fallback"}, nil
}
func main() {
rdb := redis.NewClient(&redis.Options{Addr: "localhost:6379"})
svc := &StatsService{rdb: rdb}
ctx := context.Background()
go svc.PrecomputeHotRanking(ctx)
ranking, _ := svc.GetHotRanking(ctx, 10)
fmt.Println("热门排行:", ranking)
}
⚠️ 新手必踩的坑: 预计算的缓存一定要设置过期时间。如果定时任务挂了,没有过期时间的缓存会一直返回旧数据,用户看到的永远是过期统计。设置过期时间后,即使定时任务挂了,缓存过期后会触发回源查询,至少能返回最新数据。
四、缓存雪崩:热数据防护
4.1 用生活类比先建立直觉
缓存雪崩就像超市促销引发踩踏: 一家超市把特价商品信息写在门口黑板上(缓存)。某天黑板被擦了(缓存过期),所有顾客同时涌进超市问"特价商品在哪"(请求打到 DB),超市被挤垮(DB 崩溃)。更糟糕的是,超市重新写黑板的速度跟不上(缓存重建慢),后面的顾客继续涌入,形成恶性循环。
如果是爬虫热数据就更可怕:爬虫以每秒万次的频率访问同一个 key,当这个 key 过期的瞬间,万次请求同时打到 DB,DB 瞬间崩溃。
graph TB
K[热key过期] --> A1[请求1打到DB]
K --> A2[请求2打到DB]
K --> A3[请求3打到DB]
K --> A4[请求N打到DB]
A1 --> DB[(数据库
瞬间崩溃)]
A2 --> DB
A3 --> DB
A4 --> DB
DB -->|崩溃后缓存无法重建| K
DB -.->|形成恶性循环| K桥接: 超市类比映射到缓存雪崩——黑板擦除对应热 key 过期,顾客涌入对应请求穿透到 DB,恶性循环对应 DB 崩溃后缓存无法重建。解决方案:只让一个人去问特价信息然后写到黑板上(互斥锁),或者特价信息永不擦除(热点 key 永不过期),或者每个区域都有小黑板(多级缓存)。
4.2 工程要点
缓存雪崩 vs 缓存击穿 vs 缓存穿透
| 问题 | 原因 | 解决方案 |
|---|---|---|
| 缓存雪崩 | 大量 key 同时过期 | 过期时间加随机值、多级缓存 |
| 缓存击穿 | 单个热 key 过期 | 互斥锁重建、热点 key 永不过期 |
| 缓存穿透 | 查询不存在的数据 | 布隆过滤器、缓存空值 |
缓存雪崩解决方案
graph TB
subgraph 解决方案
S1[互斥锁重建
只放一个请求查DB]
S2[热点key永不过期
逻辑过期后台更新]
S3[多级缓存
本地缓存+Redis+DB]
S4[布隆过滤器
过滤不存在的key]
S5[过期时间加随机
避免同时失效]
end
S1 --> R[防止DB被压垮]
S2 --> R
S3 --> R
S4 --> R
S5 --> RGo 互斥锁防缓存击穿代码:
package main
import (
"context"
"fmt"
"math/rand"
"sync"
"time"
"github.com/redis/go-redis/v9"
)
type CacheService struct {
rdb *redis.Client
mu sync.Mutex
}
// 步骤1:带互斥锁的缓存查询
func (s *CacheService) GetWithMutex(ctx context.Context, key string) (string, error) {
// 步骤2:先查Redis缓存
val, err := s.rdb.Get(ctx, key).Result()
if err == nil {
return val, nil // 缓存命中
}
// 步骤3:缓存未命中,加互斥锁(只放一个请求去查DB)
s.mu.Lock()
defer s.mu.Unlock()
// 步骤4:双重检查(可能其他协程已经重建了缓存)
val, err = s.rdb.Get(ctx, key).Result()
if err == nil {
return val, nil
}
// 步骤5:查数据库
val = s.queryFromDB(key)
// 步骤6:回写缓存,过期时间加随机值防止雪崩
ttl := 300 + rand.Intn(60) // 300到360秒随机
s.rdb.Set(ctx, key, val, time.Duration(ttl)*time.Second)
return val, nil
}
// 模拟数据库查询
func (s *CacheService) queryFromDB(key string) string {
time.Sleep(100 * time.Millisecond) // 模拟DB查询耗时
return fmt.Sprintf("db_value_for_%s", key)
}
func main() {
rdb := redis.NewClient(&redis.Options{Addr: "localhost:6379"})
svc := &CacheService{rdb: rdb}
ctx := context.Background()
// 模拟并发请求
var wg sync.WaitGroup
for i := 0; i < 100; i++ {
wg.Add(1)
go func(id int) {
defer wg.Done()
val, _ := svc.GetWithMutex(ctx, "hotkey:product:1001")
fmt.Printf("goroutine %d 得到: %s\n", id, val)
}(i)
}
wg.Wait()
}
⚠️ 新手必踩的坑: 互斥锁要用进程级的
sync.Mutex,不是 Redis 的SETNX。因为SETNX加锁后如果持锁进程崩溃,锁不会自动释放(虽然有过期时间但期间其他请求都在等)。进程内sync.Mutex更轻量,同一进程内只有一个请求去查 DB。如果是多进程部署,可以结合 Redis 分布式锁,但要给锁设短超时(如 3 秒)。
五、长短链转换
5.1 用生活类比先建立直觉
长短链转换就像快递柜取件码: 你有一个很长的快递地址(长链接 https://www.example.com/products/detail?id=1001&source=wechat&campaign=summer_sale),不方便发给别人。快递公司给你一个 6 位取件码(短链 https://s.cn/Ab3x9K),别人输入取件码就能查到完整地址。快递公司内部有一个映射表(Redis),记录取件码和完整地址的对应关系。
301 vs 302 就像永久搬家 vs 临时出差:
- 301 永久重定向:你永久搬家了,邮局把你的地址更新到通讯录。以后所有信都直接寄到新地址,不再经过老地址。浏览器会缓存这个映射,下次直接访问新地址,不再请求短链服务器。
- 302 临时重定向:你临时出差了,每次有人寄信,邮局都先查一下你现在的地址再转发。浏览器不缓存,每次都经过短链服务器。
为什么短链用 302 不用 301?因为短链服务需要每次都统计访问量(UV/PV/来源),如果用 301 浏览器缓存了映射,后续访问就不经过短链服务器了,统计数据不准。
graph LR
subgraph 短链生成
L[长链接] --> G[发号器
生成唯一ID]
G --> E[Base62编码
ID转短码]
E --> R[Redis存储映射
短码到长链]
R --> S[返回短链
s.cn/Ab3x9K]
end
subgraph 短链跳转
U[用户访问短链] --> Q[查Redis映射]
Q --> D{查到?}
D -->|是| H[302重定向到长链]
D -->|否| E404[返回404]
H --> ST[记录访问统计]
end桥接: 快递柜类比映射到短链转换——取件码生成对应发号器+Base62编码,映射表对应 Redis 存储,取件查询对应短链重定向。301/302 的选型取决于是否需要每次都经过服务器统计。
5.2 工程要点
Base62 编码原理
Base62 用 62 个字符(0-9, A-Z, a-z)表示数字,比 Base10 紧凑得多:
| 进制 | 字符集 | ID=1000000000 的编码 | 长度 |
|---|---|---|---|
| Base10 | 0-9 | 1000000000 | 10 位 |
| Base62 | 0-9A-Za-z | 15NYdN | 6 位 |
6 位 Base62 可表示 62^6 = 568 亿个短链,足够大多数场景。
Go 实现长短链转换:
package main
import (
"context"
"fmt"
"net/http"
"github.com/redis/go-redis/v9"
)
const base62Chars = "0123456789ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz"
type ShortURLService struct {
rdb *redis.Client
}
// 步骤1:将长ID编码为Base62短码
func EncodeBase62(id int64) string {
if id == 0 {
return string(base62Chars[0])
}
var result []byte
for id > 0 {
remainder := id % 62
result = append([]byte{base62Chars[remainder]}, result...)
id = id / 62
}
return string(result)
}
// 步骤2:生成短链
func (s *ShortURLService) Generate(ctx context.Context, longURL string) (string, error) {
// 步骤3:先检查长链是否已有短链映射(避免重复生成)
existKey := fmt.Sprintf("long:%s", longURL)
if shortCode, err := s.rdb.Get(ctx, existKey).Result(); err == nil {
return "https://s.cn/" + shortCode, nil
}
// 步骤4:通过发号器获取全局唯一ID(Redis INCR模拟)
id, err := s.rdb.Incr(ctx, "shorturl🆔seq").Result()
if err != nil {
return "", err
}
// 步骤5:Base62编码生成短码
shortCode := EncodeBase62(id)
// 步骤6:存储双向映射
s.rdb.Set(ctx, fmt.Sprintf("short:%s", shortCode), longURL, 0)
s.rdb.Set(ctx, existKey, shortCode, 0)
return "https://s.cn/" + shortCode, nil
}
// 步骤7:短链重定向
func (s *ShortURLService) Redirect(w http.ResponseWriter, r *http.Request) {
shortCode := r.URL.Path[1:] // 去掉前导/
ctx := context.Background()
// 步骤8:查Redis获取长链
longURL, err := s.rdb.Get(ctx, fmt.Sprintf("short:%s", shortCode)).Result()
if err != nil {
http.NotFound(w, r)
return
}
// 步骤9:记录访问统计(异步)
go s.rdb.Incr(ctx, fmt.Sprintf("stats:%s", shortCode))
// 步骤10:302临时重定向(不用301,确保每次都经过服务器统计)
http.Redirect(w, r, longURL, http.StatusFound)
}
func main() {
rdb := redis.NewClient(&redis.Options{Addr: "localhost:6379"})
svc := &ShortURLService{rdb: rdb}
ctx := context.Background()
// 测试生成短链
shortURL, _ := svc.Generate(ctx, "https://www.example.com/products/detail?id=1001&source=wechat")
fmt.Println("短链:", shortURL)
// 启动HTTP服务
http.HandleFunc("/", svc.Redirect)
fmt.Println("短链服务启动 :9090")
http.ListenAndServe(":9090", nil)
}
⚠️ 新手必踩的坑: 发号器不能用数据库自增 ID 直接暴露给用户——短码如果是连续的(Ab3x9K 到 Ab3x9L 到 Ab3x9M),竞品可以通过遍历短码爬取你的全部内容。解决方案:用雪花算法生成不连续的 ID,或者在 ID 上做位混淆(如高低位交换)。
六、扫码登录
6.1 用生活类比先建立直觉
扫码登录就像用门禁卡开门:
- 你走到门前,门禁机显示一个二维码(PC 端生成二维码)。
- 你掏出手机扫码,手机上弹出"确认开门?"(手机扫码确认)。
- 你点"确认",门禁机屏幕显示"开门成功"(PC 端收到登录成功通知)。
关键问题:门禁机怎么知道你扫了码?两种方式:
- 轮询:门禁机每隔 2 秒问一次服务器"有人扫码了吗?"——简单但浪费资源。
- WebSocket:门禁机和服务器保持长连接,有人扫码时服务器主动通知——实时但复杂。
sequenceDiagram
participant PC as PC浏览器
participant S as 服务器
participant M as 手机APP
PC->>S: 请求登录二维码
S->>PC: 返回二维码含临时token
Note over S: 存储token状态=等待扫码
M->>S: 扫码提交token+用户身份
S->>M: 确认登录请输入密码
M->>S: 确认登录
Note over S: token状态=已确认
PC->>S: 轮询token登录了吗
S->>PC: 已确认返回登录凭证
Note over PC: 登录成功桥接: 门禁类比映射到扫码登录——二维码对应临时 token,扫码确认对应手机端授权,门禁机显示结果对应 PC 端轮询或 WebSocket 通知。核心是 PC 端和手机端通过服务器中转,通过共享的临时 token 关联两个会话。
6.2 工程要点
扫码登录状态机
| 状态 | 含义 | 触发条件 |
|---|---|---|
| WAITING_SCAN | 等待扫码 | PC 端生成二维码 |
| SCANNED | 已扫码待确认 | 手机端扫码成功 |
| CONFIRMED | 已确认登录 | 手机端点击确认 |
| EXPIRED | 已过期 | 超过 60 秒未操作 |
| CANCELED | 已取消 | 手机端取消登录 |
Go 实现扫码登录核心逻辑:
package main
import (
"context"
"encoding/json"
"fmt"
"net/http"
"time"
"github.com/google/uuid"
"github.com/redis/go-redis/v9"
)
type QrLoginService struct {
rdb *redis.Client
}
// 步骤1:PC端请求生成二维码
func (s *QrLoginService) GenerateQrCode(ctx context.Context) (string, error) {
// 步骤2:生成临时token(唯一标识本次登录会话)
token := uuid.New().String()
// 步骤3:存储token状态,设为"等待扫码",60秒过期
statusKey := fmt.Sprintf("qrlogin:%s:status", token)
s.rdb.Set(ctx, statusKey, "WAITING_SCAN", 60*time.Second)
return token, nil
}
// 步骤4:手机端扫码
func (s *QrLoginService) Scan(ctx context.Context, token string, userId string) error {
statusKey := fmt.Sprintf("qrlogin:%s:status", token)
// 步骤5:检查token是否有效
status, err := s.rdb.Get(ctx, statusKey).Result()
if err != nil {
return fmt.Errorf("二维码已过期")
}
if status != "WAITING_SCAN" {
return fmt.Errorf("二维码状态异常: %s", status)
}
// 步骤6:更新状态为"已扫码待确认",记录用户ID
s.rdb.Set(ctx, statusKey, "SCANNED", 60*time.Second)
s.rdb.Set(ctx, fmt.Sprintf("qrlogin:%s:user", token), userId, 60*time.Second)
return nil
}
// 步骤7:手机端确认登录
func (s *QrLoginService) Confirm(ctx context.Context, token string) error {
statusKey := fmt.Sprintf("qrlogin:%s:status", token)
status, _ := s.rdb.Get(ctx, statusKey).Result()
if status != "SCANNED" {
return fmt.Errorf("请先扫码")
}
// 步骤8:更新状态为"已确认"
s.rdb.Set(ctx, statusKey, "CONFIRMED", 30*time.Second)
return nil
}
// 步骤9:PC端轮询登录状态
func (s *QrLoginService) Poll(ctx context.Context, token string) (string, string, error) {
statusKey := fmt.Sprintf("qrlogin:%s:status", token)
status, err := s.rdb.Get(ctx, statusKey).Result()
if err != nil {
return "EXPIRED", "", fmt.Errorf("二维码已过期")
}
// 步骤10:如果已确认,返回登录凭证
if status == "CONFIRMED" {
userId, _ := s.rdb.Get(ctx, fmt.Sprintf("qrlogin:%s:user", token)).Result()
// 生成登录token(JWT等)
authToken := generateAuthToken(userId)
return status, authToken, nil
}
return status, "", nil
}
func generateAuthToken(userId string) string {
return fmt.Sprintf("jwt_token_for_%s", userId)
}
func main() {
rdb := redis.NewClient(&redis.Options{Addr: "localhost:6379"})
svc := &QrLoginService{rdb: rdb}
ctx := context.Background()
// 模拟完整流程
token, _ := svc.GenerateQrCode(ctx)
fmt.Println("1. PC端生成二维码, token:", token)
svc.Scan(ctx, token, "user:1001")
fmt.Println("2. 手机端扫码成功")
status, _, _ := svc.Poll(ctx, token)
fmt.Println("3. PC端轮询状态:", status)
svc.Confirm(ctx, token)
fmt.Println("4. 手机端确认登录")
status, authToken, _ := svc.Poll(ctx, token)
fmt.Printf("5. PC端轮询状态: %s, 登录凭证: %s\n", status, authToken)
// HTTP接口示例
http.HandleFunc("/qr/generate", func(w http.ResponseWriter, r *http.Request) {
token, _ := svc.GenerateQrCode(r.Context())
json.NewEncoder(w).Encode(map[string]string{"token": token})
})
http.HandleFunc("/qr/poll", func(w http.ResponseWriter, r *http.Request) {
token := r.URL.Query().Get("token")
status, authToken, _ := svc.Poll(r.Context(), token)
json.NewEncoder(w).Encode(map[string]string{
"status": status,
"token": authToken,
})
})
}
⚠️ 新手必踩的坑: 二维码的临时 token 必须设置过期时间(通常 60 秒)。如果不过期,攻击者可以截获二维码图片后任意时间扫码登录。同时 PC 端轮询间隔建议 2 秒,太短浪费带宽,太长体验差。如果用 WebSocket 替代轮询,要注意 WebSocket 断线重连逻辑。
七、设计微信的海量数据存储系统
7.1 用生活类比先建立直觉
微信的数据存储就像一座超大型城市档案馆: 一个 10 亿用户的城市,每个人每天都在产生信件(聊天消息)、照片(朋友圈图片)、户口本(好友关系)、个人档案(资料)。档案馆不能把全市 10 亿人的所有资料堆在一个大仓库——仓库会塌(单机存不下),找一份资料要翻遍整个仓库(查询慢),仓库着火全市数据都没了(无容灾)。
于是档案馆做了几件事:
- 按街区建分馆(数据分片):每个用户的资料固定存放在某个分馆,按用户 ID 编号决定去哪个分馆(一致性哈希 / 取模分片)。
- 冷热库房分离(冷热分离):最近三个月的活跃信件放恒温库随时查(热数据用 SSD + 内存),十年前的旧信打包进地下冷库(冷数据用廉价大容量存储)。
- 每个分馆都留副本(多副本):主分馆着火,隔壁分馆立刻顶上(多副本 + 主从切换)。
- 照片单独存照相馆(对象存储):聊天里的图片视频不塞进信箱,而是存到专门的照相馆,信里只留一个取件码(对象存储 + CDN)。
graph TB
subgraph 数据分类
D1[聊天消息
高频写/按会话读]
D2[好友关系链
读多写少]
D3[朋友圈Feed
写扩散读多]
D4[图片视频
大对象低频]
D5[用户资料
读极多写少]
end
D1 --> S1[分布式KV
按会话分片]
D2 --> S2[关系/图存储
按用户ID分片]
D3 --> S3[收件箱模型
推模式ZSet]
D4 --> S4[对象存储+CDN]
D5 --> S5[内存KV
多级缓存]桥接: 城市档案馆类比映射到微信存储——分馆对应数据分片(按用户 ID 取模或一致性哈希),冷热库房对应冷热分离(近期消息 SSD、历史消息压缩归档),副本对应多副本容灾,照相馆对应对象存储。核心矛盾是"单库存不下 + 单点故障 + 查询热点",所以海量存储的第一步永远是"先给数据分类"。
7.2 工程要点
数据类型与存储选型对照
| 数据类型 | 访问特征 | 存储选型 | 分片方式 |
|---|---|---|---|
| 聊天消息 | 高频写入、按会话拉取 | 分布式 KV(自研 Peanut / QuorumKV) | 按会话 ID 哈希 |
| 好友关系链 | 读多写少 | 关系型 / 图存储(按 user 分库) | 按 userID 取模分 1000 库 |
| 朋友圈 Feed | 写扩散、读多 | Redis 收件箱 + 离线存储 | 推模式写 inbox |
| 图片 / 视频 | 大对象、低频 | 对象存储(COS)+ CDN | 按上传时间分桶 |
| 用户资料 | 读极多、写少 | 内存 KV + 多级缓存 | 按 userID 哈希 |
核心设计原则
- 数据分片:单表 / 单库上限约千万级,微信按
userID % 1000拆成 1000 个库,每个库再按sessionID拆表。分片键选择要避免热点(好友关系按 userID 分片,保证单用户的关系都在同一分片,减少跨分片查询)。 - 一致性哈希:扩容时只用迁移约
1/N的数据,而非全量重哈希,节点变化影响最小。 - 冷热分离:三个月内消息存高速 KV,超期消息压缩后迁移到廉价存储(列式 / 归档),查询时按需回源。IM 场景"最近会话"访问占 90%,冷数据很少被读。
- 多副本 + 异地多活:每份数据 3 副本(2 副本同城不同机房 + 1 副本异地),写入走 Quorum(如 2/3 成功即返回),保证单机房故障不丢数据。
- 对象存储独立:图片视频走对象存储,消息体只存 URL,避免大对象拖慢消息 KV。
Go 实现一致性哈希分片路由示例:
package main
import (
"hash/crc32"
"sort"
"strconv"
)
// 步骤1:一致性哈希环(简化版,用虚拟节点减少数据倾斜)
type ConsistentHash struct {
ring []uint32 // 有序的哈希值列表
nodes map[uint32]int // 哈希值 -> 节点编号
vnodes int // 每个物理节点的虚拟节点数
}
// 步骤2:添加节点(每个物理节点映射到 vnodes 个虚拟点)
func (c *ConsistentHash) AddNode(nodeID int) {
for i := 0; i < c.vnodes; i++ {
key := crc32.ChecksumIEEE([]byte(strconv.Itoa(nodeID) + "#" + strconv.Itoa(i)))
c.ring = append(c.ring, key)
c.nodes[key] = nodeID
}
sort.Slice(c.ring, func(i, j int) bool { return c.ring[i] < c.ring[j] })
}
// 步骤3:根据用户ID找到归属的存储分片
func (c *ConsistentHash) GetNode(userID string) int {
h := crc32.ChecksumIEEE([]byte(userID))
// 步骤4:在有序环上二分查找第一个 >= h 的虚拟节点
idx := sort.Search(len(c.ring), func(i int) bool { return c.ring[i] >= h })
if idx == len(c.ring) {
idx = 0 // 环形回绕
}
return c.nodes[c.ring[idx]]
}
func main() {
ch := &ConsistentHash{vnodes: 100}
for i := 0; i < 10; i++ {
ch.AddNode(i) // 10 个存储分片(实际是 1000 个库)
}
// 步骤5:user:1001 的聊天消息固定路由到某个分片
node := ch.GetNode("user:1001")
println("user:1001 的消息存到分片", node)
}
⚠️ 新手必踩的坑: 分片键一旦选定几乎无法更改。如果选了"按群组 ID 分片"存好友关系,那"查某个用户的所有好友"就变成跨全部分片的查询(扇出爆炸)。IM 关系链务必按
userID分片,把"单用户维度"的查询收敛到单分片。扩容时的数据迁移也要评估——一致性哈希能减少迁移量,但仍有1/N的数据需要搬迁,必须灰度进行,避免一次性迁移打垮集群。
面试加分点(开放性问题的表达框架)
回答"海量数据存储"开放题时按"四步走"组织:① 先分类数据(消息 / 关系 / 媒体 / 资料),不同数据用不同存储;② 讲分片与路由(分片键、一致性哈希、扩容迁移);③ 讲可靠与容灾(多副本、Quorum 写、异地多活);④ 讲成本与冷热分离。不要上来就堆技术名词,先证明你"分得清数据",再谈"怎么存"。
八、爬虫系统设计与实现流程
8.1 用生活类比先建立直觉
爬虫系统就像"图书馆巡架采编员": 想象一座超大的图书馆,你需要把所有书架上的书登记进电脑(抓取全站数据)。一个人跑断腿也干不完(单机吞吐不够),而且同一本书可能被多个书架重复摆着(重复 URL),登记两遍纯属浪费。
于是图书馆把工作拆成了流水线:
- 调度器(Scheduler):馆长手里有一张"待采编书单"(待抓 URL 队列),决定下一本抓哪本。他先看"已采编目录"(去重器)里有没有这本书,没有才放进书单。
- 下载器(Downloader):跑腿员按书单去书架上把书取回来(发 HTTP 请求拿到网页)。
- 解析器(Parser):录入员翻开书,把正文抄下来(抽取数据),再把书里引用的"参考书目"写回书单(抽取新 URL 继续抓)。
- 去重器(Dedup):门口的查重机,扫一下书号就知道这本书采编过没(Bloom Filter / Redis Set)。
- 存储器(Storage):档案室,把抄好的正文归档入库(DB / 对象存储)。
graph TB
subgraph 单机爬虫流水线
S1[种子URL] --> SCH[调度器
URL队列+去重]
SCH -->|未抓过的URL| D1[下载器
发HTTP请求]
D1 --> P1[解析器
抽数据+抽新URL]
P1 -->|新URL回灌| SCH
P1 -->|正文数据| ST1[存储器
DB/对象存储]
end
subgraph 分布式爬虫-消息队列解耦
MS[Master调度节点] -->|推送任务| MQ[(消息队列
Kafka/NSQ)]
MQ --> W1[Worker1
下载+解析]
MQ --> W2[Worker2
下载+解析]
MQ --> W3[WorkerN
下载+解析]
W1 -->|去重查询| RED[(Redis Set
/Bloom Filter)]
W2 --> RED
W3 --> RED
W1 -->|数据| ST2[(分布式存储)]
W2 --> ST2
W3 --> ST2
end桥接: 图书馆类比映射到爬虫系统——待采编书单对应 URL 队列,查重机对应去重器,跑腿员对应下载器,录入员对应解析器,档案室对应存储。分布式爬虫的核心是用消息队列把"调度"和"抓取"解耦:Master 只管往队列里塞任务,N 个 Worker 各自消费,谁挂了任务还在队列里,不影响整体进度。
8.2 工程要点
整体架构与模块职责
| 模块 | 技术选型 | 设计要点 |
|---|---|---|
| 调度器 | 待抓队列(Redis List / Kafka) | 控制抓取顺序、频率、优先级;与去重器联动 |
| 下载器 | Go net/http + 连接池 | 复用 TCP 连接;设置超时与重试;承载反爬策略(见第九节) |
| 解析器 | 正则 / goquery / json | 抽取目标字段与下一层 URL;.xpath 或 CSS 选择器 |
| 去重器 | Bloom Filter / Redis Set | 海量 URL 用 Bloom Filter 省内存;精确去重用 Redis Set |
| 存储器 | MySQL / ES / 对象存储 | 结构化数据入库,大文件走对象存储 |
分布式爬虫的关键设计
- 多机调度:Master 节点从种子 URL 出发,把待抓任务投递到消息队列;多个 Worker 节点竞争消费。天然支持水平扩容——加机器就是加 Worker。
- 消息队列解耦:调度与抓取通过队列解耦,Worker 崩溃不丢任务(队列持久化),Master 也不用等某个 Worker 慢响应。还能做优先级队列(重要站点优先)。
- 去重:Bloom Filter vs Redis Set
- Redis Set 精确去重,但百亿 URL 内存爆炸(每个 URL 几十字节 × 百亿 ≈ 数百 GB)。
- Bloom Filter 用位数组 + 多个哈希函数,判断"一定没抓过 / 可能抓过",有极小的误判率(把没抓过的判成抓过,漏抓一条不影响大局),但内存只需 Set 的几十分之一。实战中常用 Redis Set 做第一层精确去重、Bloom Filter 做前置快速过滤。
Go 实现单机爬虫(Worker Pool + Bloom Filter 去重):
package main
import (
"fmt"
"hash/fnv"
"net/url"
"sync"
)
// 步骤1:极简 Bloom Filter(生产用 redisbloom 或 roaring bitmap)
type BloomFilter struct {
bits []uint64
size uint64
}
func NewBloomFilter(size uint64) *BloomFilter {
return &BloomFilter{bits: make([]uint64, (size+63)/64), size: size}
}
// 步骤2:两个哈希函数映射到位数组
func (b *BloomFilter) hashes(s string) (uint64, uint64) {
h1 := fnv.New32a()
h1.Write([]byte(s))
v1 := uint64(h1.Sum32()) % b.size
h2 := fnv.New32a()
h2.Write([]byte(s + "salt"))
v2 := uint64(h2.Sum32()) % b.size
return v1, v2
}
func (b *BloomFilter) Add(s string) {
h1, h2 := b.hashes(s)
b.bits[h1/64] |= 1 << (h1 % 64)
b.bits[h2/64] |= 1 << (h2 % 64)
}
// 步骤3:可能返回 false(一定没抓过)/ true(可能抓过)
func (b *BloomFilter) MaybeContains(s string) bool {
h1, h2 := b.hashes(s)
if b.bits[h1/64]&(1<<(h1%64)) == 0 {
return false
}
if b.bits[h2/64]&(1<<(h2%64)) == 0 {
return false
}
return true
}
// 步骤4:爬虫 Worker
type Crawler struct {
bf *BloomFilter
visited sync.Map // 精确去重兜底(演示用)
mu sync.Mutex
}
func (c *Crawler) ShouldCrawl(rawURL string) bool {
// 步骤5:先问 Bloom Filter 快速过滤
if !c.bf.MaybeContains(rawURL) {
return true // 一定没抓过
}
// 步骤6:可能抓过,再查精确表兜底
if _, ok := c.visited.Load(rawURL); ok {
return false
}
return true
}
func (c *Crawler) Crawl(rawURL string) {
c.mu.Lock()
c.bf.Add(rawURL)
c.visited.Store(rawURL, struct{}{})
c.mu.Unlock()
fmt.Println("下载并解析:", rawURL)
}
func main() {
c := &Crawler{bf: NewBloomFilter(1 << 20)}
seed := []string{
"https://example.com/a",
"https://example.com/b",
"https://example.com/a", // 重复 URL
}
// 步骤7:多 Worker 消费 URL 队列(用 channel 模拟)
queue := make(chan string, len(seed))
for _, u := range seed {
queue <- u
}
close(queue)
var wg sync.WaitGroup
for i := 0; i < 3; i++ { // 3 个 Worker
wg.Add(1)
go func(id int) {
defer wg.Done()
for u := range queue {
if c.ShouldCrawl(u) {
c.Crawl(u)
} else {
fmt.Printf("Worker%d 跳过重复: %s\n", id, u)
}
}
}(i)
}
wg.Wait()
// 步骤8:校验 URL 合法性(解析器前过滤非法链接)
if parsed, err := url.Parse("https://example.com/c"); err == nil {
fmt.Println("解析出下一层URL:", parsed.Host)
}
}
⚠️ 新手必踩的坑: 爬虫必须设置
robots.txt遵守与抓取限速,否则既可能违法也可能被目标站拉黑。另外解析出的新 URL 一定要做"同域名 / 白名单"过滤,否则一个页面里含有指向外站的链接,你的爬虫就会"爬出国",流量和存储都失控。分布式场景下去重用 Redis 集中存储,避免每个 Worker 各记各的、互相不知道对方抓过没。
九、反爬技术实现
9.1 用生活类比先建立直觉
反爬就像"小区门禁 + 保安盘查": 你是快递员(爬虫)想进小区送包裹(抓数据)。小区为了防止被滥用,设了一道道关卡:
- UA 伪装:保安看你"穿什么制服"(User-Agent)。你如果穿着奇怪的工作服(默认
Go-http-client/1.1),保安直接拦下。你得换上"顺丰快递"制服(伪装成浏览器 UA)才放行进。 - IP 代理池:保安发现"这个门牌号的人一天来 1000 次"(同一 IP 高频访问),拉黑。你得多备几套"假门牌号"(代理 IP 池),每次换一个进。
- Cookie 池:进楼还要刷门禁卡(Cookie/登录态)。你养一批"已注册住户"的门禁卡(Cookie 池),轮流刷。
- 请求限速:再合法的快递员也不能一秒敲 100 次门(限速),否则扰民被投诉。
- 验证码识别:偶尔保安让你做道题"证明你是人"(验证码),你用 OCR / 打码平台答。
- 字体反爬:最阴的——小区把"门牌号 3"故意印成看起来像"8"的特殊字体(自定义字体映射),你抄下来的"8"其实是"3",数据全错。
graph LR
C[爬虫请求] --> UA{UA是否浏览器?}
UA -->|否,Go默认| B1[直接拉黑]
UA -->|是| IP{IP是否被限频?}
IP -->|高频| B2[封IP]
IP -->|正常| CK{有无有效Cookie?}
CK -->|无| B3[要求登录/验证码]
CK -->|有| RL{是否超限速?}
RL -->|超| B4[429限流]
RL -->|正常| OK[返回数据
可能字体加密]
OK --> FONT{是否字体反爬?}
FONT -->|是| DEC[需本地映射表解密]
FONT -->|否| DATA[拿到真实数据]桥接: 小区门禁类比映射到反爬——UA 伪装是换"制服",代理池是换"门牌号",Cookie 池是刷"门禁卡",限速是别"敲门太狠",验证码是"做人机验证",字体反爬是"抄错门牌号"。面试官常问"你项目里反爬几项技术怎么实现的",回答时按"UA → IP → Cookie → 限速 → 验证码 → 字体"逐层讲,再点出哪层最容易被突破。
9.2 工程要点
各项反爬技术的实现思路
| 反爬手段 | 爬虫应对 | Go 实现要点 |
|---|---|---|
| UA 伪装 | 随机浏览器 UA 池 | 请求前随机选一个 UA 写入 req.Header |
| IP 代理池 | 代理 IP 列表轮询 / 健康探测 | 自定义 http.Transport 的 Proxy 字段 |
| Cookie 池 | 预登录批量养号 | 多账号 Cookie 存 Redis,按 key 轮询取用 |
| 请求限速 | 令牌桶 / 滑动窗口限速 | golang.org/x/time/rate 每个域名一个 limiter |
| 验证码 | OCR / 打码平台 / 人工 | 截图送识别服务,拿回结果回填表单 |
| 字体反爬 | 下载字体文件解析映射 | 解析 woff 提取 unicode→glyph 映射表还原 |
结合"项目中反爬的几项技术怎么实现的"
实战回答框架(照着讲即可):
- UA 池 + 代理池:启动时从配置加载几百个浏览器 UA 和代理 IP,每次请求用
rand随机挑,代理池后台定时探测可用性,失效的剔除。 - 限速:用
golang.org/x/time/rate,对每个域名建一个 limiter(如每秒 2 次),避免把单个站点打死也避免被封。 - Cookie 池:登录态过期前批量刷新,存 Redis Hash,抓需要登录的页面时取一个 Cookie 用。
- 验证码:遇验证码把页面截图丢给打码平台(或自建 OCR),拿回 token 继续;频率不高时也可降级为人工介入。
- 字体反爬兜底:遇到页面数字是乱码,先下载其 woff 字体,解析出"字形→真实字符"的映射表,抓取后统一做一次翻译还原。
Go 实现带 UA 池 + 代理池 + 限速的下载器:
package main
import (
"crypto/tls"
"fmt"
"math/rand"
"net/http"
"net/url"
"time"
"golang.org/x/time/rate"
)
// 步骤1:UA 池
var uaPool = []string{
"Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0 Safari/537.36",
"Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/605.1.15 (KHTML, like Gecko) Version/17.0 Safari/605.1.15",
"Mozilla/5.0 (X11; Linux x86_64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/119.0 Safari/537.36",
}
// 步骤2:代理池(演示用,实际从代理服务拉取并探活)
var proxyPool = []string{
"http://1.2.3.4:8080",
"http://5.6.7.8:8080",
"http://9.10.11.12:8080",
}
// 步骤3:每域限速器(域名 -> limiter)
var limiters = map[string]*rate.Limiter{}
func getLimiter(host string) *rate.Limiter {
if l, ok := limiters[host]; ok {
return l
}
l := rate.NewLimiter(rate.Limit(2), 1) // 每秒 2 次,突发 1
limiters[host] = l
return l
}
// 步骤4:构造带反爬策略的 Client
func newAntiCrawlClient() *http.Client {
return &http.Client{
Timeout: 10 * time.Second,
Transport: &http.Transport{
TLSClientConfig: &tls.Config{InsecureSkipVerify: true},
// 步骤5:每次请求随机选代理
Proxy: func(req *http.Request) (*url.URL, error) {
p := proxyPool[rand.Intn(len(proxyPool))]
return url.Parse(p)
},
},
}
}
// 步骤6:带限速 + UA 伪装的一次抓取
func fetch(client *http.Client, target string) (*http.Response, error) {
u, _ := url.Parse(target)
// 步骤7:先等限速器放行(按域名限速)
if err := getLimiter(u.Host).Wait(nil); err != nil {
return nil, err
}
req, _ := http.NewRequest("GET", target, nil)
// 步骤8:随机 UA 伪装
req.Header.Set("User-Agent", uaPool[rand.Intn(len(uaPool))])
req.Header.Set("Accept", "text/html,application/xhtml+xml")
req.Header.Set("Accept-Language", "zh-CN,zh;q=0.9")
return client.Do(req)
}
func main() {
client := newAntiCrawlClient()
resp, err := fetch(client, "https://example.com/page/1")
if err != nil {
fmt.Println("抓取失败:", err)
return
}
defer resp.Body.Close()
fmt.Println("状态码:", resp.StatusCode, "已应用UA伪装+代理+限速")
}
⚠️ 新手必踩的坑: 代理 IP 质量参差不齐,一定要做"探活 + 失败剔除 + 自动补充",否则一堆死代理只会让抓取超时率飙升。限速器必须按域名区分——所有站点共用一个 limiter 会让低频小站也被限死,而热门站却没限住。另外字体反爬最隐蔽:页面上肉眼看到的是数字,HTML 源码里是乱码字符,必须下载字体文件做映射还原,否则存进库的全是错的。
十、电商图片过多导致带宽过高
10.1 用生活类比先建立直觉
图片带宽问题就像"仓库每天发几百万个包裹": 一家电商每天有千万用户刷商品页,每个页面塞了二三十张高清大图。如果每张原图都 3MB 直接发给用户,等于每天发几百万个"巨型包裹",仓库出口带宽被挤爆(带宽费用飙升、用户加载慢、手机流量哭)。
于是仓库做了几件事来"减重提速":
- CDN 加速:在全国各地开分仓(CDN 边缘节点),用户就近从分仓取图,不用每次都回总仓(回源)。
- 图片压缩(WebP/AVIF):把包裹里的棉花抽掉,同样的东西体积更小(WebP 比 JPEG 小 25%~35%,AVIF 更小),肉眼几乎看不出差别。
- 懒加载:首屏只发用户能看到的图,往下划才发下面的(用户没滚到的图先不发,省带宽)。
- 缩略图:商品列表用"小包裹"(缩略图 200×200),点进去详情才发原图。
- 对象存储:图片不存应用服务器硬盘,存到专门的"大件寄存处"(对象存储 OSS/COS),应用服务器只回一个 URL。
- HTTP/2 多路复用:一条高速专线上同时发很多个小包裹(多路复用),不用像 HTTP/1.1 那样排长队(队头阻塞)。
- 防盗链:只允许自家域名引用图片,别人盗链直接拒发(防止别人白嫖你的带宽)。
graph TB
U[用户浏览器] -->|请求商品页| WEB[应用服务器]
WEB -->|返回HTML+图片URL| U
U -->|按需加载图片| CDN[CDN边缘节点
就近缓存]
CDN -->|命中| U
CDN -->|未命中回源| OSS[(对象存储
原图+缩略图 WebP/AVIF)]
OSS -->|返回压缩图| CDN
U -.->|懒加载:滚动到才请求| CDN
REF[盗链域名请求] -->|Referer校验失败| DENY[拒绝返回 403]桥接: 仓库发货类比映射到图片带宽优化——分仓对应 CDN 边缘缓存,抽棉花对应 WebP/AVIF 压缩,小包裹对应缩略图,滚动才发对应懒加载,大件寄存处对应对象存储,高速专线对应 HTTP/2,门禁对应防盗链。优化顺序建议:先压缩 + 缩略图(立省 60%~80% 体积)→ 再上 CDN(减少回源)→ 最后懒加载 + HTTP/2(削减首屏请求)。
10.2 工程要点
优化手段对照表
| 手段 | 原理 | 收益 | 注意点 |
|---|---|---|---|
| CDN | 边缘节点缓存,就近回源 | 回源率下降、延迟低 | 缓存刷新策略(改图要清缓存) |
| WebP/AVIF | 现代编码更小 | 体积降 25%~50% | 老浏览器兼容(<picture> 回退) |
| 懒加载 | 滚动到视口才加载 | 首屏请求大减 | 需占位避免布局抖动 |
| 缩略图 | 不同场景不同尺寸 | 列表页省流量 | 原图/缩略图分开存储 |
| 对象存储 | 图片与计算分离 | 应用服务器减负 | 直接 URL 访问,运维简单 |
| HTTP/2 | 单连接多路复用 | 并发请求不排队 | 需 TLS,连接复用 |
| 防盗链 | Referer / 签名校验 | 防带宽被白嫖 | 合法 CDN 域名要放行 |
关键设计点
- 压缩与格式协商:服务端根据
Accept头判断是否支持avif/webp,优先返回更小格式;原图上传后异步转码出多档缩略图(如 200/400/800 宽)。 - CDN 缓存刷新:图片更新时通过 API 主动刷新 CDN 缓存,避免用户看到旧图;用 URL 带版本号(如
?v=2)也能强制回源拿新图。 - 防盗链:在对象存储或 CDN 层配置 Referer 白名单,或给图片 URL 加时效签名(token + 过期时间),无签名的外站请求直接 403。
Go 实现防盗链(Referer 校验)中间件 + 图片格式协商:
package main
import (
"fmt"
"net/http"
"strings"
)
// 步骤1:防盗链中间件(Referer 白名单)
func antiLeech(allowedHosts []string, next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
referer := r.Header.Get("Referer")
// 步骤2:无 Referer 一般放行(直接访问),或按策略拒绝
if referer != "" {
ok := false
for _, h := range allowedHosts {
if strings.Contains(referer, h) {
ok = true
break
}
}
if !ok {
http.Error(w, "403 防盗链: 不允许外站引用", http.StatusForbidden)
return
}
}
next.ServeHTTP(w, r)
})
}
// 步骤3:根据 Accept 头协商返回最佳格式
func negotiateImage(w http.ResponseWriter, r *http.Request) {
accept := r.Header.Get("Accept")
switch {
case strings.Contains(accept, "image/avif"):
// 步骤4:优先返回 AVIF(体积最小)
w.Header().Set("Content-Type", "image/avif")
fmt.Fprint(w, "<AVIF 二进制数据>")
case strings.Contains(accept, "image/webp"):
w.Header().Set("Content-Type", "image/webp")
fmt.Fprint(w, "<WEBP 二进制数据>")
default:
// 步骤5:回退到 JPEG(兼容老浏览器)
w.Header().Set("Content-Type", "image/jpeg")
fmt.Fprint(w, "<JPEG 二进制数据>")
}
}
func main() {
mux := http.NewServeMux()
mux.HandleFunc("/img/", negotiateImage)
// 步骤6:套上防盗链,仅允许自家域名引用
handler := antiLeech([]string{"myshop.com", "cdn.myshop.com"}, mux)
fmt.Println("图片服务启动 :8080/img/xxx")
http.ListenAndServe(":8080", handler)
}
⚠️ 新手必踩的坑: 压缩不是越狠越好——商品图压过头会糊,影响转化。建议用"质量参数 + 人眼抽检"定档。另外 CDN 缓存刷新经常被忘:改了图但用户还看旧图,多半是 CDN 没刷或 URL 没带版本号。防盗链配置要记得把自家 CDN 域名加进白名单,否则"自己盗自己"也会被 403。
十一、秒杀系统设计
11.1 用生活类比先建立直觉
类比:秒杀就像超市"限时免费领鸡蛋"——100 个名额,却涌来 10 万人。如果让 10 万人直接冲进仓库抢,仓库必被踩塌(DB 崩溃)。所以超市要设层层关卡:①门口先发号/答题筛掉机器人(前端+网关);②只放真正有资格的人进缓冲区(限流);③仓库门口挂一块"剩余 100 个"的黑板,每人进门前先看黑板、减 1,减到 0 就关门(Redis 预扣库存);④进门的人把"我要鸡蛋"的纸条投进箱子(消息队列),后台慢慢按纸条发货(异步下单);⑤发货时再正式从仓库台账扣(数据库最终扣减)。
核心思想就一句话:把"瞬时海量写"拆成"前端拦、网关限、Redis 快扣、队列缓、DB 慢落"五段,让数据库永远只处理它扛得住的量。
flowchart LR
U[海量用户请求] --> F[前端
答题/验证码/按钮置灰]
F --> GW[网关
鉴权+限流+熔断]
GW -->|"通过且名额内"| R[Redis 预扣库存
Lua原子 DECR]
R -->|"扣减成功"| Q[消息队列
削峰缓冲]
Q --> C[异步下单Worker
创建订单]
C --> DB[(数据库
最终扣减+落单)]
R -.->|"扣减失败
已售罄"| Rej[直接返回
"秒杀结束"]桥接:每一层都在"减少到达下一层的请求数"——前端和网关挡掉无效/超额流量,Redis 用单点原子扣减保证不超卖,队列把瞬时尖峰摊平成平稳的下游写入。面试时把这条链路讲清楚,比背任何参数都加分。
11.2 工程要点
前端:答题 / 验证码 / 按钮置灰
- 按钮置灰:点击后立即禁用,防止用户狂点产生重复请求(也防手抖刷接口)。
- 答题 / 验证码:让请求多一次人机交互,既挡掉脚本机器人,又人为拉长请求时间,把瞬时并发从"同一毫秒 10 万"摊成"几秒内 10 万"。
- 动静分离 + CDN:秒杀页面静态部分走 CDN,别让抢购请求顺带把页面 HTML 也打回源站。
网关:鉴权 + 限流
- 鉴权:校验登录态/券资格,没资格的请求在网关层直接 401 拒绝,不往下传。
- 限流:单机用令牌桶(如
golang.org/x/time/rate),集群用 Redis + Lua 做全局滑动窗口。超过阈值的请求直接返回"排队中/已售罄"。 - 熔断降级:下游 Redis/DB 异常时快速失败,别让请求堆积把网关拖垮。
库存:Redis 预扣 + 数据库最终扣减(防超卖核心)
超卖的根因是"多个请求同时读库存、都判断还有、都去扣"——竞态。解决之法是用 Redis 单线程 + Lua 脚本做"判断+扣减"的原子操作:只有 DECR 后的值 ≥ 0 才算抢到,杜绝并发超卖。
package main
import (
"context"
"fmt"
"time"
"github.com/redis/go-redis/v9"
)
// 步骤1:Lua 脚本——原子地"判断库存>0 并扣减 1",避免并发超卖
// KEYS[1] = 库存 key,ARGV[1] = 本次扣减数量(秒杀通常为 1)
var deductScript = `
local stock = tonumber(redis.call("GET", KEYS[1]))
if stock == nil or stock < tonumber(ARGV[1]) then
return -1 -- 库存不足,秒杀失败
end
return redis.call("DECRBY", KEYS[1], ARGV[1]) -- 原子扣减,返回剩余库存
`
// TrySecKill 预扣库存,返回 true 表示抢到名额
func TrySecKill(ctx context.Context, rdb *redis.Client, userID, itemID string) (bool, error) {
stockKey := "seckill:stock:" + itemID
// 步骤2:用 Lua 原子扣减,单线程串行执行,绝不会超卖
res, err := rdb.Eval(ctx, deductScript, []string{stockKey}, 1).Int()
if err != nil {
return false, err
}
if res < 0 {
return false, nil // 没抢到
}
fmt.Printf("用户 %s 抢到商品 %s,剩余库存 %d\n", userID, itemID, res)
return true, nil
}
func main() {
rdb := redis.NewClient(&redis.Options{Addr: "localhost:6379"})
ctx := context.Background()
// 步骤3:初始化库存(实际由运营后台在秒杀开始前写入)
rdb.Set(ctx, "seckill:stock:iphone15", 100, time.Hour)
// 步骤4:模拟 200 个并发抢 100 个库存,Lua 原子保证恰好 100 人成功
for i := 0; i < 200; i++ {
go func(id int) {
ok, _ := TrySecKill(ctx, rdb, fmt.Sprintf("u%d", id), "iphone15")
_ = ok
}(i)
}
time.Sleep(time.Second)
}
抢到名额后,请求并不直接写订单库(那会被瞬时打爆),而是把下单消息投进队列,由后端 Worker 慢慢消费——这就是"异步下单"。
队列削峰 + 异步下单
// 步骤5:抢到名额后,向队列投递下单消息(削峰,保护 DB)
func EnqueueOrder(ctx context.Context, rdb *redis.Client, userID, itemID string) error {
// 用 Stream 或 List 当队列;这里用 List 简化演示
msg := fmt.Sprintf("%s:%s:%d", userID, itemID, time.Now().UnixNano())
return rdb.LPush(ctx, "seckill:orders", msg).Err()
}
// 步骤6:异步 Worker——从队列取消息,创建订单并"最终扣减"数据库库存
func OrderWorker(ctx context.Context, rdb *redis.Client, db *sql.DB) {
for {
// BRPOP 阻塞取,没有消息时 goroutine 休眠不占 CPU
res, err := rdb.BRPop(ctx, 0, "seckill:orders").Result()
if err != nil {
return
}
orderMsg := res[1] // "userID:itemID:ts"
// 步骤7:真正落库——创建订单 + 数据库库存 UPDATE ... WHERE stock>0
// 数据库层面用行锁/乐观锁兜底:UPDATE inventory SET stock=stock-1 WHERE item=? AND stock>0
_ = createOrderInDB(db, orderMsg)
}
}
⚠️ 新手必踩的坑: Redis 预扣成功 ≠ 订单一定成功。用户抢到名额后可能不付款,所以库存要设计"回补"机制:订单超时未支付(如 15 分钟)就
INCR把库存加回去。否则会出现"Redis 显示售罄、实际一堆未支付订单占用名额"的假售罄。
保证公平性
| 手段 | 做法 | 目的 |
|---|---|---|
| 随机打散请求 | 答题/验证码人为拉长请求时间 | 避免同一毫秒洪峰 |
| 排队有序 | 队列 FIFO,谁先抢到名额谁先下单 | 先到先得 |
| 防刷 | 同一用户限领 1 份(Redis SET 去重 userID) | 防止一人囤货 |
| 库存回补 | 超时未支付归还库存 | 名额不被僵尸订单占用 |
完整时序图:
sequenceDiagram
participant U as 用户(前端)
participant G as 网关
participant R as Redis(库存)
participant Q as 消息队列
participant W as 下单Worker
participant D as 数据库
U->>G: 提交秒杀(答题/验证码)
G->>G: 鉴权+限流
G->>R: Lua原子预扣库存
alt 库存>0 扣减成功
R-->>G: 剩余库存
G->>Q: 投递下单消息
Q->>W: 异步消费
W->>D: 创建订单+最终扣减
D-->>W: 成功
W-->>U: 秒杀成功(轮询/推送)
else 库存不足
R-->>G: -1
G-->>U: 已售罄
end一句话总结:秒杀的命门是"数据库扛不住瞬时写",所有设计都围绕"在请求到达数据库之前,用前端、网关、Redis、队列四道关卡把写入量压到数据库能承受的程度",而 Redis 原子预扣 + 异步落库是防超卖、保稳定的双保险。
十二、多活容灾与备份恢复(云原生进阶)
说明:第十一节讲了单集群秒杀等业务的"抗峰"设计,但架构师还要回答一个更硬的问题——“机房炸了怎么办”。本章覆盖单元化/同城双活/异地多活架构、数据库容灾(MGR / Redis 异地)、备份恢复(PITR / Velero)、故障演练与流量调度切换。这是云原生容灾面试的集大成章节。
12.1 跨 Region Kubernetes 集群灾备
用生活类比先建立直觉
跨 Region 灾备像"两座城市的双中心电网":A 城是主供电(主集群),B 城是备用供电(灾备集群)。平时 B 城冷备或轻载;A 城整个变电站炸了(API Server 不可达),调度中心立刻把负荷切到 B 城,用户几乎无感。但切换的前提是——B 城手里也握着最新的"用户档案副本"(数据备份),否则切过去也是空城。
graph TB
subgraph 主集群[Region A 主]
A1[API Server] --> A2[(etcd + PV 数据)]
end
subgraph 灾备集群[Region B 备]
B1[API Server] --> B2[(备份恢复的数据)]
end
DNS[全局 DNS / GSLB] -->|主集群存活| A1
DNS -.->|主集群 API 不可达 自动切| B1
BK[Velero 增量备份] -->|持续同步 PV| B2桥接: “主供电/备用供电"对应主备双集群;“调度中心切负荷"对应 DNS/GSLB 在探测到主集群 API Server 不可达时自动把流量指到灾备;“用户档案副本"对应 Velero 对 PV 的增量备份与恢复。
工程要点
1. Velero + Restic 对 >1TB PV 增量备份并验证完整性
PV 超过 1TB 时,全量备份既慢又占存储。Velero 配合 Restic(或 CSI 快照)只备份"上次备份之后变化的数据块”。备份完不能"传完就完事”,要用 velero backup describe --details 和校验和(checksum)确认每块都完整到达后端存储。
# 步骤1:创建增量备份(Restic 默认就是增量,只传变化块)
velero backup create prod-pv-bak \
--include-namespaces prod \
--default-volumes-to-restic \
--snapshot-volumes
# 步骤2:查看备份详情,确认每个卷的状态都是 Completed
velero backup describe prod-pv-bak --details
# 步骤3:恢复前先校验完整性(校验和比对,防止静默损坏)
velero backup logs prod-pv-bak | grep -i "checksum\|completed"
2. Crossplane + RDS Global Database 跨云 MySQL 灾备
Crossplane 把"云资源"变成 K8s 里的 CR,用同一套 GitOps 流程管跨云数据库。RDS Global Database 让主区域的写入自动异步复制到异地只读节点,跨云(如 AWS + 阿里云)灾备时,用 Crossplane 在两地各声明一个 RDS 实例并组成 Global Database。
# 步骤1:用 Crossplane Composition 声明跨云 MySQL Global Database
apiVersion: database.aws.crossplane.io/v1alpha1
kind: GlobalCluster
metadata:
name: mysql-global
spec:
forProvider:
region: ap-east-1
# 步骤2:主区域写入自动异步复制到灾备区域(通常 <1s 延迟)
engine: aurora-mysql
writeConnectionSecretToRef:
name: mysql-global-secret
namespace: prod
3. 主集群 API Server 不可达触发 DNS 自动切换
最可靠的"集群挂了"判据是"API Server 连不上”。用一个外部健康探测(如独立探针或 GSLB 的健康检查)持续打 API Server 的 /healthz,连续失败就通过 DNS(或 Anycast)把域名解析从主集群 VIP 切到灾备集群 VIP。
# 步骤1:外部探针周期性探测主集群 API Server 健康端点
curl -k -s https://<主集群API>:6443/healthz || echo "UNHEALTHY"
# 步骤2:连续 N 次不健康则调用 DNS/GSLB 切换(示意)
if [ "$(probe_fail_count)" -ge 3 ]; then
# 把 example.com 的解析从主集群 VIP 切到灾备集群 VIP
update_dns_record example.com --target <灾备集群VIP>
echo "主集群 API 不可达,已触发 DNS 切换至灾备"
fi
⚠️ 考点总结: 跨 Region 灾备三要素——“数据有副本(Velero/Crossplane)、切换有判据(API Server 健康)、切换有手段(DNS/GSLB)"。增量备份别忘了校验完整性,否则恢复时才发现备份是坏的。
12.2 有状态服务数据一致性
用生活类比先建立直觉
有状态服务的数据一致性像"两地账簿对账”:总店(主库)和分店(从库)各记一本账。最怕的是"网络抖动让两地都以为自己是总店",各自收钱各自记,最后对账发现两边差了十万八千里(双主写入冲突)。解决是"谁拿到的’令牌’(选主)才有权记账",网络恢复后落伍的那本账自动认输对齐。
graph TB
subgraph 主
P[PostgreSQL 主] -->|Patroni 基于 etcd 选主| E[(etcd)]
end
subgraph 备
S[PostgreSQL 备] -->|同看 etcd 令牌| E
end
P -->|WAL 流复制| S
NET[网络分区] -.->|误判导致双主风险| E桥接: “拿令牌才有权记账"对应 Patroni + etcd 的分布式选主——只有被 etcd 选为 Leader 的节点才接受写;“落伍的账本对齐"对应备库通过 WAL 流复制追平主库。
工程要点
1. Patroni + etcd 防网络分区双主写
Patroni 用 etcd(或 Consul/ZooKeeper)做选主仲裁:启动时谁先在 etcd 抢到 Leader 锁,谁就是主。网络分区时,少数派节点因拿不到锁自动降级为只读,杜绝双主。关键是 etcd 本身要奇数节点、跨故障域部署,避免仲裁本身脑裂。
# 步骤1:Patroni 配置——用 etcd 做选主后端
scope: postgres-cluster
name: pg-node-1
restapi:
listen: 0.0.0.0:8008
etcd:
hosts:
- etcd-1:2379
- etcd-2:2379
- etcd-3:2379
bootstrap:
dcs:
# 步骤2:定义选主与故障切换参数,防止双主
ttl: 30
loop_wait: 10
retry_timeout: 10
maximum_lag_on_failover: 1048576 # 步骤3:备库落后超过 1MB 不参与接管
2. Kafka MirrorMaker 2.0 双向复制解决循环复制
两个 Kafka 集群互相镜像(A→B、B→A)做异地容灾,但朴素双向复制会把"从 A 复制到 B 的消息"再复制回 A,无限循环。MirrorMaker 2.0 给每条消息打上"来源集群"标记(topic 前缀 A. / B.),复制时跳过"自己产生的消息”,从而打破循环。
# 步骤1:MirrorMaker 2.0 双向复制配置(截断循环的关键在 replication.policy)
clusters:
- name: cluster-a
bootstrapServers: kafka-a:9092
- name: cluster-b
bootstrapServers: kafka-b:9092
mirrors:
- sourceCluster: cluster-a
targetCluster: cluster-b
# 步骤2:IdentityReplicationPolicy 让消息保留来源标记,避免被反向再复制
replicationPolicy: org.apache.kafka.connect.mirror.IdentityReplicationPolicy
- sourceCluster: cluster-b
targetCluster: cluster-a
replicationPolicy: org.apache.kafka.connect.mirror.IdentityReplicationPolicy
3. etcd 快照恢复后校验 Patroni 集群视图与实际 PG 一致
从 etcd 快照恢复后,Patroni 的"集群视图”(谁主谁备)可能和实际 PostgreSQL 的 pg_is_in_recovery() 状态对不上——比如 etcd 认为 A 是主,但 A 实际已从库。必须校验:etcd 里的 Leader 锁指向的节点,其 pg_is_in_recovery() 必须为 false(真主)。
# 步骤1:看 etcd 里 Patroni 记录的 Leader 是谁
etcdctl get /service/postgres-cluster/leader
# 步骤2:到该节点执行,必须为 f(false=主);若是 t 说明视图与实际不一致
psql -c "SELECT pg_is_in_recovery();" # 期望返回 f
# 步骤3:不一致时让 Patroni 重新选主对齐视图
patronictl reinit postgres-cluster <视图错误的节点>
⚠️ 考点总结: 有状态一致性的命门是"选主仲裁 + 防双主"。etcd/Patroni 负责谁有权写,WAL/流复制负责追平数据,MirrorMaker 2.0 的 replication policy 负责打破双向复制循环。恢复后务必校验"控制面视图"与"数据面实际角色"一致。
12.3 蓝绿发布与快速回切
用生活类比先建立直觉
蓝绿发布像"商场换招牌":蓝组(旧版本)正常营业,绿组(新版本)在后台偷偷布置好。等到切换时刻,把门口指引牌从"去蓝区"改成"去绿区"——顾客(用户)顺着指引瞬间全部去绿区,蓝区空着待命。万一绿区出问题,把指引牌改回蓝区,顾客立刻回到旧版本,绿区当晚再修。
sequenceDiagram
participant U as 用户流量
participant B as 蓝组(旧 v1)
participant G as 绿组(新 v2)
U->>B: 切换前 100% 流量
Note over G: 绿组已部署就绪 但不接流量
U->>G: 切换后 流量整体切到绿组
Note over U,G: 会话保持 同一用户不跳来跳去
alt 绿组异常
U->>B: 一键回切 流量回蓝组 绿组保留7天
end桥接: “换指引牌"对应流量整体切换(权重 0/100 或全切);“绿区后台布置"对应新旧两套环境并存;“指引牌改回"对应一键回切;“绿组保留7天"对应旧版本不立即销毁、留作快速回退。
工程要点
1. 蓝绿基于权重 + 会话保持用户无感知切换
切换不是"硬切”,而是用网关/Service Mesh 把权重从蓝 100% 调到绿 100%,同时按用户会话 ID 做会话保持(sticky session),保证同一个用户在一次会话内不被蓝绿之间来回弹,体验无感。
# 步骤1:用 Istio VirtualService 做蓝绿权重切换 + 会话保持
apiVersion: networking.istio.io/v1beta1
kind: VirtualService
metadata:
name: checkout
spec:
hosts: ["checkout.example.com"]
http:
- route:
- destination: { host: checkout-blue }
weight: 0 # 步骤2:切换前绿 0 / 蓝 100
- destination: { host: checkout-green }
weight: 100
# 步骤3:基于 cookie 的会话保持,同一用户始终落同一版本
headers:
request:
set:
x-canary: "green"
2. Argo Rollouts 一键蓝绿回滚并保留旧版本 7 天
Argo Rollouts 的 BlueGreen 策略原生支持"切换 + 自动分析 + 回滚”。切换后把旧 ReplicaSet(蓝)保留 7 天(scaleDownDelay / 保留 revision 数),期间随时可一键回切,7 天后才真正回收资源。
apiVersion: argoproj.io/v1alpha1
kind: Rollout
metadata:
name: checkout
spec:
strategy:
blueGreen:
# 步骤1:切换后保留旧版本(蓝)7 天再缩容,便于快速回切
scaleDownDelaySeconds: 604800
activeService: checkout-active
previewService: checkout-preview
revisionHistoryLimit: 7 # 步骤2:保留 7 个历史 revision,支持一键回滚
3. 数据库慢查询致绿环境异常自动阻断流量
绿环境接的是同一份生产库,如果新版本 SQL 写得差引发慢查询拖垮数据库,会连累蓝组也跟着慢。防护:在流量切换的同时挂一个"数据库健康探针”,慢查询数/连接数超阈值就自动把流量切回蓝组,阻断绿组继续放大故障。
# 步骤1:监控绿环境引发的数据库慢查询数
SLOW=$(mysql -e "SHOW GLOBAL STATUS LIKE 'Slow_queries'" | awk '{print $2}')
# 步骤2:超过阈值则阻断绿组流量,自动回切蓝组
if [ "$SLOW" -gt 100 ]; then
echo "绿环境引发慢查询风暴,自动阻断并回切蓝组"
kubectl patch virtualservice checkout -p '{"spec":{"http":[{"route":[{"destination":{"host":"checkout-blue"},"weight":100},{"destination":{"host":"checkout-green"},"weight":0}]}]}}'
fi
⚠️ 考点总结: 蓝绿发布的优势是"几乎零停机、切换干净、可一键回切”。但必须配套"会话保持"保证体验、配合"旧版本保留期"留退路、并盯住"共享数据库"这个共同风险点——绿组出 SQL 问题会连累蓝组。
12.4 备份数据加密与密钥管理
用生活类比先建立直觉
备份数据加密像"把重要文件锁进保险柜再放进仓库":文件(备份)本身放进了带锁的仓库(S3),但这还不够——你还要把钥匙(加密密钥)单独交给一个更可信的"钥匙保管员"(KMS/Vault),而且这把钥匙定期更换(轮换)。即使仓库被盗,没有钥匙打不开;钥匙泄露了,换一把就能让旧钥匙失效。
flowchart LR
B[备份数据] -->|SSE-C 客户端密钥加密| S3[(S3 存储)]
B -->|KMS 托管密钥加密| S3
K[(KMS 密钥)] -->|信封加密| V[Vault Transit]
V -->|加密 etcd 快照| ES[(etcd 快照)]
RK[密钥轮换] -->|旧备份批量重加密| S3桥接: “带锁仓库"对应 S3 + 服务端加密;“钥匙保管员"对应 KMS;“信封加密"对应用 KMS 的数据密钥加密数据、再用 Vault 保护数据密钥,两层保险;“换钥匙"对应密钥轮换后批量重加密历史备份。
工程要点
1. Velero 到 S3 启用 KMS + SSE-C 双重加密
Velero 备份到 S3 时,先由 Velero 用自己管理的密钥做客户端加密(SSE-C),到 S3 端再由 S3 用 KMS 托管密钥再做一层服务端加密。双保险:即便某层密钥泄露,还有另一层。
# 步骤1:Velero 备份存储位置配置——KMS + SSE-C 双重加密
apiVersion: velero.io/v1
kind: BackupStorageLocation
metadata:
name: s3-backup
spec:
provider: aws
objectStorage:
bucket: my-velero-bucket
config:
# 步骤2:服务端用 KMS 托管密钥加密(SSE-KMS)
sse-kms-key-id: arn:aws:kms:ap-east-1:123:key/abcd
# 步骤3:客户端用指定密钥加密(SSE-C),数据在出 Velero 前已加密
sse-customer-algorithm: AES256
sse-customer-key: <base64-encoded-key>
2. Vault Transit 引擎对 etcd 快照信封加密 Shell 脚本
etcd 快照是集群的"全部家当”,落地前用 Vault Transit 引擎加密:Vault 用一一个"密钥加密密钥(KEK)“把数据加密密钥(DEK)包起来(信封加密),DEK 随数据走、KEK 只驻留在 Vault。这样轮换 KEK 时不用重新加密全部数据。
#!/bin/bash
# 步骤1:先 snapshot 出 etcd 数据
ETCDCTL_API=3 etcdctl --endpoints=:2379 snapshot save /tmp/etcd-snap.db
# 步骤2:用 Vault Transit 加密快照(信封加密:Vault 管理 KEK)
vault write transit/encrypt/etcd-key \
plaintext=$(base64 /tmp/etcd-snap.db) \
> /tmp/etcd-snap.enc.json
# 步骤3:把密文落到备份存储,明文快照立即销毁
rm -f /tmp/etcd-snap.db
echo "etcd 快照已信封加密并安全落盘"
3. KMS 密钥轮换批量重加密历史备份不重传
KMS 密钥轮换后,旧备份仍用旧密钥加密。若想统一到新密钥,不需要把数据重新上传一遍——用"重加密"接口就地换新密钥的密文(KMS 在服务端用旧密钥解、新密钥加密新的数据密钥),省带宽、不重写数据体。
# 步骤1:列出需要轮换密钥的历史备份
aws s3 ls s3://my-velero-bucket/backups/ > /tmp/old_backups.txt
# 步骤2:对每个对象调用 S3 重加密(KMS 服务端就地换密钥,无需重传)
while read obj; do
aws s3api copy-object \
--bucket my-velero-bucket \
--key "$obj" \
--copy-source "my-velero-bucket/$obj" \
--sse-kms-key-id arn:aws:kms:ap-east-1:123:key/NEW-KEY \
--metadata-directive REPLACE
done < /tmp/old_backups.txt
⚠️ 考点总结: 备份加密的核心是"信封加密 + 密钥与数据分离 + 可轮换”。Vault Transit 管 KEK、S3/KMS 管 DEK,轮换 KMS 用
copy-object就地重加密而不是重新上传。永远不要让加密密钥和数据躺在同一个桶里。
12.5 合规与审计日志留存
用生活类比先建立直觉
合规审计日志像"银行的监控录像 + 盖章台账”:录像(审计日志)必须"只写一次、谁都不能改”(WORM,一次写多次读),否则出事后有人偷偷删了录像就死无对证。而且监管(等保 2.0)要求这类记录至少保存 7 年,所以老录像要自动从"常温硬盘"挪到"冷库”(Glacier)长期封存,同时给录像做"指纹链"(Merkle Tree),谁改了一帧立刻被发现。
flowchart LR
E[集群事件] -->|Falco 捕获| F[Fluent Bit]
F -->|写入 WORM 存储| W[(不可篡改存储)]
W -->|生命周期策略| G[(Glacier 7年归档)]
L[审计日志] -->|Merkle Tree 指纹| M[校验服务]
M -->|检测篡改| A[告警]桥接: “监控录像"对应审计日志;“不可改的台账"对应 WORM 存储;“冷库封存"对应 S3 Glacier 7 年生命周期;“指纹链"对应 Merkle Tree——任何一条记录被改,根哈希对不上,立即告警。
工程要点
1. Falco + Fluent Bit 审计日志写 WORM 满足等保 2.0
Falco 在节点侧捕获异常系统调用与 Kubernetes 审计事件,Fluent Bit 把它转发到支持 WORM(一次写多次读、不可删除/改写)的对象存储。等保 2.0 要求审计记录"不可篡改、可溯源”,WORM 正是合规落点。
# 步骤1:Fluent Bit 输出到 WORM 存储(以 S3 Object Lock 为例)
apiVersion: v1
kind: ConfigMap
metadata:
name: fluent-bit-config
data:
fluent-bit.conf: |
[OUTPUT]
Name s3
Match falco.*
bucket audit-log-worm
# 步骤2:开启 Object Lock,写入后合规期内不可删改
s3_key_format /logs/$TAG/%Y/%m/%d/%H/%M/%S
# 步骤3:配合存储桶的 Object Lock 合规模式,满足等保 2.0 不可篡改要求
2. Loki + S3 Glacier 生命周期 7 年归档
热日志存 Loki + S3 标准层供实时查询,老日志通过 S3 生命周期策略自动沉降到 Glacier 深度归档,保留 7 年,既省钱又满足长期合规留存。
# 步骤1:S3 生命周期规则——热转冷、保留 7 年
apiVersion: s3.aws.crossplane.io/v1alpha1
kind: Bucket
metadata:
name: audit-log-archive
spec:
forProvider:
lifecycleConfiguration:
rules:
- id: to-glacier-7y
status: Enabled
# 步骤2:30 天后从标准层转到 Glacier 深度归档
transitions:
- days: 30
storageClass: GLACIER
# 步骤3:保留 7 年(2555 天)后自动过期清理
expiration:
days: 2555
3. 审计日志篡改用 Merkle Tree 校验并告警
给审计日志按时间分块构建 Merkle Tree,把根哈希定期(如每小时)写进不可篡改的存储或区块链式账本。事后校验时,只要任意一条记录被篡改,根哈希就会变化,比对失败即触发告警,定位到具体被改的块。
#!/bin/bash
# 步骤1:对当日审计日志分块计算 Merkle 根哈希
ROOT=$(build_merkle_root /var/log/audit/$(date +%F)/*.log)
# 步骤2:把根哈希写入不可篡改账本(WORM / 区块链)
write_root_to_ledger "$ROOT" "$(date +%F-%H)"
# 步骤3:校验时重新计算并比对,不一致即告警
CURRENT=$(build_merkle_root /var/log/audit/$(date +%F)/*.log)
if [ "$CURRENT" != "$ROOT" ]; then
echo "ALERT: 审计日志检测到篡改,根哈希不匹配" >&2
trigger_alert "audit-log-tampered"
fi
⚠️ 考点总结: 合规审计三板斧——“采集用 Falco、存储用 WORM 不可篡改、留存用 Glacier 生命周期 7 年”。Merkle Tree 是篡改检测的"指纹锁”:根哈希一对不上,立刻知道有人动过日志。等保 2.0 的核心就是"留得住、改不了、查得到”。
十三、自测题与动手练习
自测题
直播弹幕系统中,为什么用消息队列削峰? 如果不用消息队列,直接让 WebSocket 节点处理弹幕,在高并发场景下会出现什么问题?
朋友圈 Feed 流的推模式和拉模式各自的优缺点是什么? 什么场景下适合推拉结合?大 V 的阈值如何确定?
短链重定向用 301 还是 302?为什么? 如果用 301 会出现什么问题?发号器生成的 ID 为什么不能是连续的?
缓存雪崩和缓存击穿的区别是什么? 各自的解决方案有何不同?互斥锁重建为什么要做"双重检查”?
扫码登录中,PC 端轮询和 WebSocket 通知各有什么优缺点? 二维码的临时 token 为什么必须设置过期时间?
设计微信海量数据存储系统,第一步应该做什么? 为什么好友关系链必须按 userID 分片而不是按群组 ID 分片?一致性哈希相比取模分片在扩容时有什么优势?冷热分离能解决什么问题?
秒杀系统为什么不能让用户请求直接写订单数据库? 请说出"前端 → 网关 → Redis 预扣 → 队列 → 异步落库"这条链路上,每一层分别帮数据库挡掉了什么压力。
用 Redis 做秒杀库存预扣时,为什么必须用 Lua 脚本而不是"先 GET 判断再 DECR"?只做 Redis 预扣够不够,还需要什么机制兜住"抢到但不付款"的情况?
动手练习
实现多房间弹幕系统: 用 Go + gorilla/websocket 实现一个支持多房间的弹幕广播服务。要求:支持创建/加入/离开房间,消息按房间隔离广播,每秒弹幕量超过 100 条时自动限流(令牌桶算法),用 Redis Pub/Sub 支持跨节点广播。
实现长短链转换服务: 用 Go 实现完整的长短链转换服务。要求:用 Redis INCR 实现发号器,Base62 编码生成 6 位短码,支持长短链双向映射,302 重定向时记录访问统计(UV/PV),为同一长链返回相同短链(去重)。
实现扫码登录 Demo: 用 Go + Redis 实现完整的扫码登录流程。要求:PC 端生成二维码(含临时 token),手机端扫码并确认,PC 端用轮询方式获取登录状态,token 60 秒过期,包含完整的状态机(等待扫码到已扫码到已确认/已过期/已取消)。
实现分布式爬虫去重模块: 用 Go 实现基于 Redis Set 的精确去重器,封装
Add(url)/Exists(url)两个方法,支持设置已抓 URL 的过期时间(避免无限膨胀)。对比 Bloom Filter 方案,说明在 100 亿 URL 量级下两者内存占用的差异,并给出"Bloom Filter 前置过滤 + Redis Set 兜底"的混合实现思路。实现带反爬策略的下载器: 用 Go 实现支持 UA 随机伪装、代理 IP 池轮询、按域名令牌桶限速(每秒 2 次)的 HTTP 下载器。要求代理池后台定时探活、剔除失效代理;遇到 429 状态码时指数退避重试 3 次。说明命中字体反爬时如何定位并还原真实字符。
实现图片防盗链 + 格式协商服务: 用 Go 实现一个图片服务:① 中间件校验
Referer白名单,非白名单域名引用返回 403;② 根据请求头Accept协商返回 AVIF / WebP / JPEG;③ 给图片 URL 加时效签名(token + 过期时间)以支持无 Referer 的安全外链。画出"原图上传→转码缩略图→CDN 刷新"的整体流程。
十四、本章小结
本章围绕四个高频业务场景,系统讲解了从需求分析到架构设计的完整思路:
直播弹幕系统:WebSocket 长连接实现双向通信,消息队列(Kafka)削峰防止后端过载,Redis Pub/Sub 实现跨节点广播,令牌桶限流控制弹幕展示速率。核心是"写少读多"的读扩散模型——一条弹幕写入,N 个用户读取。
朋友圈 Feed 流:推模式(写扩散)适合粉丝少的用户,拉模式(读扩散)适合大 V,推拉结合兼顾两者。用 Redis ZSet 存储时间线,score 为时间戳,支持按时间排序和范围查询。inbox 须设上限和过期时间防止内存膨胀。
数据统计页面加速:预计算定时聚合结果写入缓存是核心手段,配合 CDN 静态化、前端异步加载、数据库物化视图等多层加速。缓存必须设过期时间,防止定时任务故障后返回永久过期数据。
缓存雪崩防护:热 key 过期导致大量请求穿透到 DB 是核心问题。互斥锁重建(双重检查)、热点 key 永不过期(逻辑过期 + 后台更新)、多级缓存(本地缓存 + Redis + DB)、过期时间加随机值是四大防护手段。
长短链转换:发号器生成全局唯一 ID,Base62 编码压缩为 6 位短码,Redis 存储双向映射。重定向用 302 而非 301,确保每次访问都经过服务器统计。ID 不能连续,防止竞品遍历爬取。
- 弹幕系统核心:WebSocket长连接保实时性 + Kafka削峰防雪崩 + Redis Pub/Sub跨节点广播——三件套缺一不可。
- Feed流推拉结合:小V用推模式(写扩散),大V用拉模式(读扩散),混合策略平衡写入与读取开销。
- 缓存过期策略:定时任务+TTL双保险,过期时间加随机值防雪崩,重要数据设逻辑过期后台异步刷新。
- 长短链安全:短码用雪花算法+Base62编码保证全球唯一,302重定向保留统计能力,防盗链配合Referer白名单。
发布时把内容推给所有粉丝的时间线。
问题:大V(如千万粉丝)发帖时,要写入库的数量巨大,写入压力大。
纯拉模式(读扩散):
关注列表存粉丝关注的人,取 feed 时逐个查每个被关注者的内容再合并排序。
问题:粉丝多时每次刷新要查很多库,读取压力大。
推拉结合:
① 小V用推模式:粉丝少,写入成本低
② 大V用拉模式:粉丝多,避免写入爆炸
③ 以粉丝数阈值(如1000)划分,平衡读写压力
面试加分点:提到Instagram、微博等真实产品都用推拉结合方案,且会根据用户活跃度动态调整策略。
- 扫码登录:PC 端生成含临时 token 的二维码,手机端扫码并确认,PC 端轮询或 WebSocket 获取登录状态。核心是通过服务器中转关联 PC 端和手机端两个会话,token 设置 60 秒过期保证安全。