vault backup: 2026-05-18 00:17:59
This commit is contained in:
+503
-110
@@ -1,186 +1,574 @@
|
||||
---
|
||||
tags: [microservice, database, sharding, replication]
|
||||
create time: 2026-05-05
|
||||
create time: 2026-05-17 14:30
|
||||
---
|
||||
|
||||
# 数据库拆分
|
||||
|
||||
## 概述
|
||||
|
||||
微服务的核心设计原则是 **"每个服务拥有独立数据库"**,这意味着每个服务的表结构、数据存储、甚至数据库类型都可以不同。但当单表数据量持续增长时,就面临拆分的需求。
|
||||
想象一下:你的电商系统刚上线时只有一张 `orders` 表,数据量不大,一条 SQL 就搞定。一年后,日订单量涨到 100 万,同样的查询慢到了 5 秒——数据库成了瓶颈,你不得不开始考虑:**怎么拆?**
|
||||
|
||||
微服务的核心设计原则是 **"每个服务拥有独立数据库"**。这意味着每个服务的表结构、数据存储、甚至数据库类型都可以不同。但当单表数据量持续增长时,你就面临另一个维度的拆分需求。
|
||||
|
||||
```mermaid
|
||||
graph LR
|
||||
S1[订单服务 DB]
|
||||
S2[库存服务 DB]
|
||||
|
||||
subgraph BAD["反模式:共享数据库"]
|
||||
S1 --- SHARED[(共享 DB)]
|
||||
S2 --- SHARED
|
||||
subgraph "理想状态:微服务 = 独立数据库"
|
||||
S1["📦 订单服务<br/>MySQL"]
|
||||
S2["👤 用户服务<br/>PostgreSQL"]
|
||||
S3["📦 库存服务<br/>Redis + MySQL"]
|
||||
end
|
||||
|
||||
|
||||
style S1 fill:#e8f5e9,stroke:#4caf50
|
||||
style S2 fill:#e8f5e9,stroke:#4caf50
|
||||
style S3 fill:#e8f5e9,stroke:#4caf50
|
||||
```
|
||||
|
||||
> [!failure] 反面教材:共享数据库
|
||||
>
|
||||
> 如果两个微服务连接同一个数据库的同一张表,它们就不再是独立的微服务——你得到的是 **分布式单体**。服务可以随意互相查对方的数据,失去了边界和自治性。这在早期开发中很常见(方便联调),但一定要在正式微服务化之前拆掉。
|
||||
|
||||
```mermaid
|
||||
graph LR
|
||||
ORDER["📦 订单服务"] --- SHARED[(❌ 共享 DB)]
|
||||
USER["👤 用户服务"] --- SHARED
|
||||
|
||||
style SHARED fill:#ffebee,stroke:#ef5350
|
||||
style ORDER fill:#fff3e0,stroke:#ff9800
|
||||
style USER fill:#fff3e0,stroke:#ff9800
|
||||
```
|
||||
|
||||
> [!failure] 反模式警告
|
||||
> 如果两个服务连接同一个数据库的同一张表,它们就不再是独立的微服务——你得到的是**分布式单体**。服务可以随意互相查询彼此的数据,失去了边界和自治性。
|
||||
## 两种拆分思路:先"纵向切",再"横向切"
|
||||
|
||||
## 垂直拆分 vs 水平拆分
|
||||
拆分不是选择题,而是**两步走**:第一步按业务域垂直拆分(微服务化的标配),第二步当单个表太大时再水平拆分(分库分表)。
|
||||
|
||||
### 垂直拆分(按业务域)
|
||||
### 第一步:垂直拆分 — 按业务域独立建库
|
||||
|
||||
按 **微服务边界** 拆库——这是微服务的标配。
|
||||
> [!question] 思考一下
|
||||
>
|
||||
> 如果一个"用户下单"操作需要同时访问订单表和商品信息表——这说明这两个表应该属于同一个数据库吗?
|
||||
>
|
||||
> **答案是否定的。** 商品不会因为你改了订单就跟着变。真正的判断标准是:**哪些表经常一起被修改?** 如果一组表的变更总是由同一个业务逻辑触发,它们就属于同一个限界上下文,应该放在同一个库里。
|
||||
|
||||
核心做法:以 **微服务边界** 为基准,将相关表打包到一个数据库中。
|
||||
|
||||
```mermaid
|
||||
graph TB
|
||||
subgraph "DB-Order"
|
||||
orders[orders]
|
||||
order_items[order_items]
|
||||
order_status[order_status]
|
||||
flowchart TB
|
||||
subgraph DBOrder["DB-Order:订单域"]
|
||||
direction LR
|
||||
T1[orders]
|
||||
T2[order_items]
|
||||
T3[order_status_log]
|
||||
end
|
||||
|
||||
subgraph "DB-User"
|
||||
users[users]
|
||||
user_profiles[user_profiles]
|
||||
user_addresses[user_addresses]
|
||||
|
||||
subgraph DBUser["DB-User:用户域"]
|
||||
direction LR
|
||||
U1[users]
|
||||
U2[user_profiles]
|
||||
U3[user_addresses]
|
||||
end
|
||||
|
||||
subgraph "DB-Product"
|
||||
products[products]
|
||||
categories[categories]
|
||||
product_images[product_images]
|
||||
|
||||
subgraph DBProduct["DB-Product:商品域"]
|
||||
direction LR
|
||||
P1[products]
|
||||
P2[categories]
|
||||
P3[product_images]
|
||||
end
|
||||
|
||||
style DBOrder fill:#e3f2fd,stroke:#1976d2,rx:8
|
||||
style DBUser fill:#e3f2fd,stroke:#1976d2,rx:8
|
||||
style DBProduct fill:#e3f2fd,stroke:#1976d2,rx:8
|
||||
```
|
||||
|
||||
| 特点 | 说明 |
|
||||
| 优势 | 说明 |
|
||||
|------|------|
|
||||
| 每个服务独占一个数据库 | 物理隔离,互不干扰 |
|
||||
| 可异构选型 | 订单用 MySQL,用户用 PostgreSQL,搜索用 Elasticsearch |
|
||||
| 天然解耦 | 服务间不能直接查对方库 |
|
||||
| **物理隔离,互不干扰** | 订单库崩了不影响用户登录 |
|
||||
| **异构选型** | 订单用 MySQL(事务强),搜索用 Elasticsearch,缓存用 Redis |
|
||||
| **天然解耦** | 服务之间不能直连对方数据库,只能通过 API 通信 |
|
||||
|
||||
### 水平拆分(分库分表)
|
||||
> [!tip] 拆分的粒度
|
||||
>
|
||||
> 不要把一张表里的字段都拆到不同库里——那叫"过度拆分"。**以表为单位**进行垂直拆分是最常见的做法。一个微服务对应一个库,一个库包含多个相关表。
|
||||
|
||||
当单个表的记录量达到千万级以上,需要进一步拆分。
|
||||
### 垂直拆分的判断标准(实操 checklist)
|
||||
|
||||
当你在犹豫某张表应该留在这个库还是挪到另一个库时,用下面几个维度来判断:
|
||||
|
||||
```mermaid
|
||||
flowchart TD
|
||||
START["开始判断"] --> A["这张表和当前库里的表<br/>是否经常被同一个事务修改?"]
|
||||
A -->|是| SAME_DB["留在同一库 ✅"]
|
||||
A -->|否| B["它们是否属于同一个<br/>业务限界上下文?"]
|
||||
B -->|是| SAME_CTX["考虑放在同一库<br/>(降低跨库调用成本)"]
|
||||
B -->|否| DIFF_CTX["必须独立建库 ✅"]
|
||||
SAME_CTX --> C["变更频率差异大吗?"]
|
||||
C -->|高频 vs 低频| SPLIT["建议拆分<br/>(避免互相影响)✅"]
|
||||
C -->|同频| SAME_DB
|
||||
DIFF_CTX --> END["完成"]
|
||||
SAME_DB --> END
|
||||
SPLIT --> END
|
||||
```
|
||||
|
||||
| 判断维度 | 放同一库的信号 | 拆开的信号 |
|
||||
|---------|--------------|----------|
|
||||
| **事务耦合度** | 经常在一个 `BEGIN...COMMIT` 里一起改 | 各自有独立的写入路径 |
|
||||
| **读取热点** | 总是被一起查询、一起展示 | 访问模式完全不同 |
|
||||
| **团队归属** | 同一小组维护 | 不同团队负责(Conway 定律) |
|
||||
| **数据增长率** | 增长速度接近 | 一个涨得快、一个基本不变 |
|
||||
|
||||
> [!note] 经验法则
|
||||
>
|
||||
> 如果一组表之间的跨库 JOIN 操作占了日常 SQL 的 **80% 以上**,把它们放在一起通常更合理。反过来,如果大部分关联查询都只需要一次 LEFT JOIN,说明拆分时机已经成熟——因为 JOIN 已经在两个数据库之间产生网络开销了。
|
||||
|
||||
## 第二步:水平拆分 — 一张表太大了怎么办?
|
||||
|
||||
垂直拆分解决的是"谁来负责什么"的问题。但当 **单个服务内的单表** 达到千万级甚至亿级记录时,就需要水平拆分(也叫分片/Sharding)。
|
||||
|
||||
> [!note] 为什么要分?
|
||||
>
|
||||
> - 单表超过 2000 万行后,索引效率急剧下降
|
||||
> - 单机 MySQL 的写入 QPS 通常在 5000~20000,超出后成为瓶颈
|
||||
> - InnoDB 缓冲池装不下全部索引数据,大量磁盘 IO
|
||||
|
||||
水平拆分的核心思想:**把一张大表按规则切成多张小表,分散到不同的数据库实例中。**
|
||||
|
||||
```mermaid
|
||||
graph TB
|
||||
subgraph "分库策略"
|
||||
DB1[(DB-01)]
|
||||
DB2[(DB-02)]
|
||||
DB3[(DB-03)]
|
||||
subgraph "拆分前:一张表扛所有"
|
||||
BIG_TABLE[(order 表 1 亿行)]
|
||||
end
|
||||
|
||||
subgraph "order_0001 表"
|
||||
Row1[userId=1 → order_0001]
|
||||
Row2[userId=4 → order_0001]
|
||||
|
||||
subgraph "拆分后:按 user_id 散列"
|
||||
subgraph "DB-01"
|
||||
T1[(order_0001 2500 万行)]
|
||||
end
|
||||
subgraph "DB-02"
|
||||
T2[(order_0002 2500 万行)]
|
||||
end
|
||||
subgraph "DB-03"
|
||||
T3[(order_0003 2500 万行)]
|
||||
end
|
||||
subgraph "DB-04"
|
||||
T4[(order_0004 2500 万行)]
|
||||
end
|
||||
end
|
||||
|
||||
subgraph "order_0002 表"
|
||||
Row3[userId=2 → order_0002]
|
||||
Row4[userId=5 → order_0002]
|
||||
end
|
||||
|
||||
userId_mod["ORDER BY user_id % 2"] --> DB1
|
||||
userId_mod --> DB2
|
||||
|
||||
ROUTE["user_id % 4"] -->|"uid=1001"| T1
|
||||
ROUTE -->|"uid=2002"| T2
|
||||
ROUTE -->|"uid=3003"| T3
|
||||
ROUTE -->|"uid=4004"| T4
|
||||
|
||||
style BIG_TABLE fill:#ffebee,stroke:#ef5350
|
||||
style T1 fill:#e8f5e9,stroke:#4caf50
|
||||
style T2 fill:#e8f5e9,stroke:#4caf50
|
||||
style T3 fill:#e8f5e9,stroke:#4caf50
|
||||
style T4 fill:#e8f5e9,stroke:#4caf50
|
||||
```
|
||||
|
||||
### 常见分片策略
|
||||
#### 常见分片策略对比
|
||||
|
||||
| 策略 | 哈希公式 | 优点 | 缺点 |
|
||||
|------|---------|------|------|
|
||||
| **Hash Mod** | `user_id % N` | 简单高效,路由确定 | 扩缩容困难,数据迁移成本高 |
|
||||
| **Range** | `user_id BETWEEN x AND y` | 范围查询友好 | 热点用户集中到单分片 |
|
||||
| **Time-based** | `year_month` | 按生命周期管理 | 新分片写入压力大 |
|
||||
| **Geo-based** | 按地域分片 | 本地化访问,延迟低 | 跨区域操作复杂 |
|
||||
| 策略 | 怎么分 | 适合场景 | 代价 |
|
||||
|------|--------|---------|------|
|
||||
| **哈希取模** `hash % N` | `user_id % 4`,余数决定去哪个库 | 均匀分布,写入均衡 | 扩库时需要迁移大部分数据 |
|
||||
| **范围划分** `BETWEEN x AND y` | userId 1~10000 → DB-A,10001~20000 → DB-B | 范围查询友好 | 热点账号集中在一个分片 |
|
||||
| **时间分区** `year_month` | 2024_01 → DB-Jan, 2024_02 → DB-Feb | 按生命周期管理(冷数据归档) | 最新月份写入压力大 |
|
||||
| **地理位置** | 华东用户 → 杭州 DB,华南 → 广州 DB | 降低跨地域延迟 | 跨区域操作复杂 |
|
||||
|
||||
### ShardingSphere / MyCat
|
||||
> [!important] 扩容陷阱
|
||||
>
|
||||
> Hash Mod 最容易踩坑:当你从 4 个分片扩展到 8 个分片时,`% 4` 变成 `% 8`,**几乎所有数据的新归属都变了**,需要大规模数据迁移。这是一个需要提前规划的重大决策。
|
||||
>
|
||||
> 如果扩容是高频需求,建议从一开始就用一致性哈希或预留足够多的槽位。
|
||||
|
||||
> [!tip] 推荐中间件
|
||||
>
|
||||
> **Apache ShardingSphere** 是国内使用最广泛的分库分表方案:
|
||||
> - 支持 JDBC / Proxy / Sidecar 三种部署模式
|
||||
> - 内置分片算法:Mod、Range、Hash、Tag
|
||||
> - 分布式主键生成器(Snowflake)原生集成
|
||||
> - 读写分离、强制路由、广播表等高级特性
|
||||
### 分片扩容方案:不停机迁移
|
||||
|
||||
#### ShardingSphere 配置示例
|
||||
生产环境的扩容不能停服停机——你需要一个 **双写 + 历史数据回迁** 的渐进式流程:
|
||||
|
||||
```mermaid
|
||||
flowchart TD
|
||||
A["当前状态: N 个分片<br/>hash % N"] --> B["第一步: 新增 M 个分片<br/>总容量变为 N+M"]
|
||||
B --> C["第二步: 双写阶段<br/>新写入同时写到旧分片和新分片"]
|
||||
C --> D["第三步: 历史数据回迁<br/>按分片逐个搬移存量数据"]
|
||||
D --> E{"全部搬完?"}
|
||||
E -->|否| D
|
||||
E -->|是| F["第四步: 校验数据一致性<br/>checksum 对比"]
|
||||
F --> G["第五步: 切读流量<br/>新请求读新分片"]
|
||||
G --> H["第六步: 关闭双写<br/>恢复到单写单读"]
|
||||
|
||||
style B fill:#fff3e0,stroke:#ff9800
|
||||
style C fill:#fff3e0,stroke:#ff9800
|
||||
style D fill:#e3f2fd,stroke:#1976d2
|
||||
style F fill:#e8f5e9,stroke:#4caf50
|
||||
style G fill:#e8f5e9,stroke:#4caf50
|
||||
style H fill:#f3e5f5,stroke:#7b1fa2
|
||||
```
|
||||
|
||||
| 阶段 | 核心操作 | 关键风险 | 规避方法 |
|
||||
|------|---------|---------|---------|
|
||||
| **双写** | 新数据同时写入新旧两套分片 | 数据不一致、写入性能下降 | 用消息队列保证异步双写;设置标记位可快速回滚 |
|
||||
| **回迁** | 按 user_id 范围分批迁移历史数据 | 迁移期间持续写入导致数据漂移 | 迁移后对已迁移范围做一次增量同步 |
|
||||
| **切读** | 将读取路由切换到新分片 | 漏读、脏数据 | 先灰度 1% 流量验证,逐步放量 |
|
||||
| **关双写** | 停止向旧分片写入 | 遗漏最后一段增量数据 | 双写关之前做一次全量 checksum 校验 |
|
||||
|
||||
> [!warning] 平滑扩容的时间成本
|
||||
>
|
||||
> 假设你有 1 亿条订单数据,网络带宽 1Gbps,压缩后约 50GB。**理论传输时间不到 1 分钟**——但实际中还要考虑锁竞争、慢查询、监控告警等因素。建议给每个分片预留 **2~4 小时** 的迁移窗口期,夜间低峰期执行。
|
||||
|
||||
## 业界方案:ShardingSphere
|
||||
|
||||
> [!tip] 为什么选 ShardingSphere?
|
||||
>
|
||||
> Apache ShardingSphere 是国内使用最广泛的分库分表中间件。它提供三种部署模式:
|
||||
> - **JDBC**:嵌入应用,零运维(最常用)
|
||||
> - **Proxy**:独立代理服务,语言无关
|
||||
> - **Sidecar**:Kubernetes 侧车模式
|
||||
|
||||
#### 配置示例
|
||||
|
||||
以下是 ShardingSphere-JDBC 的核心配置片段(YAML 格式):
|
||||
|
||||
```yaml
|
||||
# sharding-jdbc 配置
|
||||
sharding jdbc:
|
||||
data-sources:
|
||||
ds0: { type: HikariCP, ... }
|
||||
ds1: { type: HikariCP, ... }
|
||||
|
||||
sharding-jdbc:
|
||||
data-sources: # 定义数据源
|
||||
ds0: { type: com.zaxxer.hikari.HikariDataSource, ... }
|
||||
ds1: { type: com.zaxxer.hikari.HikariDataSource, ... }
|
||||
|
||||
sharding:
|
||||
tables:
|
||||
orders:
|
||||
orders: # 逻辑表名
|
||||
actual-data-nodes: ds$->{0..1}.orders$->{0..1}
|
||||
# 上面的表达式展开后是:ds0.orders0, ds0.orders1, ds1.orders0, ds1.orders1
|
||||
table-strategy:
|
||||
standard:
|
||||
sharding-column: user_id
|
||||
sharding-column: user_id # 分片键
|
||||
sharding-algorithm-name: user-id-mod
|
||||
key-generate-strategy:
|
||||
column: order_id
|
||||
key-generator-name: snowflake
|
||||
|
||||
column: order_id # 主键生成
|
||||
key-generator-name: snowflake # Snowflake 雪花算法
|
||||
|
||||
sharding-algorithms:
|
||||
user-id-mod:
|
||||
type: MOD
|
||||
props:
|
||||
sharding-count: 2
|
||||
sharding-count: 2 # 分成 2 个分片
|
||||
```
|
||||
|
||||
## 跨库查询方案
|
||||
关键点:
|
||||
- **逻辑表名 vs 实际数据节点**:应用层看到的永远是 `orders`,底层路由到 `ds0.orders0` 等是由中间件透明的完成的
|
||||
- **分片键选择**:一定要选写入频率高且用于查询条件的字段(如 `user_id`),否则每次查询都要扫全部分片
|
||||
|
||||
### 分布式 ID 生成:为什么不能用自增主键?
|
||||
|
||||
分库后每个库的 `AUTO_INCREMENT` 是独立的——两个库可能都生成了 `id = 100`。你必须用一种方式保证 **全局唯一**。
|
||||
|
||||
#### Snowflake 雪花算法
|
||||
|
||||
Twitter 开源的雪花算法是目前最主流的分布式 ID 方案:
|
||||
|
||||
```mermaid
|
||||
graph LR
|
||||
subgraph "64-bit 长整型 ID"
|
||||
S["符号位<br/>1 bit"]
|
||||
T["时间戳<br/>41 bits"]
|
||||
D["机器 ID<br/>10 bits"]
|
||||
SQ["序列号<br/>12 bits"]
|
||||
end
|
||||
|
||||
style S fill:#f5f5f5,stroke:#9e9e9e
|
||||
style T fill:#e3f2fd,stroke:#1976d2
|
||||
style D fill:#e8f5e9,stroke:#4caf50
|
||||
style SQ fill:#fff3e0,stroke:#ff9800
|
||||
```
|
||||
|
||||
| 字段 | 长度 | 作用 | 范围 |
|
||||
|------|------|------|------|
|
||||
| **符号位** | 1 bit | 恒为 0(保证 ID 为正数) | - |
|
||||
| **时间戳** | 41 bits | 毫秒级时间戳 | 可支撑约 69 年 |
|
||||
| **机器 ID** | 10 bits | 区分部署实例 | 1024 个节点 |
|
||||
| **序列号** | 12 bits | 同一毫秒内的递增序号 | 每毫秒 4096 个 ID |
|
||||
|
||||
> [!note] 核心特性
|
||||
>
|
||||
> - **单调递增**:基于时间戳保证整体趋势递增,适合 InnoDB 聚簇索引的 append-only 写入模式
|
||||
> - **高吞吐**:单实例每秒可生成 4096 × 1000 = 数百万个 ID
|
||||
> - **无中心节点**:不依赖 ZooKeeper 或数据库,挂掉一个机器不影响其他节点
|
||||
>
|
||||
> 如果同一毫秒内生成的 ID 超过 4096 个,算法会 **等待下一毫秒** 再重试——这在实际场景中极少发生(通常 QPS < 10 万)。
|
||||
|
||||
#### 备选方案对比
|
||||
|
||||
| 方案 | 唯一性保证 | 性能 | 复杂度 | 适用场景 |
|
||||
|------|----------|------|--------|---------|
|
||||
| **Snowflake** | 时间戳+机器ID+序列号组合 | 极高(本地生成) | 中 | 通用首选 |
|
||||
| **数据库号段模式** | 每次从 DB 批量拉取一段 ID | 高 | 中高 | 有现成 DB 基础设施 |
|
||||
| **UUID** | 随机 128 位 | 高 | 极低 | 对有序性无要求的场景 |
|
||||
| **Redis INCR** | Redis 原子递增 | 高 | 低 | 已有 Redis 集群 |
|
||||
|
||||
> [!warning] UUID 的坑
|
||||
>
|
||||
> UUID 虽然简单,但在 MySQL InnoDB 中是 **灾难性的**:因为 UUID 无序,插入位置随机分布,导致大量的页分裂和碎片化。如果必须用 UUID,建议存为 `BINARY(16)` 并用 `UUID_TO_BIN(uuid, 1)` 转换,让它在索引中保持局部有序。
|
||||
|
||||
## 跨库查询:拆分之后最难解决的问题之一
|
||||
|
||||
> [!question] 经典难题
|
||||
> 订单服务需要展示用户的姓名和手机号来做收货地址。但用户信息在用户库,订单数据在订单库——怎么办?
|
||||
>
|
||||
> 用户打开"我的订单"页面。订单服务拿到 `user_id` 查出订单列表,但现在要展示用户的头像和昵称——这些信息在用户库里。数据库已经拆开了,怎么做关联查询?
|
||||
|
||||
### 方案对比
|
||||
这是拆分后必然遇到的挑战。没有银弹,只有 **权衡后的取舍**。
|
||||
|
||||
| 方案 | 描述 | 性能 | 复杂度 | 适用场景 |
|
||||
|------|------|------|--------|---------|
|
||||
| **冗余字段** | 订单表存用户名字段 | ⭐⭐⭐⭐⭐ | 低 | 只读字段,变更频率低 |
|
||||
| **接口组装** | 先查订单,再调用户服务补全 | ⭐⭐⭐ | 中 | 偶尔需要关联的场景 |
|
||||
| **CQRS / 宽表** | 异步同步一份宽表用于查询 | ⭐⭐⭐⭐ | 高 | 高频关联查询 |
|
||||
| **搜索引擎** | ES/Kibana 做关联查询 | ⭐⭐⭐⭐ | 中高 | 复杂搜索 + 聚合 |
|
||||
### 四种方案对比
|
||||
|
||||
### 冗余字段实践
|
||||
```mermaid
|
||||
mindmap
|
||||
root((跨库查询方案))
|
||||
冗余字段
|
||||
最简单最直接
|
||||
少量字段
|
||||
快照语义
|
||||
一致性问题
|
||||
接口组装
|
||||
按需查询
|
||||
链路长时性能差
|
||||
N+1 查询陷阱
|
||||
适合低频关联
|
||||
CQRS / 宽表
|
||||
异步最终一致
|
||||
写入有额外开销
|
||||
适合高频关联
|
||||
搜索引擎
|
||||
ES 做多维聚合
|
||||
架构重
|
||||
适合复杂搜索
|
||||
```
|
||||
|
||||
| 方案 | 一句话描述 | 性能 | 复杂度 |
|
||||
|------|-----------|------|--------|
|
||||
| **冗余字段** | 在订单表里直接存用户名字段 | ⭐⭐⭐⭐⭐ | 低 |
|
||||
| **接口组装** | 先查订单,循环调用户服务补全信息 | ⭐⭐⭐ | 中 |
|
||||
| **CQRS / 宽表** | 异步同步一份含用户信息的宽表 | ⭐⭐⭐⭐ | 高 |
|
||||
| **搜索引擎** | 把数据推到 ES,用 ES 做关联查询 | ⭐⭐⭐⭐ | 中高 |
|
||||
|
||||
### 实战一:冗余字段(最推荐的首选方案)
|
||||
|
||||
```sql
|
||||
-- 订单表冗余关键字段
|
||||
-- 订单表中冗余关键字段(快照模式)
|
||||
CREATE TABLE orders (
|
||||
id BIGINT PRIMARY KEY,
|
||||
user_id BIGINT NOT NULL,
|
||||
username VARCHAR(64), -- 冗余用户名(快照,不参与编辑)
|
||||
phone VARCHAR(20), -- 冗余手机号(脱敏存储)
|
||||
username VARCHAR(64), -- 下单时的用户名快照
|
||||
phone VARCHAR(20), -- 下单时的手机号(已脱敏)
|
||||
created_at TIMESTAMP DEFAULT NOW(),
|
||||
INDEX idx_user (user_id)
|
||||
);
|
||||
```
|
||||
|
||||
> [!note] 冗余数据的维护
|
||||
>
|
||||
> 用户改名字了怎么办?**不改订单表**。订单上的 username 是该时刻的"快照"——它反映的是下单时的状态,不是当前状态。这符合业务语义。
|
||||
>
|
||||
> 如果需要批量更新冗余字段(如用户头像),通过消息队列异步通知订单服务更新。
|
||||
> [!note] 关键理解:快照语义
|
||||
>
|
||||
> 用户后来改了名字、换了手机号,**不改订单表**。订单上的 `username` 反映的是下单那一刻的状态,而不是"当前"状态。这不仅是合理的,而且是正确的——用户查看历史订单时,看到的是当时的信息。
|
||||
>
|
||||
> 如果需要主动更新冗余字段(如用户更换了头像),通过消息队列通知订单服务批量更新。
|
||||
|
||||
## 数据库迁移工具
|
||||
Go 代码实现示例:
|
||||
|
||||
```mermaid
|
||||
graph LR
|
||||
Dev["开发环境"] -->|"flyway migrate"| Stage["Staging"]
|
||||
Stage -->|"人工审批"| Prod["Production"]
|
||||
Prod -->|"flyway validate"| Check["校验版本一致性"]
|
||||
```go
|
||||
// 下单时将用户信息快照写入订单表
|
||||
func (s *orderSvc) CreateOrder(ctx context.Context, req *CreateOrderReq) (*Order, error) {
|
||||
// 1. 查用户基本信息
|
||||
user, err := s.userClient.GetByID(ctx, req.UserID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// 2. 构造订单,携带快照字段
|
||||
order := &Order{
|
||||
UserID: req.UserID,
|
||||
Username: user.Username, // ✅ 快照:锁定下单时的值
|
||||
Phone: maskPhone(user.Phone),
|
||||
Items: req.Items,
|
||||
}
|
||||
return s.orderRepo.Save(ctx, order)
|
||||
}
|
||||
```
|
||||
|
||||
| 工具 | 语言 | 特点 |
|
||||
|------|------|------|
|
||||
| **Flyway** | Java | 基于文件命名,SQL 脚本方式,简单直观 |
|
||||
| **Liquibase** | Java | XML/YAML/JSON 格式,支持回滚生成 |
|
||||
| **golang-migrate** | Go | CLI 工具,轻量,适合 Go 项目 |
|
||||
### 实战二:接口组装(轻量场景够用)
|
||||
|
||||
### Flyway 迁移流程
|
||||
适用于关联查询不频繁的场景,比如后台管理系统的偶尔查看详情。
|
||||
|
||||
```go
|
||||
// 查订单 + 补齐用户信息
|
||||
func (s *orderSvc) GetOrderWithUser(ctx context.Context, orderID int64) (*OrderDetail, error) {
|
||||
// 1. 先查订单
|
||||
order, err := s.orderRepoFindByID(ctx, orderID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// 2. 再调用户服务补齐
|
||||
user, err := s.userClient.GetByID(ctx, order.UserID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &OrderDetail{
|
||||
Order: order,
|
||||
UserName: user.Username,
|
||||
UserHead: user.AvatarURL,
|
||||
}, nil
|
||||
}
|
||||
```
|
||||
|
||||
> [!warning] N+1 陷阱
|
||||
>
|
||||
> 如果是列表查询(一次性返回 20 条订单),逐条调用户服务会导致 20 次 RPC 调用。正确做法是:先收集所有 `user_id`,**批量查询**用户信息,再拼回去。
|
||||
|
||||
```go
|
||||
// ❌ 错误:N 次 RPC
|
||||
for _, o := range orders {
|
||||
u, _ := userClient.GetByID(ctx, o.UserID)
|
||||
}
|
||||
|
||||
// ✅ 正确:1 次批量 RPC
|
||||
ids := collectIDs(orders) // [1, 3, 7, 12, ...]
|
||||
users, _ := userClient.GetByIds(ctx, ids) // 一次拿回所有
|
||||
byID := indexBy(users, func(u *User) int64 { return u.ID })
|
||||
for _, o := range orders {
|
||||
o.UserInfo = byID[o.UserID]
|
||||
}
|
||||
```
|
||||
|
||||
### 实战三:CQRS / 宽表(重度关联查询必备)
|
||||
|
||||
当跨库关联是高频操作(如运营后台的多维度筛选),冗余字段就不够了——你需要一套完整的 **事件驱动宽表同步机制**:
|
||||
|
||||
```mermaid
|
||||
sequenceDiagram
|
||||
participant U as 用户服务
|
||||
participant MQ as 消息队列
|
||||
participant O as 订单服务
|
||||
participant W as 宽表 (Read DB)
|
||||
|
||||
U->>MQ: UserUpdated 事件 (userId, newAvatar)
|
||||
MQ->>O: 消费事件
|
||||
O->>W: 更新宽表中对应用户的头像
|
||||
|
||||
Note over W: 宽表包含了订单 + 用户 + 商品的冗余字段<br/>只读,专为查询优化
|
||||
```
|
||||
|
||||
Go 代码实现——事件消费者(订单服务侧):
|
||||
|
||||
```go
|
||||
// Order宽表同步消费者:监听来自各服务的业务事件,维护一张可跨维度查询的宽表
|
||||
type WideTableSyncer struct {
|
||||
wideDB *sql.DB // 独立的读库连接
|
||||
}
|
||||
|
||||
// OnUserUpdated 消费用户变更事件,更新宽表中对应用户的信息
|
||||
func (s *WideTableSyncer) OnUserUpdated(ctx context.Context, evt UserUpdatedEvent) error {
|
||||
tx, err := s.wideDB.BeginTx(ctx, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer tx.Rollback()
|
||||
|
||||
// 批量更新宽表中的用户快照字段
|
||||
_, err = tx.ExecContext(ctx,
|
||||
`UPDATE order_wide_table
|
||||
SET username = ?, phone_masked = ?, avatar_url = ?
|
||||
WHERE user_id = ?`,
|
||||
evt.Username, evt.PhoneMasked, evt.AvatarURL, evt.UserID,
|
||||
)
|
||||
return tx.Commit()
|
||||
}
|
||||
|
||||
// OnOrderCreated 订单创建时写入宽表(包含完整的关联信息)
|
||||
func (s *WideTableSyncer) OnOrderCreated(ctx context.Context, evt OrderCreatedEvent) error {
|
||||
_, err := s.wideDB.ExecContext(ctx,
|
||||
`INSERT INTO order_wide_table
|
||||
(order_id, user_id, username, product_name, amount, status, created_at)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?)`,
|
||||
evt.OrderID, evt.UserID, evt.Username,
|
||||
evt.ProductName, evt.Amount, evt.Status, evt.CreatedAt,
|
||||
)
|
||||
return err
|
||||
}
|
||||
```
|
||||
|
||||
> [!tip] 宽表设计的三个原则
|
||||
>
|
||||
> 1. **反范式化**:宽表故意违反第一范式——同一个用户的名字可能出现在几百条订单记录中。这正是它的价值所在。
|
||||
> 2. **最终一致性延迟可控**:通过 MQ 保证,通常在 1~3 秒内同步完成。前端加一个 loading 状态即可掩盖这短暂的延迟。
|
||||
> 3. **写多读少时才考虑**:如果宽表的写入放大超过原始数据的 3 倍,说明你的场景用冗余字段就够了,不需要上宽表。
|
||||
|
||||
## 何时该拆?——不要过早拆分
|
||||
|
||||
拆分是有 **代价** 的:复杂度上升、运维成本增加、跨服务调用变慢。在决定拆之前,先看这些硬指标是否已经触达:
|
||||
|
||||
```mermaid
|
||||
flowchart TD
|
||||
START["你的数据量到了什么级别?"]
|
||||
|
||||
START --> SMALL["< 100 万行<br/>单库单机足够"]
|
||||
START --> MEDIUM["100 万 ~ 2000 万行<br/>考虑读写分离"]
|
||||
START --> LARGE["> 2000 万行<br/>考虑分片"]
|
||||
|
||||
SMALL --> S1["✅ 先优化索引"]
|
||||
S1 --> S2["✅ 加缓存层 Redis"]
|
||||
S2 --> S3["✅ 读写分离<br/>一主多从"]
|
||||
|
||||
MEDIUM --> M1["✅ 先做垂直拆分"]
|
||||
M1 --> M2["按业务域独立建库"]
|
||||
|
||||
LARGE --> L1["✅ 再做水平拆分"]
|
||||
L1 --> L2["分库分表 + 宽表同步"]
|
||||
|
||||
style SMALL fill:#e8f5e9,stroke:#4caf50
|
||||
style MEDIUM fill:#fff3e0,stroke:#ff9800
|
||||
style LARGE fill:#ffebee,stroke:#ef5350
|
||||
style S3 fill:#e8f5e9,stroke:#4caf50
|
||||
style M2 fill:#e8f5e9,stroke:#4caf50
|
||||
style L2 fill:#ffebee,stroke:#ef5350
|
||||
```
|
||||
|
||||
> [!quote] 一条经验法则
|
||||
>
|
||||
> **"在数据量还没到 2000 万行之前,不要做任何形式的数据水平拆分。"**
|
||||
>
|
||||
> 绝大多数系统通过索引优化 + 缓存 + 读写分离就能撑到日活百万级。见过太多团队在项目刚上线就搞分库分表——六个月后回头看,完全是提前踩雷。
|
||||
|
||||
### 拆与不拆的判断矩阵
|
||||
|
||||
| 场景 | 推荐方案 | 预期寿命 |
|
||||
|------|---------|---------|
|
||||
| DAU < 1 万,QPS < 500 | 单库单机,专心做好索引 | 半年~1年 |
|
||||
| DAU 1~50 万,QPS 500~5000 | 加 Redis 缓存 + 读写分离 | 1~2 年 |
|
||||
| DAU 50~200 万,单表 > 2000 万行 | 垂直拆分(按微服务) | 2~3 年 |
|
||||
| DAU > 200 万,单表 > 5000 万行 | 水平拆分 + 宽表同步 | 长期 |
|
||||
|
||||
## 数据库迁移工具:告别手动执行 SQL
|
||||
|
||||
随着服务拆分增多,SQL 脚本的管理变得复杂。手工执行、口头传达"我跑过 V3 了"的方式不再可行——你需要**版本化的、可重复执行的数据库迁移工具**。
|
||||
|
||||
```mermaid
|
||||
flowchart LR
|
||||
Dev["💻 开发环境<br/>git commit SQL 文件"] -->|"CI/CD 自动执行"| Stage["🧪 Staging<br/>flyway migrate"]
|
||||
Stage -->|"人工审批"| Prod["🚀 Production<br/>flyway migrate"]
|
||||
Prod -->|"校验"| Check["🔒 flyway validate<br/>确认无漂移"]
|
||||
|
||||
style Dev fill:#e3f2fd,stroke:#1976d2
|
||||
style Stage fill:#fff3e0,stroke:#ff9800
|
||||
style Prod fill:#e8f5e9,stroke:#4caf50
|
||||
style Check fill:#f3e5f5,stroke:#7b1fa2
|
||||
```
|
||||
|
||||
### 主流工具对比
|
||||
|
||||
| 工具 | 生态 | 特点 |
|
||||
|------|------|------|
|
||||
| **Flyway** | Java / Go (via CLI) | 基于文件名命名,纯 SQL 脚本,简单直接 |
|
||||
| **Liquibase** | Java | 支持 XML/YAML/JSON,自带回滚生成能力 |
|
||||
| **golang-migrate** | Go | 轻量 CLI,Go 项目首选 |
|
||||
|
||||
### Flyway 工作流
|
||||
|
||||
```
|
||||
db/migration/
|
||||
@@ -190,7 +578,12 @@ db/migration/
|
||||
└── V4__add_order_status_enum.sql
|
||||
```
|
||||
|
||||
执行顺序:V1 → V2 → V3 → V4。Flyway 内部维护 `_schema_version` 表追踪已执行的迁移。
|
||||
执行顺序严格遵循版本号:V1 → V2 → V3 → V4。Flyway 会在数据库中维护一张 `_schema_version` 表,记录每个迁移的版本和执行状态。
|
||||
|
||||
核心要点:
|
||||
- **文件名即版本**:`V` 开头 + 序号 + `__` + 描述
|
||||
- **不可修改已执行的文件**:改了的话 Flyway 会报验证失败(`validate` 阶段)
|
||||
- **永远不要写"反向"SQL**:迁移脚本只做"升级",回滚通过发版解决
|
||||
|
||||
## 关联笔记
|
||||
|
||||
|
||||
+320
-21
@@ -1,6 +1,6 @@
|
||||
---
|
||||
tags: [microservice, distributed-transactions, saga, tcc, outbox]
|
||||
create time: 2026-05-05
|
||||
tags: [microservice, distributed-transactions, saga, tcc, outbox, seata, idempotency, dlq]
|
||||
create time: 2026-05-17 15:00
|
||||
---
|
||||
|
||||
# 分布式事务
|
||||
@@ -101,25 +101,6 @@ func outboxWorker(ctx context.Context, ticker *time.Ticker) {
|
||||
}
|
||||
```
|
||||
|
||||
#### 事务消息(RocketMQ 原生支持)
|
||||
|
||||
如果使用的是 RocketMQ,可以绕过 Outbox 模式直接用事务消息:
|
||||
|
||||
```go
|
||||
// 发送事务消息
|
||||
txMsg := rocketmq.NewTransactionMessage("order-created", payload)
|
||||
localTx := &MyLocalTxChecker{}
|
||||
|
||||
// half 消息发送 → 本地事务执行 → 提交/回查
|
||||
res, _ := producer.SendMessageInTransaction(txMsg, localTx)
|
||||
```
|
||||
|
||||
RocketMQ 的事务消息流程:
|
||||
1. 生产者发送 "half 消息" 到 MQ(消费者不可见)
|
||||
2. 执行本地事务
|
||||
3. 根据结果 Commit(消费者可见)或 Rollback(丢弃)
|
||||
4. 如果步骤 2 超时,MQ 回查本地事务状态
|
||||
|
||||
### 消费幂等性
|
||||
|
||||
> [!warning] 关键保障
|
||||
@@ -148,6 +129,57 @@ if !isUniqueViolation(err) {
|
||||
// 执行业务逻辑——到这里说明消息是新到达的
|
||||
```
|
||||
|
||||
#### 事务消息 vs Outbox
|
||||
|
||||
> [!question] 思考一下
|
||||
> Outbox 模式有一个天然缺点:定时轮询有延迟(通常秒级)。如果你需要**毫秒级的延迟敏感型通知**(如支付结果实时推送到用户端),该怎么办?
|
||||
|
||||
如果使用的是 RocketMQ,可以绕过 Outbox 模式直接用事务消息:
|
||||
|
||||
```go
|
||||
// 发送事务消息
|
||||
txMsg := rocketmq.NewTransactionMessage("order-created", payload)
|
||||
localTx := &MyLocalTxChecker{}
|
||||
|
||||
// half 消息发送 → 本地事务执行 → 提交/回查
|
||||
res, _ := producer.SendMessageInTransaction(txMsg, localTx)
|
||||
```
|
||||
|
||||
RocketMQ 事务消息的完整交互流程:
|
||||
|
||||
```mermaid
|
||||
sequenceDiagram
|
||||
participant Producer as 生产者
|
||||
participant Broker as MQ Broker
|
||||
participant DB as 业务数据库
|
||||
participant Consumer as 消费者
|
||||
|
||||
Producer->>Broker: ① 发送 Half Message (消费者不可见)
|
||||
Broker-->>Producer: ACK
|
||||
|
||||
Producer->>DB: ② 执行本地事务
|
||||
alt 本地事务成功
|
||||
Producer->>Broker: ③a Commit (消息对消费者可见)
|
||||
else 本地事务失败
|
||||
Producer->>Broker: ③b Rollback (消息丢弃)
|
||||
end
|
||||
|
||||
Note over Producer,Broker: ④ 如果步骤 2 超时未回复
|
||||
Broker->>Producer: ④ 回查请求 (CheckLocalTx)
|
||||
Producer->>DB: 查询本地事务状态
|
||||
DB-->>Producer: 返回 COMMIT / ROLLBACK
|
||||
Producer->>Broker: 回查结果
|
||||
```
|
||||
|
||||
> [!important] 回查机制的意义
|
||||
> 当生产者执行本地事务过程中发生崩溃或网络中断,Broker 收不到 Commit/Rollback 指令,就会通过回查主动询问生产者。这确保了 **half 消息最终一定会收敛到可见或被丢弃**,不会出现悬空状态。
|
||||
|
||||
**可靠性要点清单:**
|
||||
- ✅ 业务数据库和 Outbox 必须在**同一个本地事务**中写入
|
||||
- ✅ 消费者必须有幂等保护(否则至少一次投递 = 无限次重复)
|
||||
- ✅ 定时扫描任务需要设置最大重试次数,超过后进入人工处理
|
||||
- ✅ 生产环境建议配置 **DLQ(死信队列)**,将连续消费失败的消息归档
|
||||
|
||||
## 方案二:Saga 模式
|
||||
|
||||
Saga 适用于**跨多个服务的长流程操作**,将大事务拆成一系列本地小事务,每个步骤都有对应的补偿操作。
|
||||
@@ -198,6 +230,81 @@ graph TB
|
||||
> [!warning] 补偿操作的幂等性
|
||||
> 补偿操作也必须幂等——CancelOrder 可能被触发多次。用订单状态的流转来保证(如只有 PENDING 才能转到 CANCELLED)。
|
||||
|
||||
### 编排式实现示例 (Go)
|
||||
|
||||
编排式的核心是一个 **Saga Coordinator**,它维护着当前步骤的执行状态和已完成的正向/反向操作列表:
|
||||
|
||||
```go
|
||||
// Step 定义一个 Saga 步骤
|
||||
type Step struct {
|
||||
Name string
|
||||
Execute func(ctx context.Context) error // 正向操作
|
||||
Compensation func(ctx context.Context) error // 补偿操作
|
||||
}
|
||||
|
||||
// Orchestrator 编排器
|
||||
type Orchestrator struct {
|
||||
steps []Step
|
||||
executed []string // 已完成步骤名,用于补偿时逆向遍历
|
||||
}
|
||||
|
||||
func (o *Orchestrator) AddStep(name string, exec, compens func(context.Context) error) {
|
||||
o.steps = append(o.steps, Step{Name: name, Execute: exec, Compensation: compens})
|
||||
}
|
||||
|
||||
// Execute 按顺序执行所有步骤,任一步骤失败则逆向补偿
|
||||
func (o *Orchestrator) Execute(ctx context.Context) error {
|
||||
for i, step := range o.steps {
|
||||
if err := step.Execute(ctx); err != nil {
|
||||
log.Error("step failed", z.String("step", step.Name), z.Err(err))
|
||||
// 逆向补偿:从后往前执行已完步骤的补偿操作
|
||||
o.rollback(ctx, i-1)
|
||||
return fmt.Errorf("saga failed at %s: %w", step.Name, err)
|
||||
}
|
||||
o.executed = append(o.executed, step.Name)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (o *Orchestrator) rollback(ctx context.Context, upTo int) {
|
||||
for i := len(o.executed) - 1; i >= 0 && i <= upTo; i-- {
|
||||
// 找到对应步骤的补偿函数并执行
|
||||
step := o.steps[i]
|
||||
if err := step.Compensation(ctx); err != nil {
|
||||
log.Warn("compensation also failed — need manual intervention", z.String("step", step.Name))
|
||||
// 补偿失败必须告警!人工介入处理
|
||||
}
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
使用方式——构建一个下单 Saga:
|
||||
|
||||
```go
|
||||
saga := &Orchestrator{}
|
||||
|
||||
saga.AddStep("create_order",
|
||||
func(ctx context.Context) error { return orderSvc.Create(ctx, req) }, // 正向
|
||||
func(ctx context.Context) error { return orderSvc.Cancel(ctx, orderID) }, // 反向
|
||||
)
|
||||
saga.AddStep("reserve_stock",
|
||||
func(ctx context.Context) error { return stockSvc.Reserve(ctx, productIDs) },
|
||||
func(ctx context.Context) error { return stockSvc.Release(ctx, productIDs) },
|
||||
)
|
||||
saga.AddStep("charge_payment",
|
||||
func(ctx context.Context) error { return paySvc.Charge(ctx, amount) },
|
||||
func(ctx context.Context) error { return paySvc.Refund(ctx, orderID) },
|
||||
)
|
||||
|
||||
err := saga.Execute(ctx)
|
||||
```
|
||||
|
||||
> [!note] 关键理解:补偿的局限性
|
||||
>
|
||||
> - **补偿不是撤销**:取消订单不会把商品退回到原始状态,而是创建一条反向业务记录
|
||||
> - **补偿可能失败**:如果退款接口也挂了,你需要重试 + 告警 + 人工介入机制
|
||||
> - **时间窗口问题**:如果正向前向到了第 5 步、补偿到了第 2 步时第 2 步也失败了,第 3~4 步的正向操作已经发生且无法撤销
|
||||
|
||||
## 方案三:TCC (Try-Confirm-Cancel)
|
||||
|
||||
TCC 在每个事务步骤中实现三个接口:
|
||||
@@ -242,6 +349,198 @@ sequenceDiagram
|
||||
| 性能 | 中(需要两阶段提交) | 中高(单阶段本地事务) |
|
||||
| 适用场景 | 资金、库存等高敏感业务 | 订单流程、审批流等业务链 |
|
||||
|
||||
## 方案四:AT 模式 (Seata)
|
||||
|
||||
> [!question] 什么时候该用 Seata?
|
||||
>
|
||||
> 如果你的团队**不愿或没有能力为每个业务方法编写 TCC 接口**,但又有跨库事务需求——Seata 的 AT 模式就是为此设计的。它是"零侵入"的最强候选方案。
|
||||
|
||||
Seata AT 模式的核心原理是 **全局锁 + 二阶段提交**,业务代码不需要任何改动:
|
||||
|
||||
```
|
||||
第一阶段(Global Lock):
|
||||
1. 拦截 SQL → 解析前后镜像
|
||||
2. 申请全局锁(锁住被修改的行)
|
||||
3. 本地事务执行(含 undo_log 写入)
|
||||
4. 提交本地事务(全局锁暂不释放)
|
||||
|
||||
第二阶段(Global Commit/Rollback):
|
||||
Commit: 异步删除 undo_log,释放全局锁
|
||||
Rollback: 利用 undo_log 生成反向 SQL,回滚并释放全局锁
|
||||
```
|
||||
|
||||
### AT 模式时序图
|
||||
|
||||
```mermaid
|
||||
sequenceDiagram
|
||||
participant TC as Transaction Coordinator<br/>(Seata Server)
|
||||
participant TM as Transaction Manager<br/>(应用 A)
|
||||
participant RS1 as Resource 1<br/>(订单 DB)
|
||||
participant RS2 as Resource 2<br/>(库存 DB)
|
||||
|
||||
TM->>TC: 开启全局事务 (xid)
|
||||
TM->>RS1: BEGIN; UPDATE orders SET ...
|
||||
Note over RS1: 自动捕获 BEFORE/AFTER 镜像<br/>写入 undo_log
|
||||
|
||||
TM->>RS2: BEGIN; UPDATE stock SET qty = qty - 1
|
||||
Note over RS2: 同样写 undo_log
|
||||
|
||||
TM->>RS1: COMMIT (本地事务完成)
|
||||
TM->>RS2: COMMIT (本地事务完成)
|
||||
|
||||
TM->>TC: 报告二阶段提交完成
|
||||
TC-->>RS1: 异步清理 undo_log
|
||||
TC-->>RS2: 异步清理 undo_log
|
||||
```
|
||||
|
||||
### AT 与 XA 的区别
|
||||
|
||||
很多人混淆 AT 和 XA——它们都是两阶段提交,但实现方式完全不同:
|
||||
|
||||
| 维度 | XA 模式 | AT 模式 (Seata) |
|
||||
|------|---------|----------------|
|
||||
| **锁粒度** | 数据库级(长事务锁) | 行级(短期全局锁) |
|
||||
| **隔离级别** | 需要 SERIALIZABLE | 基于脏读 + 全局锁 |
|
||||
| **性能** | 较差(锁持有时间长) | 较好(本地事务快速提交后异步清理) |
|
||||
| **侵入性** | 需配置 DataSource Proxy | 零侵入,自动代理 JDBC |
|
||||
|
||||
### 使用注意事项
|
||||
|
||||
```go
|
||||
// Seata Go SDK — 几乎零侵入
|
||||
import "github.com/seata/seata-sdk-go-client/tx"
|
||||
|
||||
func CreateOrder(ctx context.Context, req *Request) error {
|
||||
// 只需添加这一个注解
|
||||
tx.GlobalTransaction()
|
||||
|
||||
// 下面的代码完全不变——Seata 自动代理数据源
|
||||
orderRepo.Create(ctx, &order)
|
||||
stockRepo.Decrease(ctx, productID, qty)
|
||||
|
||||
return nil
|
||||
}
|
||||
```
|
||||
|
||||
> [!warning] 生产环境避坑清单
|
||||
>
|
||||
> 1. **undo_log 表必须建**:`seata_undo_log` 是 Seata 自动创建的,用于存储前后镜像。务必保证这个表可用(通常和业务库在同一实例)
|
||||
> 2. **不要混用本地事务和全局事务**:同一个方法里同时出现 `@Transactional` 和 `@GlobalTransactional` 会导致行为不可预测
|
||||
> 3. **异常处理**:全局事务中抛出的异常会被 Seata 捕获并触发回滚,确保异常能正确向上传递
|
||||
> 4. **性能考量**:全局锁虽然比 XA 短,但在高并发场景下仍然是瓶颈。**优先用 MQ 最终一致,实在做不到再上 AT**
|
||||
> 5. **只支持部分数据库**:MySQL、PostgreSQL、Oracle 支持良好;TiDB 等 NewSQL 数据库兼容性有限
|
||||
|
||||
> [!tip] 选型建议
|
||||
>
|
||||
> AT 模式最适用的场景:**遗留系统微服务化改造**。当原有单体拆分成多个服务时,原有的 `@Transactional` 跨库调用突然失效,引入 AT 模式可以快速过渡——等业务稳定后再逐步迁移到 MQ 事件驱动架构。
|
||||
|
||||
## 补充保障:重试与死信队列
|
||||
|
||||
分布式系统中,**没有哪个组件永远可靠**。消息消费失败、网络闪断、下游超时——这些都不可避免。你需要一套完整的错误处理链路:
|
||||
|
||||
```mermaid
|
||||
flowchart LR
|
||||
MQ["消息队列"] --> C["消费者"]
|
||||
C --> FAIL{消费成功}
|
||||
FAIL -- 是 --> OK["正常结束"]
|
||||
FAIL -- 否 --> RETRY{重试次数未满}
|
||||
RETRY -- 是 --> DELAY["延迟重试<br/>指数退避"]
|
||||
DELAY --> C
|
||||
RETRY -- 否 --> DLQ["进入死信队列<br/>人工介入"]
|
||||
|
||||
DLQ --> ALERT["告警通知"]
|
||||
DLQ --> REPLAY["手动重放或修复后补发"]
|
||||
```
|
||||
|
||||
### 三种重试策略对比
|
||||
|
||||
| 策略 | 适用场景 | 示例 |
|
||||
|------|---------|------|
|
||||
| **立即重试** | 瞬时故障(连接池未就绪) | 等 100ms 再试 3 次 |
|
||||
| **延迟重试** | 需要时间窗口等待恢复 | 1s → 2s → 4s → 8s |
|
||||
| **定时任务补偿** | 批量漏掉的消息 | 每 5 分钟扫描 pending 表 |
|
||||
|
||||
### 死信队列 (DLQ) 设计
|
||||
|
||||
```go
|
||||
// 消费者主循环:带重试和 DLQ 保护
|
||||
func (c *consumer) Process(ctx context.Context, msg Message) error {
|
||||
for attempt := 0; attempt < c.maxRetries; attempt++ {
|
||||
err := c.handle(ctx, msg)
|
||||
if err == nil {
|
||||
return nil // 成功
|
||||
}
|
||||
|
||||
if isRetryable(err) {
|
||||
time.Sleep(time.Duration(attempt+1) * time.Second)
|
||||
continue // 重试
|
||||
}
|
||||
|
||||
// 非可重试错误,直接进 DLQ
|
||||
return c.sendToDLQ(msg, err)
|
||||
}
|
||||
|
||||
// 超过最大重试次数,进 DLQ
|
||||
return c.sendToDLQ(msg, fmt.Errorf("exhausted %d retries", c.maxRetries))
|
||||
}
|
||||
```
|
||||
|
||||
> [!quote] 运维心态
|
||||
>
|
||||
> "一个没有死信队列的消息消费者,就像一辆没有备胎的车——任何一次不可预期的故障都会让你抛锚在路上。"
|
||||
|
||||
## 如何选择分布式事务方案?
|
||||
|
||||
> [!question]- 决策树
|
||||
>
|
||||
> 面对跨服务的数据一致性问题,你不需要每次都重新选型。按照下面的思路逐层判断即可:
|
||||
|
||||
```mermaid
|
||||
flowchart TD
|
||||
START["跨服务写操作"] --> Q1{核心资金或资产}
|
||||
|
||||
Q1 -- 是且要求强一致 --> TCC["TCC"]
|
||||
Q1 -- 否 --> Q2{流程步数超过5}
|
||||
|
||||
Q2 -- 是长流程 --> SAGA["Saga"]
|
||||
Q2 -- 否简单链路 --> Q3{容忍秒级延迟}
|
||||
|
||||
Q3 -- 可以 --> MQ["本地事务+MQ(Outbox)"]
|
||||
Q3 -- 不可以毫秒级 --> TXMSG["RocketMQ事务消息"]
|
||||
|
||||
Q3 -- 无法接受最终一致 --> SEATA["Seata AT模式"]
|
||||
|
||||
TCC --> END["完成选择"]
|
||||
SAGA --> END
|
||||
MQ --> END
|
||||
TXMSG --> END
|
||||
SEATA --> WARNING["过渡方案,后续应迁移至事件驱动架构"]
|
||||
WARNING --> END
|
||||
|
||||
style TCC fill:#e8f5e9,stroke:#4caf50,stroke-width:2px
|
||||
style SAGA fill:#fff3e0,stroke:#ff9800,stroke-width:2px
|
||||
style MQ fill:#e3f2fd,stroke:#1976d2,stroke-width:2px
|
||||
style TXMSG fill:#e3f2fd,stroke:#1976d2,stroke-width:2px
|
||||
style SEATA fill:#fce4ec,stroke:#e91e63,stroke-width:2px
|
||||
```
|
||||
|
||||
### 决策速查矩阵
|
||||
|
||||
| 你的业务特征 | 推荐方案 | 理由 |
|
||||
|-------------|---------|------|
|
||||
| 日订单量百万级,通知类场景 | Outbox + MQ | 性能最好,开发成本最低 |
|
||||
| 支付/清算/账务核心链路 | TCC | 资源占用期间不允许其他事务修改 |
|
||||
| 审批流、物流状态流转(多步) | Saga 编排式 | 流程清晰,补偿机制完备 |
|
||||
| 已有 RocketMQ,对延迟敏感 | 事务消息 | 绕过 Outbox 轮询延迟 |
|
||||
| 遗留系统快速拆分 | Seata AT | 零侵入,快速过渡 |
|
||||
| 只需要"最终一致" | 本地事务 + MQ | 80% 的场景这就是答案 |
|
||||
|
||||
> [!tip] 黄金法则
|
||||
>
|
||||
> **"能用异步事件解决的,不要用两阶段提交;能本地事务 + MQ 解决的,不要上框架。"**
|
||||
>
|
||||
> 复杂度越高,出问题的概率越大。每一层抽象都引入了新的故障点——协调器挂了怎么办?补偿逻辑写错了怎么发现?回滚又失败了谁来做?保持简单是最好的工程纪律。
|
||||
|
||||
## 关联笔记
|
||||
|
||||
- [[03-数据一致性/01-数据库拆分]] — 数据库拆分是分布式事务的前提
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
---
|
||||
tags: [microservice, id-generation, snowflake, uuid, distributed-id]
|
||||
create time: 2026-05-05
|
||||
create time: 2026-05-05 12:00
|
||||
---
|
||||
|
||||
# 分布式 ID 生成
|
||||
|
||||
Reference in New Issue
Block a user