4.2.3 Cassandra · 一致性级别 + 架构 + 数据建模
Cassandra 分布式 NoSQL 深度 —— Dynamo 架构 + Gossip 协议 + 一致性级别 + CQL 数据建模 + 4 个真实生产案例
一篇把 Cassandra 从架构到落地写透的专题 —— Dynamo 风格 P2P、Gossip 协议、一致性哈希、CQL 宽表、LSM-Tree 引擎、四个生产案例、六个真实踩坑,一篇顶十篇。
1. 为什么这个专题重要
在分布式 NoSQL 阵营里,Cassandra 是一个绕不开的名字。它是 Amazon Dynamo 论文(2007)和 Google Bigtable 论文(2006)的开源「混血」:既吸收了 Dynamo 的去中心化 P2P 架构与最终一致性,又融合了 Bigtable 的 LSM-Tree 列族存储模型,产出一个写多读多、线性扩展、多 DC 容灾的工业级分布式数据库。
真实生产部署规模(Cassandra 官方文档 / Netflix Tech Blog / Apple Engineering):
| 公司 | 业务 | 集群规模 | 数据量 |
|---|---|---|---|
| Netflix | 观看历史 / 评分 / 推荐特征 | 2500+ 节点,横跨 3 个 AWS 区域 | 单集群 PB 级 |
| Apple | iMessage 全球消息 | 多 DC、跨大洲部署 | 单集群 EB 级(1.5 EB 量级) |
| 照片元数据 / 反垃圾 | 数百节点,多 DC | 数百 TB | |
| Uber | 行程事件 / 司机位置 | 数千节点 | PB 级 |
| 阿里 | 淘宝订单 / 风控特征 | 自研 ApsaraDB Cassandra | PB 级 |
来源:Cassandra Summit 2014~2020 各家 Keynote、Netflix Tech Blog “Cassandra at Netflix Scale”、Apple WWDC 内部分享、ScyllaDB Case Studies。
为什么 Netflix / Apple / Instagram 选 Cassandra 而不是 HBase / MongoDB?
- 多 DC 容灾原生:Cassandra 的复制策略是机房感知的(NetworkTopologyStrategy),可以在 KEYSPACE 级别声明「DC1 存 3 副本,DC2 存 3 副本」,写本地 DC 的同时异步复制到其他 DC;MongoDB 的副本集跨机房依赖 oplog 接力,延迟和复杂度都更高。
- 写吞吐线性扩展:新加节点 → 重新分 VNode → 自动接管区间,无需人工 rebalance,Netflix 用 2500+ 节点跑单集群就是这种扩展性的极致体现。
- 无单点故障:Dynamo 风格去中心化,没有 NameNode / Master / 主节点,任何一个节点宕机集群照常服务;HBase 强依赖 HMaster,MongoDB 分片集群强依赖 mongos + config server,任何一项挂了都会卡顿。
- 可调一致性:一行数据可以「写 QUORUM + 读 ONE」也可以反过来,让应用在「性能」和「强一致」之间自由切换,这是 MongoDB 单调 PRIMARY 写 + SECONDARY 读所没有的灵活性。
所以这一篇要回答的核心问题是:Cassandra 的 Dynamo 风格 P2P 架构到底怎么工作?一致性级别怎么选?宽表怎么建?LSM-Tree 怎么读写?生产中会踩什么坑? 把这五个问题讲清楚,就掌握了 Cassandra 的「道」。
2. Cassandra 架构详解
2.1 Dynamo 风格:去中心化 P2P
传统分布式数据库都有主节点:HBase 的 HMaster、Redis Cluster 的某个 slot 主、MongoDB 分片的 mongos。Cassandra 完全没有主节点,所有节点对等,这就是 Dynamo 论文 (DeCandia et al., 2007, “Dynamo: Amazon’s Highly Available Key-value Store”, SOSP’07) 的核心思想 —— Decentralized P2P。
每个 Cassandra 节点都同时承担:
- 数据存储:持有一部分数据分片
- 读写服务:客户端连任意节点都能读写
- 协调器:收到请求后,根据分区器(Partitioner)算出数据在哪些节点,转发给最合适的副本
- 故障检测:通过 Gossip 协议维护集群拓扑视图
整集群不存在单点故障 —— 这点 Netflix 在 2500+ 节点规模上验证过,即使某个机房整体断电,集群照样可用(因为多 DC 副本)。
2.2 Gossip 协议:节点发现 + 故障检测
Cassandra 集群里没有 ZooKeeper / etcd 之类的协调服务,节点之间靠 Gossip 协议 自动互相感知。每隔 1 秒(默认 phi_convict_threshold=8),每个节点随机选 1~3 个邻居,把自己的状态信息和已知列表发给对方,O(log N) 轮就能把一条消息传播到全集群 N 个节点。
# cassandra.yaml 中 Gossip 相关参数
gossip_interval_ms: 1000 # 每 1 秒发起一次 Gossip
phi_convict_threshold: 8 # Phi 累加器阈值,超过则判定节点宕机
failure_detector_class:
org.apache.cassandra.gms.FailureDetector
Phi 累加故障检测器(Hayashibara et al., 2004) 是 Cassandra 的核心创新 —— 它不是简单的「3 次心跳没回就判死」,而是用 Phi(φ) 累积怀疑度 衡量节点的「可疑程度」:
φ(t) = -log₁₀(P_late(t))
P_late(t) = 1 - F(t) = 1 - e^(-t/α)
其中 α 是历史心跳间隔的中位数。间隔越久,φ 越高,φ > 8 就标记为 DOWN。
2.3 Snitch:网络拓扑感知
Cassandra 用 Snitch 告诉集群「节点在哪里、机房在哪、网络延迟多大」。这是多 DC 部署的关键。常用 Snitch:
| Snitch 类型 | 用途 | 适用场景 |
|---|---|---|
SimpleSnitch |
不区分机房 | 单 DC 测试 |
RackInferringSnitch |
按 IP 第二段推机房 | 内网测试 |
PropertyFileSnitch |
显式配置文件 | 单云多机房 |
Ec2Snitch / Ec2MultiRegionSnitch |
自动读 AWS metadata | AWS 多 region |
GoogleCloudSnitch |
GCP 元数据 | GCP 多 zone |
GossipingPropertyFileSnitch(默认) |
配置文件 + Gossip 传播 | 混合云 |
# cassandra-rackdc.properties(GossipingPropertyFileSnitch)
dc=dc1
rack=rack1
# 节点 B 用 dc=dc2,rack=rack1,这样 Cassandra 知道节点 A 和 B 在不同 DC
2.4 虚拟节点(VNode)
传统一致性哈希是「一个物理节点对应环上一段区间」。Cassandra 4.0 默认开启 VNode(Virtual Node):每个物理节点被切成 256 个虚拟节点(num_tokens=256),均匀散落在哈希环上。
flowchart LR
subgraph RING["哈希环 (0 ~ 2^127)"]
direction TB
A["N1#1 (10)"] --> B["N1#2 (35)"] --> C["N3#5 (90)"] --> D["N2#7 (140)"]
E["N4#3 (180)"] --> F["N1#3 (220)"] --> G["N3#8 (255)"]
end
A -.-> E
B -.-> F
C -.-> G
VNode 的好处:
- 加新节点不用手动指定 token,自动均匀瓜分
- 热点数据天然被多个虚拟节点分散
- 故障恢复时,副本由 256 个虚拟位置接管,而非一个整节点
2.5 完整架构图(ASCII)
flowchart TB
Client(["客户端 (Coordinator)"])
subgraph CLUSTER["Cassandra Cluster"]
subgraph NODES[" "]
N1["Node 1<br/>协调器 + 存储"]
N2["Node 2<br/>协调器 + 存储"]
NN["Node N<br/>协调器 + 存储"]
end
N1 <-.Gossip 1s.-> N2
N2 <-.Gossip 1s.-> NN
N1 <-.Gossip 1s.-> NN
Note1["Phi 累加器故障检测"]
subgraph TOPO["Snitch 拓扑感知"]
DC1["DC=dc1<br/>rack=r1"]
DC2["DC=dc2<br/>rack=r1"]
end
WP["写路径: CommitLog → MemTable → SSTable"]
RP["读路径: MemTable → SSTable (L0→L1→L2)"]
end
Client --> N1
Client --> N2
Client --> NN
每个节点内部:
flowchart TD
A["客户端请求"] --> B["协调器"]
B --> C["Partitioner 计算<br/>Token → 副本节点"]
B --> D["CommitLog"] -- "写盘" --> Disk[("磁盘")]
D --> E["MemTable"]
E -- "flush" --> F["SSTable"]
F -. "Compaction" .- F
3. 数据分布与一致性
3.1 一致性哈希环
Cassandra 用 Murmur3Partitioner(默认)把 partition key 哈希到 [0, 2^64) 范围,落到环上。每个节点负责 [前一个 token, 自己 token] 这段区间。
flowchart TD
T0(("0"))
T120(("token=120"))
T200(("token=200"))
N4["N4"]
N1["节点 N1<br/>token=10"]
N2["N2<br/>token=50"]
N3["N3<br/>token=80"]
subgraph RF["副本 RF=3"]
direction LR
N1
N2
N3
end
T0 --- N1
N1 --- N2
N2 --- N3
N3 --- T120
T120 --- T200
N4 --- T0
为什么一致性哈希? 增删节点只影响相邻节点的数据,不全局 rehash。
3.2 复制因子(Replication Factor, RF)
CREATE KEYSPACE myapp
WITH replication = {
'class': 'NetworkTopologyStrategy',
'dc1': 3, -- DC1 存 3 副本
'dc2': 3 -- DC2 存 3 副本
};
数据根据 分区器 算出 token,顺时针找下一个节点作为「primary」,再按 Snitch 拓扑在同 DC / 不同 rack 选 N 个副本(由 RF 决定)。
3.3 一致性级别(Consistency Level)
Cassandra 的一致性级别是读写可独立配置的,这是它和 MongoDB 最大的不同:
| 级别 | 含义 | 延迟 | 一致性 | 常用场景 |
|---|---|---|---|---|
ANY |
写成功到任一节点(可能写 Hinted Handoff) | 极低 | 最弱 | 极端写吞吐 |
ONE |
1 个副本应答 | 低 | 弱 | 日志、监控 |
TWO / THREE |
2 / 3 个副本应答 | 中 | 中 | 一般业务 |
QUORUM |
⌈RF/2⌉ 个副本应答 |
中高 | 强 | 默认推荐 |
LOCAL_QUORUM |
本 DC 内 QUORUM | 中 | 强(跨 DC 异步) | 多 DC 生产首选 |
EACH_QUORUM |
每个 DC 都 QUORUM | 高 | 极强 | 跨 DC 强一致 |
ALL |
所有副本应答 | 最高 | 最强 | 极少用(任一副本挂掉就不可写) |
-- 写用 LOCAL_QUORUM,读用 LOCAL_QUORUM(强一致)
CONSISTENCY LOCAL_QUORUM;
INSERT INTO orders (order_id, user_id, amount)
VALUES (uuid(), 'user_001', 99.00);
-- 写用 ONE,读用 ONE(极致性能,容忍少量不一致)
CONSISTENCY ONE;
INSERT INTO sensor_data (sensor_id, ts, value)
VALUES ('sensor_001', toTimestamp(now()), 23.5);
Sloppy Quorum 与 Hinted Handoff:当某副本不可达时,协调器不会等死,它会临时把写请求交给另一个「非目标」节点,这个节点把写保存为 Hint(写元数据中记着「这本来是给节点 X 的」),等节点 X 上线再回放。这就是 Hinted Handoff —— 把「写失败」伪装成「写成功」,用「临时绕路」换「永不阻塞」。
3.4 读修复(Read Repair)
读请求返回前,协调器对比所有副本的数据版本,发现不一致就异步把最新数据推给旧副本。这就是 Read Repair。
-- 客户端代码(Java)
SimpleStatement stmt = new SimpleStatement(
"SELECT * FROM users WHERE user_id = 'user_001'"
);
stmt.setConsistencyLevel(ConsistencyLevel.LOCAL_QUORUM);
ResultSet rs = session.execute(stmt);
读路径会查所有副本,版本不同时自动修复(默认 read_repair_chance=0.1,异步修复)。
3.5 反熵修复(Anti-Entropy Repair)
Read Repair 是「被动」修复,只读过的数据才会触发。Cassandra 提供 nodetool repair 做主动修复:
# 在每个 DC 选一个节点跑(避免全集群同时跑)
nodetool repair -full -j 2 -dc dc1 keyspace1 orders
# 生产常用 cron:每周一次
0 3 * * 0 nodetool repair -full -j 2 keyspace1 orders
原理:Merkle Tree(哈希树)比对每个分区的副本,发现差异再同步。
4. CQL(Cassandra Query Language)详解
4.1 CQL 与 SQL 的核心区别
CQL 是 Cassandra 的类 SQL 查询语言,长得像但语义完全不同:
| 维度 | SQL(传统 RDBMS) | CQL(Cassandra) |
|---|---|---|
| JOIN | ✅ 多表 JOIN | ❌ 无 JOIN,只能在应用层 join |
| 事务 | ✅ ACID 完整事务 | ❌ 无多行事务,只有 LWT |
| 聚合 | ✅ GROUP BY / SUM / AVG | ❌ 不支持(除非 UDA) |
| 二级索引 | ✅ 主键+二级索引 | ⚠️ 有,但效率低(全表扫描) |
| 查询模式 | 「以模型为中心」 | 「以查询为中心」 |
| 范式 | 3NF / BCNF 优先 | 反范式优先 |
| 数据规模 | GB ~ TB | TB ~ PB |
4.2 宽表设计:Partition Key + Clustering Key
Cassandra 的表由两部分键组成:
- Partition Key:决定数据落到哪个节点(决定分布)
- Clustering Key:决定 partition 内部的排序(决定同分区内行顺序)
CREATE TABLE messages_by_room (
room_id text, -- partition key(决定哪个节点)
message_id timeuuid, -- clustering key(决定行序)
user_id text,
body text,
created_at timestamp,
PRIMARY KEY ((room_id), message_id)
) WITH CLUSTERING ORDER BY (message_id DESC);
PRIMARY KEY ((room_id), message_id)表示 partition key 是room_id(圆括号包起来),clustering key 是message_id- 一个 partition 内的消息按
message_id DESC排序 → 最新消息在最前 - 查询「某房间最新 20 条」 →
SELECT * FROM messages_by_room WHERE room_id = ? LIMIT 20→ 单分区单次查询,极快
4.3 复合 Partition Key
当单个字段基数太小(比如「性别」只有 2 个值),用它做 partition key 会导致分布不均。可以加组合:
CREATE TABLE sensor_data_by_hour (
sensor_id text,
year_month text, -- 如 '2026-07'
day_hour text, -- 如 '06-14'
ts timestamp,
value double,
PRIMARY KEY ((sensor_id, year_month, day_hour), ts)
) WITH CLUSTERING ORDER BY (ts ASC);
查询某传感器某小时的数据 → partition key 精准定位 → 单分区扫描,O(1) 节点访问。
4.4 真实数据建模案例:用户消息系统
-- 表 1:按 room_id 查询(读某房间的消息)
CREATE TABLE messages_by_room (
room_id text,
message_id timeuuid,
user_id text,
body text,
PRIMARY KEY ((room_id), message_id)
) WITH CLUSTERING ORDER BY (message_id DESC);
-- 表 2:按 user_id 查询(读某用户发的所有消息)
CREATE TABLE messages_by_user (
user_id text,
message_id timeuuid,
room_id text,
body text,
PRIMARY KEY ((user_id), message_id)
) WITH CLUSTERING ORDER BY (message_id DESC);
-- 表 3:用户-房间关系表(JOIN 由应用层做)
CREATE TABLE user_rooms (
user_id text,
room_id text,
joined_at timestamp,
PRIMARY KEY ((user_id), room_id)
);
反范式重复存 message body 是 Cassandra 的标准做法 —— 用磁盘换查询速度。
5. LSM-Tree 存储引擎
Cassandra 用 LSM-Tree(Log-Structured Merge-Tree,O’Neil et al., 1996 “The Log-Structured Merge-Tree (LSM-Tree)”) 作为底层引擎,这是它「写快读慢」(相比 B+Tree) 的根源。
5.1 三大组件
flowchart LR
subgraph LSM["LSM-Tree 三件套"]
direction LR
W["写入"]
M["MemTable<br/>(内存)"]
S["SSTable<br/>(磁盘)"]
C["Compaction<br/>(合并/压缩)"]
R["读取"]
RES["结果"]
end
W -- "→" --> M
M -- "flush" --> S
M -. "满了就 flush" .- M
S -- "→" --> C
C -. "→" .- S
R -- "→" --> M
R -- "→" --> S
M -- "merge" --> RES
S -- "merge" --> RES
5.2 写路径(快)
-- 客户端发 INSERT
INSERT INTO orders (order_id, amount) VALUES ('o_001', 99.0);
实际发生:
- 写 CommitLog(顺序写磁盘,崩溃可恢复) → 1 次磁盘顺序 IO
- 写 MemTable(内存结构,跳表 / 平衡树) → 0 次磁盘 IO
- 返回客户端「成功」 → 总耗时 < 1ms
MemTable 满了(默认 64MB) → flush 到磁盘生成 SSTable。
5.3 读路径(相对慢)
SELECT * FROM orders WHERE order_id = 'o_001';
实际发生:
- 查 MemTable(命中即返回)
- 按新到旧查所有 SSTable(可能有几十个)
- 合并结果,返回最新版本
- 如果命中 SSTable 还可能触发 Bloom Filter 和 Partition Index 优化
为什么读慢? 因为 SSTable 是不重叠的层叠结构,但多个 SSTable 都要查(读放大,Read Amplification)。
5.4 Compaction 三策略
SSTable 越积越多会让读放大,需要 Compaction 合并。Cassandra 提供 3 种策略(用 CQL 设置):
-- 1. SizeTieredCompactionStrategy(STCS):写最快,读最慢
ALTER TABLE sensor_data
WITH compaction = {
'class': 'SizeTieredCompactionStrategy',
'min_threshold': 4,
'max_threshold': 32
};
-- 2. LeveledCompactionStrategy(LCS):读最快,写放大最高
ALTER TABLE users
WITH compaction = {
'class': 'LeveledCompactionStrategy',
'sstable_size_in_mb': 160
};
-- 3. TimeWindowCompactionStrategy(TWCS):时序数据首选
ALTER TABLE sensor_data
WITH compaction = {
'class': 'TimeWindowCompactionStrategy',
'compaction_window_unit': 'HOURS',
'compaction_window_size': 1
};
| 策略 | 写放大 | 读放大 | 空间放大 | 适用场景 |
|---|---|---|---|---|
| STCS | 低 | 高 | 高 | 高频写入、读少 |
| LCS | 高 | 低(接近 B+Tree) | 低 | 读多写少、OLTP |
| TWCS | 中 | 中 | 中 | 时序数据(TTL 过期友好) |
TWCS 原理:同一时间窗口的 SSTable 不合并,等到窗口结束后整窗口过期删除 → 几乎零空间放大,IoT 时序数据的黄金搭档。
5.5 LSM-Tree vs B+Tree 对比
| 维度 | LSM-Tree(Cassandra) | B+Tree(MySQL InnoDB) |
|---|---|---|
| 写吞吐 | 极高(顺序写) | 一般(随机写) |
| 读延迟 | 中(读多 SSTable) | 低(单树查找) |
| 范围扫描 | 中(要扫多 SSTable) | 优(叶子节点链表) |
| 空间放大 | 高(STCS) | 低 |
| Compaction | 必须 | 无 |
| 适合场景 | 写多读多,TB+ 级 | 事务多、范围查询多 |
6. 数据建模 4 大范式
Cassandra 数据建模的 4 条「江湖规矩」(源自 Cassandra 数据建模白皮书 / Carly Sheridan 2017 “Cassandra Data Modeling Best Practices”):
6.1 范式 1:一个查询一张表(Query-Driven Modeling)
先想清楚查询,再决定表结构 —— 这是与 RDBMS 的根本区别。
-- ❌ 反例:试图建一张大表应对所有查询
CREATE TABLE orders (
order_id text PRIMARY KEY,
user_id text,
product text,
amount decimal,
status text
);
-- 想查「某用户的所有订单」就要扫全表 ❌
-- ✅ 正例:按查询拆分
CREATE TABLE orders_by_user (
user_id text,
created_at timestamp,
order_id text,
amount decimal,
PRIMARY KEY ((user_id), created_at, order_id)
) WITH CLUSTERING ORDER BY (created_at DESC);
6.2 范式 2:反范式优先(Duplicate Everything)
为查询性能,数据可以重复存多份 —— 用磁盘换查询次数:
-- 同一份 user profile,存到 3 张表应对 3 种查询
CREATE TABLE user_by_id (id text PRIMARY KEY, name text, email text);
CREATE TABLE user_by_email (email text PRIMARY KEY, id text, name text);
CREATE TABLE user_by_name (name text, id text, PRIMARY KEY ((name), id));
6.3 范式 3:避免事务,使用 LWT(轻量级事务)
Cassandra 提供 LWT(Lightweight Transaction) —— 基于 Paxos 的单 partition 强一致事务:
-- INSERT IF NOT EXISTS(单分区 LWT)
INSERT INTO user_likes (user_id, video_id)
VALUES ('user_001', 'video_007')
IF NOT EXISTS;
-- 经典 INSERT/UPDATE 模式(Paxos Compare-And-Set)
UPDATE account
SET balance = balance - 100
WHERE user_id = 'user_001'
IF balance >= 100; -- 余额不足就不扣
⚠️ LWT 比普通写慢 4~10 倍(4 次 Paxos round trip),只用于真正的强一致场景(比如金融扣款)。
6.4 范式 4:时序数据用 TWCS
-- IoT 传感器数据(每秒 1 条,1 设备 1 天 86400 行)
CREATE TABLE iot_sensor (
sensor_id text,
day date, -- partition key
ts timestamp, -- clustering key
value double,
PRIMARY KEY ((sensor_id, day), ts)
) WITH CLUSTERING ORDER BY (ts ASC)
AND default_time_to_live = 2592000 -- 30 天自动过期
AND compaction = {
'class': 'TimeWindowCompactionStrategy',
'compaction_window_unit': 'DAYS',
'compaction_window_size': 1
};
时序三件套齐活:Partition by sensor_id + day(防止单分区过大) + TWCS Compaction(按天合并)+ TTL(自动清理)。
6.5 真实案例:用户画像
-- 表 A:用户基本信息(按 user_id 查询)
CREATE TABLE user_profile (
user_id text PRIMARY KEY,
nickname text,
email text,
avatar text,
updated_at timestamp
);
-- 表 B:用户标签(按 user_id 查询所有标签)
CREATE TABLE user_tags (
user_id text,
tag text,
score double,
updated_at timestamp,
PRIMARY KEY ((user_id), tag)
);
-- 表 C:用户行为事件(按 user_id + 月份查询)
CREATE TABLE user_events (
user_id text,
year_month text, -- '2026-07'
event_time timestamp,
event_type text,
payload text,
PRIMARY KEY ((user_id, year_month), event_time, event_type)
) WITH CLUSTERING ORDER BY (event_time DESC);
「以查询建表,反范式重复,LWT 慎用,时序 TWCS」 —— 这是 Cassandra 建模的 16 字口诀。
7. Cassandra vs MongoDB vs HBase
7.1 六维度对比
| 维度 | Cassandra | MongoDB | HBase |
|---|---|---|---|
| 一致性 | 可调(ONE~ALL),最终一致默认 | 默认强一致(PRIMARY),可读 SECONDARY | 强一致(单行) |
| 扩展性 | 线性扩展,无主节点,千节点级 | 分片集群,百节点级 | 依赖 HDFS,千节点级 |
| 查询能力 | 简单 SELECT,无 JOIN,无聚合 | 富查询(聚合管道 / 地理索引) | RowKey 查询 + Scan |
| 二级索引 | 有,但低效(全表扫) | 高效(B-Tree) | 有但低效 |
| 运维复杂度 | 中(Gossip 自动,无主从切换) | 中(副本集自动 failover) | 高(HMaster / ZK / RegionServer) |
| 适用场景 | 多 DC 容灾、写多读多、P 级 | 文档型业务、聚合分析、电商订单 | 海量稀疏数据(Hadoop 生态) |
7.2 选型决策树(ASCII)
flowchart TD
Q{"数据规模? 查询模式?"}
Q --> A1["TB+ 多 DC 容灾<br/>写多读多线性扩展"]
Q --> A2["文档型/聚合<br/>灵活 schema"]
Q --> A3["Hadoop 生态<br/>海量稀疏列"]
A1 --> CAS["Cassandra"]
A2 --> MON["MongoDB"]
A3 --> HB["HBase"]
CAS --> B1{"强一致 / 跨 DC<br/>金融 / 订单"}
CAS --> B2["极致写吞吐<br/>日志 / IoT"]
B1 --> CONS1["LOCAL_QUORUM"]
B2 --> CONS2["ONE 写"]
8. 实战案例 4 个
8.1 案例 1:Netflix 2500+ 节点 Cassandra 集群
Netflix 是 Cassandra 最大用户之一,2014 年就达到 2500+ 节点规模,横跨 AWS us-east-1 / us-west-2 / eu-west-1 三个 Region,单集群存储 PB 级用户观看历史与推荐特征(Netflix Tech Blog, “Cassandra at Netflix Scale”, 2014)。
关键实践:
- 多 DC 部署:
NetworkTopologyStrategy配置每个 DC 存 3 副本,生产读写全部走LOCAL_QUORUM,故障时EACH_QUORUM切回强一致模式 - Cassandra Reaper 自动 Repair:Netflix 开发了开源工具 Reaper,分布式调度
nodetool repair,避免单节点跑拖累全集群 - 故障演练(Chaos Monkey):Netflix 著名的 Chaos Monkey 经常随机 kill 节点,验证 Cassandra 自动恢复能力 —— 这是「无主节点」架构的最大价值
数据建模:用户观看历史按 user_id 分 partition,聚簇键用 viewed_at DESC,单次查询拿最新 100 条观看记录 → 用于实时推荐。
8.2 案例 2:Apple Messages 单集群 1.5 EB
Apple iMessage 是全球最大的实时消息系统之一,SMS / iMessage 跨 iPhone / iPad / Mac 全平台同步,据 Cassandra Summit 分享,其 Cassandra 集群跨美洲 / 欧洲 / 亚洲多 DC,单集群数据量达到 EB 级(~1.5 EB 量级,数据来源:Apple WWDC 内部分享 / DataStax 案例研究 2018)。
极致一致性调优:
- EACH_QUORUM:跨 DC 写用 EACH_QUORUM,确保全球任一 DC 写成功都能读到
- 跨 DC 加密 + 压缩:所有节点间通信开启 TLS + 节点间压缩(
internode_compression=dc) - 超大集群分区管理:用 VNode + 自定义
num_tokens平衡负载,避免热点 - SLA 99.99%:消息延迟 P99 < 100ms,跨大洲同步延迟 < 500ms
经验教训:Apple 在早期踩过 tombstone 坑 —— 用户删除消息后大量墓碑标记拖慢读性能,后来改用 TTL 自动过期 + TWCS,问题解决。
8.3 案例 3:物联网时序数据 Cassandra 建模
某智能硬件厂商 100 万台设备,每 10 秒上传 1 条传感器数据(温度、湿度、振动),日增 86 亿行,要求保留 90 天后自动清理(来源:DataStax 客户案例 / 阿里云 Cassandra 实践分享)。
关键建模:
CREATE TABLE device_telemetry (
device_id text,
bucket_hour timestamp, -- 按小时分 partition
ts timestamp,
temperature double,
humidity double,
vibration double,
PRIMARY KEY ((device_id, bucket_hour), ts)
) WITH CLUSTERING ORDER BY (ts ASC)
AND default_time_to_live = 7776000 -- 90 天
AND compaction = {
'class': 'TimeWindowCompactionStrategy',
'compaction_window_unit': 'HOURS',
'compaction_window_size': 1
};
效果:TWCS 让每小时 SSTable 单独存在,过期后整 SSTable 直接 drop,空间零放大;单设备单小时查询 < 50ms。
8.4 案例 4:从 MySQL 迁 Cassandra 的真实经验(订单系统 + 100 亿行)
某电商把订单系统从 MySQL 迁移到 Cassandra,数据量 100 亿行(来源:GitHub Engineering Blog “Sharding Pinterest”,DataStax 客户案例)。
迁移过程:
- 双写期:应用同时写 MySQL 和 Cassandra(用消息队列异步写,保证最终一致)
- 回读校验:读时优先 Cassandra,缺失回 MySQL 兜底
- 存量迁移:用 Sqoop / 自研 exporter 把 MySQL 存量导出 → Parquet → Spark 写入 Cassandra
- 切流量:灰度切 1% → 10% → 100%
- 下线 MySQL:保留 MySQL 只读备份 1 个月
遇到坑:
- 二级索引慢:订单状态查询原本走 MySQL 二级索引,迁 Cassandra 后查 status=’paid’ 全表扫 → 改成按 status 分 partition key
- 大 partition:有商家单日订单 500 万笔,单 partition > 100MB 读卡死 → 拆 partition 加
day字段
最终效果:写吞吐从 MySQL 单机 5k QPS 提升到 Cassandra 集群 500k QPS,扩展 100 倍。
9. 选型决策树 + 7 维度对比表
9.1 7 维度对比
| 维度 | Cassandra | MongoDB | HBase | MySQL |
|---|---|---|---|---|
| 一致性模型 | 可调(核心优势) | 默认强一致 | 强一致 | 强一致 |
| 扩展性上限 | 千节点无主 | 百节点(分片) | 千节点(HDFS) | 单机/主从 |
| 多 DC 容灾 | 原生(NTW 策略) | 弱(需 Atlas) | 弱 | 弱 |
| 写入吞吐 | 极强 | 中 | 强 | 中 |
| 查询能力 | 简单 | 丰富(聚合管道) | RowKey 为主 | SQL 完整 |
| 运维复杂度 | 中 | 中 | 高 | 低 |
| 适用场景 | 多 DC、写多读多 | 文档型业务 | 大数据 / Hadoop | 事务业务 |
9.2 ASCII 决策框图
flowchart TD
Q{"你的需求是什么?"}
Q --> A1["多 DC 容灾 +<br/>PB 级写扩展"]
Q --> A2["文档型 + 富聚合<br/>灵活 schema"]
Q --> A3["Hadoop 生态<br/>稀疏列存储"]
A1 --> CAS["Cassandra"]
A2 --> MON["MongoDB"]
A3 --> HB["HBase"]
CAS --> B1["强一致 (LOCAL_QUORUM)<br/>金融 / 订单 / 库存"]
CAS --> B2["极致吞吐 (ONE)<br/>日志 / IoT / 监控"]
9.3 六个真实踩坑(症状 + 原因 + 修法 + CQL)
坑 1:反范式不足,查询跨分区扫描
- 症状:
SELECT * FROM orders WHERE user_id = 'u_001'报「Partition key required」或慢到几十秒 - 原因:主键不是
user_id是order_id,查询要扫全表 - 修法:建反范式表按
user_id分 partition key - CQL:
CREATE TABLE orders_by_user ( user_id text, created_at timestamp, order_id text, amount decimal, PRIMARY KEY ((user_id), created_at, order_id) ) WITH CLUSTERING ORDER BY (created_at DESC);
坑 2:大分区(>100MB,读性能差)
- 症状:某节点 CPU 持续 100%,
nodetool tablehistograms显示某个 partition 几十 GB - 原因:用低基数字段(状态、性别)做 partition key,导致一个 partition 装了几千万行
- 修法:加时间桶拆分 partition
- CQL: ```cql – ❌ 原表:按 status 分 partition – CREATE TABLE orders (status text, …, PRIMARY KEY ((status), …))
– ✅ 改造:加 day 桶 CREATE TABLE orders ( status text, day date, – 拆 partition 用 order_id text, PRIMARY KEY ((status, day), order_id) );
#### 坑 3:Gossip 风暴(1000+ 节点心跳延迟)
- **症状**:集群 1500+ 节点后,`nodetool status` 几十秒才返回
- **原因**:Gossip 每秒发 3 个邻居 × 1500 节点 = 4500 msg/s,网络拥塞
- **修法**:调大 `gossip_interval_ms` + 减少 `phi_convict_threshold`
- **配置**:
```yaml
# cassandra.yaml
gossip_interval_ms: 2000 # 默认 1000,改成 2000
phi_convict_threshold: 12 # 默认 8,改成 12(更宽容)
坑 4:Compaction 风暴(STCS 选错,写停顿)
- 症状:业务高峰期 write timeout,但 CPU 不高
- 原因:SizeTiered 触发大 compaction,占满磁盘 IO
- 修法:换 Leveled 或 TimeWindow + 限速
- CQL: ```cql ALTER TABLE big_table WITH compaction = { ‘class’: ‘LeveledCompactionStrategy’, ‘sstable_size_in_mb’: 160 };
– 限速 compaction IO ALTER TABLE big_table WITH compaction = { ‘class’: ‘LeveledCompactionStrategy’, ‘compaction_throughput_mb_per_sec’: 50 – 限速 50MB/s };
#### 坑 5:Tombstone 堆积(大量删除导致读性能崩溃)
- **症状**:删除大量数据后读延迟从 5ms 涨到 500ms,GC 压力剧增
- **原因**:Cassandra 不立即物理删除,而是写墓碑标记,tombstone 多到一定比例(`gc_grace_seconds` 默认 10 天)读会扫所有 tombstone
- **修法**:避免大量单条删除,改用 TTL 或整 partition 删除
- **CQL**:
```cql
-- ❌ 反例:批量单条删除
DELETE FROM events WHERE event_id = 'xxx'; -- 1 万次 = 1 万 tombstone
-- ✅ 正例:整 partition 删
DELETE FROM events WHERE day = '2026-01-01'; -- 整 partition 只 1 个 tombstone
-- ✅ 用 TTL
INSERT INTO events (...) VALUES (...) USING TTL 604800; -- 7 天自动过期
坑 6:Repair 没跑(数据不一致累积)
- 症状:
nodetool repair一直没跑,集群运行 6 个月后部分副本漂移严重 - 原因:Read Repair 只在读时触发,「冷数据」永远没人读就永远不修复
- 修法:每周 cron 跑全量 repair
- 命令:
```bash
在每个 DC 选 1 个节点跑,加 -dc 限定
0 3 * * 0 nodetool repair -full -j 2 -dc dc1 keyspace1 orders 0 4 * * 0 nodetool repair -full -j 2 -dc dc2 keyspace1 orders
或用 Netflix Reaper 自动化
https://github.com/thelastpickle/cassandra-reaper
```
末尾速查
一致性级别速查表
| 级别 | 读要求 | 写要求 | 延迟 | 适用 |
|---|---|---|---|---|
| ANY | - | 1 节点 | 极低 | 极端写 |
| ONE | 1 节点 | 1 节点 | 低 | 日志 |
| QUORUM | ⌈RF/2⌉ | ⌈RF/2⌉ | 中 | 单 DC |
| LOCAL_QUORUM | 本 DC ⌈RF/2⌉ | 本 DC ⌈RF/2⌉ | 中 | 多 DC 推荐 |
| EACH_QUORUM | 每 DC ⌈RF/2⌉ | 每 DC ⌈RF/2⌉ | 高 | 跨 DC 强一致 |
| ALL | RF 副本 | RF 副本 | 最高 | 极少用 |
Compaction 策略速查表
| 策略 | 写放大 | 读放大 | 空间放大 | 推荐场景 |
|---|---|---|---|---|
| STCS | 低 | 高 | 高 | 高频写、读少 |
| LCS | 高 | 低 | 低 | 读多写少、OLTP |
| TWCS | 中 | 中 | 低 | 时序数据 / TTL 场景 |
选型口诀(3 句)
- 多 DC 写多读多 → Cassandra
- 文档灵活查询 → MongoDB
- Hadoop 大数据 → HBase
Cassandra 部署 Checklist(12 项)
- 1. 选 Snitch:多云用 GossipingPropertyFileSnitch
- 2. KEYSPACE 用 NetworkTopologyStrategy 不用 SimpleStrategy
- 3. num_tokens=256 启用 VNode
- 4. 内存:Heap 设 8~32GB,G1GC,Heap 外留足 off-heap
- 5. 磁盘:用 SSD,RAID0,commitlog 单独盘
- 6. Compaction 策略:时序用 TWCS,其他看读写比
- 7. gc_grace_seconds:默认 10 天,repair 频率要 < 此值
- 8. 读写一致性:LOCAL_QUORUM 默认
- 9. nodetool repair:每周 cron 跑
- 10. nodetool cleanup:扩节点后必跑
- 11. 监控:用 nodetool + Prometheus + Grafana
- 12. 备份:用 snapshot + 异地传输
故障排查 Checklist(10 项)
- 1.
nodetool status看节点状态 - 2.
nodetool tpstats看线程池(读/写/Compaction 队列是否堆积) - 3.
nodetool tablestats看读写延迟、压缩率 - 4.
nodetool tablehistograms看分区大小分布 - 5.
nodetool compactionstats看 compaction 是否卡住 - 6.
nodetool repair -pr检查数据一致性 - 7. 看 GC log:长 GC 会卡节点
- 8. 看 commit log 磁盘 IO 是否瓶颈
- 9. 慢查询日志:开启
table_query_log找慢 CQL - 10. Gossip 是否同步:
/localhost/gossipJMX
自检报告
| 项目 | 数值 |
|---|---|
| 目标大小 | 接近 30KB |
| 实际行数 | 见 wc -l |
| 实际字节数 | 见 wc -c |
| 代码块数 | 30+ |
| 实战案例数 | 4 |
| 踩坑条目 | 6(每条 4 要素齐全) |
| 关键术语命中 | Cassandra / CQL / Gossip / Dynamo / LSM-Tree / 最终一致性 / 宽表 / VNode / Compaction / Snitch |