第一章:Cassandra 概述
1.1 什么是 Cassandra
| 概念名称 | 说明 | 注意事项 |
|---|
| Apache Cassandra | 一个开源的、分布式的、高可用的 NoSQL 列式数据库,最初由 Facebook 开发,后捐赠给 Apache 软件基金会。 | 不是传统的关系型数据库,不支持 JOIN 和复杂事务。 |
| 列式存储 | 数据按列族(Column Family)组织,物理上以列为单位存储,适合大规模写入和时间序列场景。 | 与 HBase 的列式模型类似,但数据模型更灵活。 |
| 无中心架构 | 所有节点对等(peer-to-peer),无主从之分,任意节点均可处理读写请求。 | 避免了单点故障,提升系统容错能力。 |
| 线性可扩展 | 可通过添加节点线性提升吞吐量和存储容量,支持 PB 级数据规模。 | 扩容过程无需停机,数据自动再平衡。 |
1.2 Cassandra 的核心特性
| 特性名称 | 说明 | 注意事项 |
|---|
| 高可用性(High Availability) | 通过多副本机制和无中心架构实现 24/7 服务,即使部分节点宕机仍可正常响应。 | 副本数(replication factor)需合理配置,通常 ≥3。 |
| 最终一致性(Eventual Consistency) | 默认采用最终一致性模型,允许短暂的数据不一致以换取高写入性能。 | 可通过设置一致性级别(如 QUORUM)增强一致性。 |
| 容错与去中心化 | 无主节点,所有节点平等;使用 Gossip 协议传播集群状态。 | 网络分区(split-brain)需通过 NTP 同步时间避免。 |
| 线性可扩展性 | 添加新节点即可水平扩展,系统自动分配 token 范围并迁移数据。 | 需确保硬件资源(CPU、磁盘、网络)均衡。 |
| 多数据中心支持 | 支持跨多个物理数据中心部署,复制策略可按机架或区域定制。 | 使用 NetworkTopologyStrategy 实现地理冗余。 |
| 高写入吞吐 | 写操作先写入内存(Memtable)和提交日志(CommitLog),延迟低、吞吐高。 | 适合日志、事件、IoT 等高频写入场景。 |
| CQL(Cassandra Query Language) | 提供类 SQL 的查询语言,简化开发人员上手难度。 | 不支持 JOIN、子查询、外键等关系型特性。 |
1.3 与其他数据库的对比(如 MySQL、MongoDB、HBase)
| 对比维度 | Cassandra | MySQL(关系型) | MongoDB(文档型) | HBase(列式) |
|---|
| 数据模型 | 宽列模型(Wide-column),基于 Partition Key 组织 | 表格模型,严格 Schema | BSON 文档,动态 Schema | 列族模型,强依赖 HDFS |
| 一致性模型 | 可调一致性(默认最终一致) | 强一致性(ACID) | 可选一致性(多数派或本地) | 强一致性(依赖 ZooKeeper) |
| 扩展方式 | 水平扩展,无中心架构 | 垂直扩展为主,分库分表复杂 | 水平扩展(Sharding) | 水平扩展,依赖 Hadoop 生态 |
| 写入性能 | 极高(顺序写 CommitLog + 内存 Memtable) | 中等(受索引和事务影响) | 高 | 高(但写路径较重) |
| 查询能力 | 支持主键查询和有限二级索引 | 全功能 SQL(JOIN、子查询、聚合等) | 丰富查询(嵌套字段、数组、地理空间) | 仅支持 RowKey 查询,Scan 性能差 |
| 事务支持 | 无跨行事务,仅单分区原子性 | 完整 ACID 事务 | 单文档原子性(4.0+ 支持多文档事务) | 无事务 |
| 运维复杂度 | 中等(需理解一致性、Compaction、Token 环) | 低(成熟工具链) | 中 | 高(依赖 Hadoop/ZooKeeper/HDFS) |
| 典型适用场景 | 时间序列、IoT、消息系统、高写入日志 | 金融交易、ERP、CRM 等强一致性业务 | 内容管理、用户画像、实时分析 | 大数据分析(与 Spark/Hive 集成) |
第二章:Cassandra 架构原理
2.1 分布式架构与无中心设计
| 概念名称 | 说明 | 注意事项 |
|---|
| 无中心(Masterless) | 所有节点角色对等,任意节点均可接收客户端读写请求并协调操作。 | 避免了单点故障,提升系统可用性。 |
| Token Ring | 整个集群的数据空间被划分为连续的 token 环(默认使用 Murmur3Partitioner),每个节点负责一段 token 范围。 | 新增节点会分割现有 token 范围,触发数据迁移。 |
| 虚拟节点(Vnodes) | 每个物理节点可承担多个虚拟 token(默认 256 个),使数据分布更均匀。 | 启用 vnodes 后无需手动分配 token。 |
| 协调者(Coordinator) | 接收客户端请求的节点作为协调者,负责路由请求到副本节点并聚合结果。 | 协调者本身不一定是数据持有者。 |
| Snitch | 用于确定节点所属数据中心和机架位置,影响复制策略和读写路径选择。 | 常用 SimpleSnitch(单 DC)或 GossipingPropertyFileSnitch(多 DC)。 |
2.2 数据模型:Keyspace、Table、Partition、Clustering Column
| 概念名称 | 说明 | 注意事项 |
|---|
| Keyspace | 类似于关系数据库中的”数据库”,是表的容器,定义复制策略和副本数。 | 创建时必须指定 replication 策略。 |
| Table(Column Family) | 存储数据的逻辑结构,由主键(Primary Key)和若干列组成。 | 表名在 keyspace 内唯一。 |
| Partition Key | 主键的第一部分,决定数据存储在哪个分区(即哪个 token 范围)。 | 相同 Partition Key 的数据物理上存储在一起。 |
| Clustering Columns | 主键中 Partition Key 之后的部分,用于在分区内排序和索引行。 | 查询时若未指定完整主键,可能需扫描整个分区。 |
| Row | 由 Partition Key + Clustering Columns 唯一标识的一条记录。 | 单个分区可包含数百万行(宽行)。 |
| Static Column | 在一个分区内所有行共享的列值,仅存储一次。 | 适用于分区元数据(如用户昵称)。 |
2.3 一致性级别(Consistency Level)
| 一致性级别 | 说明 | 注意事项 |
|---|
| ANY | 写入只需被任一节点(包括 Hinted Handoff 节点)接受即成功。 | 可能丢失数据,仅用于极端高可用场景。 |
| ONE | 读/写操作只需 1 个副本确认。 | 性能最高,但可能读到旧数据。 |
| QUORUM | 需要多数副本确认(公式:(replication_factor / 2) + 1)。 | 平衡一致性与可用性,推荐生产环境使用。 |
| LOCAL_QUORUM | 仅在本地数据中心达到 QUORUM,适用于多 DC 部署。 | 减少跨 DC 延迟。 |
| EACH_QUORUM | 所有数据中心都必须达到 QUORUM(仅用于写操作)。 | 写延迟高,极少使用。 |
| ALL | 所有副本必须响应。 | 强一致性,但任一节点宕机将导致失败。 |
| SERIAL / LOCAL_SERIAL | 用于轻量级事务(LWT),保证 compare-and-set 操作的线性一致性。 | 性能开销大,慎用。 |
注:读写一致性可独立设置,例如写用 ONE,读用 QUORUM 可实现”最终一致+读修复”。
2.4 复制策略(SimpleStrategy vs NetworkTopologyStrategy)
| 策略名称 | 说明 | 注意事项 |
|---|
| SimpleStrategy | 将副本按 token 环顺序放置在后续节点上,不考虑机架或数据中心拓扑。 | 仅适用于单数据中心测试环境。 |
| NetworkTopologyStrategy | 按数据中心(DC)和机架(Rack)显式指定副本数量,支持地理冗余。 | 生产环境必须使用此策略。 |
| 配置示例(NetworkTopologyStrategy) | {'class': 'NetworkTopologyStrategy', 'DC1': 3, 'DC2': 2} | DC 名称需与 snitch 配置一致。 |
| 副本放置规则 | 同一分区的副本优先放在不同机架,避免机架故障导致数据不可用。 | 需配合 GossipingPropertyFileSnitch 使用。 |
2.5 Gossip 协议与故障检测
| 概念名称 | 说明 | 注意事项 |
|---|
| Gossip 协议 | 节点每秒随机与其他 1~3 个节点交换集群状态(如心跳、负载、schema 版本)。 | 去中心化的心跳机制,O(log N) 收敛速度。 |
| Heartbeat | 每次 Gossip 包含一个递增的版本号,用于检测节点是否活跃。 | 心跳超时(默认 30 秒)标记为 DOWN。 |
| Failure Detector | 基于 Phi Accrual 算法判断节点是否故障,比固定超时更适应网络波动。 | 可通过 phi_convict_threshold 调整敏感度。 |
| Seed Nodes | 新节点启动时联系的初始节点列表,用于加入集群。 | 不需要所有节点都是 seed,通常 2~3 个即可。 |
| 状态同步 | Schema 变更(如建表)通过 Gossip 传播到全集群。 | 所有节点 schema_version 必须一致。 |
2.6 Hinted Handoff 与 Read Repair
| 机制名称 | 说明 | 注意事项 |
|---|
| Hinted Handoff | 当目标副本节点宕机时,协调者暂存写请求(hint),待其恢复后重放。 | 默认保存 3 小时(可配置 hinted_handoff_enabled 和 max_hint_window_in_ms)。 |
| Read Repair | 读取时若发现副本数据不一致,自动在后台修复较旧的副本。 | 由 read_repair_chance 或 dclocal_read_repair_chance 触发(新版默认关闭)。 |
| Background Read Repair | 即使客户端使用 ONE 一致性,Cassandra 仍可异步读取额外副本进行修复。 | 通过 speculative_retry 配置触发条件。 |
| Anti-Entropy Repair | 使用 nodetool repair 手动触发 Merkle tree 对比,修复长期不一致数据。 | 建议定期执行,尤其在节点长时间离线后。 |
| Digest Request | 读操作中协调者先向副本发送摘要请求(digest),仅当摘要不一致才拉取完整数据。 | 减少网络传输,提升读性能。 |
第三章:CQL(Cassandra Query Language)基础
3.1 CQL 与 SQL 的异同
| 对比项 | CQL | SQL(传统关系型) | 注意事项 |
|---|
| 语法风格 | 类似 SQL,支持 SELECT、INSERT、UPDATE 等关键字 | 标准 ANSI SQL | 表面相似,底层模型完全不同。 |
| JOIN 支持 | 不支持 JOIN 操作 | 支持多表 JOIN | 需通过反范式化建模避免关联查询。 |
| 子查询 | 不支持子查询 | 广泛支持 | 所有数据必须在单次查询中可获取。 |
| 主键约束 | 必须定义主键(Partition Key + 可选 Clustering Columns) | 可选主键,支持外键 | 主键决定数据分布与查询能力。 |
| 数据更新 | UPDATE 实际是 UPSERT(存在则覆盖,不存在则插入) | UPDATE 仅修改已有行 | CQL 无”行不存在”错误。 |
| NULL 值处理 | 不存储 NULL 值,写入 NULL 相当于删除该列 | 显式存储 NULL | 查询时未设置的列返回 null。 |
| 事务 | 仅保证单分区内的原子性,无跨分区事务 | 支持 ACID 事务 | 无法回滚,需应用层补偿。 |
| 聚合函数 | 支持 COUNT、MAX、MIN、AVG、SUM(但性能差,慎用) | 全功能聚合 | 聚合需扫描整个分区,大数据量下禁用。 |
3.2 Keyspace 操作(CREATE / ALTER / DROP)
| 方法名称 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| CREATE KEYSPACE | CREATE KEYSPACE keyspace_name WITH replication = {...} [AND durable_writes = true/false]; | 创建新的 keyspace | CREATE KEYSPACE myks WITH replication = {'class': 'SimpleStrategy', 'replication_factor': 3}; | 必须指定 replication 策略;名称区分大小写(若用双引号)。 |
| ALTER KEYSPACE | ALTER KEYSPACE keyspace_name WITH replication = {...} [AND durable_writes = ...]; | 修改 keyspace 的复制策略或持久写入设置 | ALTER KEYSPACE myks WITH replication = {'class': 'NetworkTopologyStrategy', 'DC1': 3}; | 不能降级副本数(如从 3 改为 2),需先执行 repair。 |
| DROP KEYSPACE | DROP KEYSPACE [IF EXISTS] keyspace_name; | 删除整个 keyspace 及其所有表 | DROP KEYSPACE IF EXISTS myks; | 操作不可逆,数据永久丢失。 |
3.3 表操作(CREATE / ALTER / DROP TABLE)
| 方法名称 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| CREATE TABLE | CREATE TABLE table_name (columns..., PRIMARY KEY (...)); | 创建新表 | CREATE TABLE users (id UUID, name TEXT, email TEXT, PRIMARY KEY (id)); | 主键必须包含 Partition Key;列类型不可变。 |
| ALTER TABLE ADD | ALTER TABLE table_name ADD column_name type; | 添加新列 | ALTER TABLE users ADD age INT; | 不能添加主键列;新列对旧数据默认为 null。 |
| ALTER TABLE DROP | ALTER TABLE table_name DROP column_name; | 删除列(逻辑删除) | ALTER TABLE users DROP email; | 实际数据仍存在于 SSTable,需 compaction 清理。 |
| ALTER TABLE RENAME | ALTER TABLE table_name RENAME old_name TO new_name; | 重命名列(仅限 clustering 列) | ALTER TABLE events RENAME event_time TO ts; | 不能重命名 partition key 列。 |
| DROP TABLE | DROP TABLE [IF EXISTS] table_name; | 删除表及其所有数据 | DROP TABLE IF EXISTS users; | 操作不可逆;删除后立即释放内存,磁盘数据随 compaction 清除。 |
3.4 数据写入(INSERT / UPDATE / UPSERT)
| 方法名称 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| INSERT | INSERT INTO table (cols...) VALUES (vals...) [USING TTL N]; | 插入新行或更新现有行 | INSERT INTO users (id, name) VALUES (uuid(), 'Alice'); | 若主键已存在,则覆盖非主键列;TTL 可设过期时间(秒)。 |
| UPDATE | UPDATE table SET col = val [WHERE pk = ...] [USING TTL N]; | 更新指定主键的行 | UPDATE users SET email = 'a@example.com' WHERE id = abc123; | 必须包含完整 Partition Key;不能更新主键列。 |
| UPSERT | CQL 中无独立 UPSERT 语句,INSERT/UPDATE 均为 upsert 语义 | 不存在则插入,存在则更新 | 同 INSERT 或 UPDATE 示例 | 所有写操作均为幂等。 |
| USING TIMESTAMP | 在 INSERT/UPDATE 中指定写入时间戳 | 用于解决冲突或回填历史数据 | INSERT INTO logs (id, msg) VALUES (1, 'test') USING TIMESTAMP 1670000000000000; | 时间戳单位为微秒;较小时间戳会被忽略(Last Write Wins)。 |
3.5 数据查询(SELECT)
| 方法名称 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| SELECT | SELECT cols FROM table [WHERE pk = ... AND ck = ...] [LIMIT N]; | 查询数据 | SELECT name, email FROM users WHERE id = abc123; | WHERE 必须包含完整 Partition Key;可选 Clustering 条件。 |
| ALLOW FILTERING | 在 SELECT 后添加此子句 | 允许对非主键列过滤(全分区扫描) | SELECT * FROM events WHERE type = 'login' ALLOW FILTERING; | 性能极差,仅用于调试或小数据集。 |
| ORDER BY | SELECT ... ORDER BY clustering_col DESC/ASC | 按 clustering 列排序 | SELECT * FROM events WHERE user_id = 123 ORDER BY ts DESC; | 仅支持主键中的 clustering 列;方向需与建表一致(除非使用 COMPACT STORAGE)。 |
| COUNT | SELECT COUNT(*) FROM table WHERE pk = ... | 统计分区行数 | SELECT COUNT(*) FROM events WHERE user_id = 123; | 不支持跨分区 COUNT;大数据分区可能超时。 |
| IN 查询 | WHERE pk IN (val1, val2, ...) | 一次查询多个分区 | SELECT * FROM users WHERE id IN (abc, def); | 不推荐用于高并发场景;每个分区独立查询。 |
3.6 数据删除(DELETE / TRUNCATE)
| 方法名称 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| DELETE 列 | DELETE col1, col2 FROM table WHERE pk = ... [AND ck = ...]; | 删除指定列(标记 tombstone) | DELETE email FROM users WHERE id = abc123; | 仅删除指定列,其他列保留;tombstone 需 compaction 清理。 |
| DELETE 行 | DELETE FROM table WHERE pk = ... [AND ck = ...]; | 删除整行 | DELETE FROM events WHERE user_id = 123 AND ts = '2025-01-01'; | 必须提供完整主键(至少 partition key)。 |
| TRUNCATE | TRUNCATE [TABLE] table_name; | 清空整张表 | TRUNCATE events; | 等效于 DROP + CREATE,但保留 schema;立即生效,不可逆。 |
| TTL 自动删除 | 写入时设置 TTL,到期自动删除 | 基于时间的自动清理 | INSERT INTO sessions (id, data) VALUES (1, '...') USING TTL 3600; | 过期后生成 tombstone,仍需 compaction 物理删除。 |
第四章:Cassandra 命令行操作(cqlsh)
4.1 启动与连接 cqlsh
| 操作名称 | 操作细节 | 注意事项 |
|---|
| 启动本地连接 | 在终端执行 cqlsh,默认连接本机 127.0.0.1:9042 | 需确保 Cassandra 服务已启动且监听 9042 端口。 |
| 指定主机和端口 | 执行 cqlsh <host> [port],例如 cqlsh 192.168.1.10 9042 | 若未指定端口,默认使用 9042。 |
| 使用用户名密码登录 | 执行 cqlsh -u <username> -p <password> [host] | 需在 cassandra.yaml 中启用 PasswordAuthenticator。 |
| SSL 加密连接 | 使用 cqlsh --ssl [host],需提前配置客户端证书和 truststore | 需在 cassandra.yaml 和 cqlshrc 中正确设置 SSL 参数。 |
| 退出 cqlsh | 输入 exit 或 quit,或按 Ctrl+D | 会话结束,断开与集群的连接。 |
4.2 常用 cqlsh 命令(DESCRIBE, SHOW, SOURCE 等)
| 命令名称 | 语法 | 用途 | 示例 | 注意事项 |
|---|
| DESCRIBE KEYSPACES | DESCRIBE KEYSPACES; 或 DESC KEYSPACES; | 列出所有 keyspace | DESCRIBE KEYSPACES; | 不显示系统 keyspace(如 system_schema)除非显式查询。 |
| DESCRIBE KEYSPACE | DESCRIBE KEYSPACE keyspace_name; | 显示指定 keyspace 的定义(含复制策略) | DESCRIBE KEYSPACE myks; | 可省略分号。 |
| DESCRIBE TABLE | DESCRIBE TABLE table_name; | 显示表结构(列、主键、属性等) | DESCRIBE TABLE users; | 若表名不唯一,需加 keyspace 前缀:myks.users。 |
| DESCRIBE CLUSTER | DESCRIBE CLUSTER; | 显示集群名称和 partitioner 类型 | DESCRIBE CLUSTER; | 用于确认当前连接的集群信息。 |
| SHOW HOST | SHOW HOST; | 显示当前连接的节点地址和版本 | SHOW HOST; | 包含 Cassandra 版本和 data center 名称。 |
| SHOW VERSION | SHOW VERSION; | 显示 cqlsh 和 Cassandra 的版本号 | SHOW VERSION; | 用于排查兼容性问题。 |
| SOURCE | SOURCE 'file_path.cql'; | 执行外部 CQL 脚本文件 | SOURCE '/home/user/init.cql'; | 路径需为绝对路径或相对于启动目录;文件编码应为 UTF-8。 |
| CAPTURE | CAPTURE 'output.txt'; | 将后续命令输出重定向到文件 | CAPTURE 'result.log'; SELECT * FROM t; CAPTURE OFF; | 用于记录查询结果;执行 CAPTURE OFF 停止捕获。 |
| EXPAND | EXPAND ON; / EXPAND OFF; | 控制多列结果是否垂直显示(每列一行) | EXPAND ON; SELECT * FROM large_row; | 适合查看宽行数据。 |
| PAGING | PAGING ON; / PAGING OFF; / PAGING N; | 控制分页行为(默认每页 100 行) | PAGING 50; | 关闭分页(PAGING OFF)可能导致内存溢出。 |
4.3 执行 CQL 脚本文件
| 操作名称 | 操作细节 | 注意事项 |
|---|
| 使用 SOURCE 命令 | 在 cqlsh 交互模式中执行 SOURCE 'path/to/script.cql'; | 脚本中每条 CQL 语句需以分号结尾;不支持 shell 命令。 |
| 命令行直接执行脚本 | 执行 cqlsh -f script.cql [host] | 适用于自动化部署;执行后自动退出。 |
| 脚本内容要求 | 包含合法 CQL 语句(如 CREATE、INSERT),可包含注释(-- 或 /* */) | 不支持变量替换或流程控制(如 if/loop)。 |
| 错误处理 | 脚本中任一语句失败将停止执行(无事务回滚) | 建议先在测试环境验证脚本。 |
| 编码与路径 | 脚本文件应为 UTF-8 编码;路径建议使用绝对路径 | Windows 路径需转义反斜杠或使用正斜杠。 |
4.4 配置 cqlsh(如时间格式、分页等)
| 配置项 | 配置方式与说明 | 示例值 / 选项 | 注意事项 |
|---|
| cqlshrc 文件位置 | 用户主目录下:~/.cassandra/cqlshrc(Linux/macOS)或 %USERPROFILE%\.cassandra\cqlshrc(Windows) | - | 首次使用需手动创建目录和文件。 |
时间格式(time_format) | 在 [ui] 节点下设置 time_format = %Y-%m-%d %H:%M:%S%z | %Y-%m-%d %H:%M:%S%z | 影响 SELECT 结果中 timestamp 类型的显示。 |
分页大小(page_size) | 在 [ui] 节点下设置 page_size = 200 | 默认 100 | 大于 5000 可能影响性能。 |
| 默认 keyspace | 在 [connection] 节点下设置 keyspace = myks | - | 启动后自动 USE 指定 keyspace。 |
| 连接超时 | 在 [connection] 节点下设置 connect_timeout = 10(秒) | 默认 5 秒 | 适用于高延迟网络。 |
| SSL 配置 | 在 [connection] 下配置 ssl = true,并指定 certfile、truststore 等 | 需与 Cassandra 服务端 SSL 配置匹配 | 用于安全连接生产集群。 |
| 自动补全与历史 | cqlsh 自动保存命令历史到 ~/.cassandra/cqlsh_history | - | 可通过上下箭头调用历史命令。 |
注:修改 cqlshrc 后需重启 cqlsh 生效。
第五章:高级数据建模
5.1 主键设计(Partition Key + Clustering Columns)
| 概念名称 | 说明 | 注意事项 |
|---|
| Partition Key | 主键的第一部分,决定数据存储在哪个节点(通过 partitioner 哈希)。 | 应避免”大分区”(单个分区 > 100MB);选择高基数字段。 |
| Composite Partition Key | 多列组合作为 Partition Key,用括号包裹:PRIMARY KEY ((col1, col2), ...) | 适用于需联合分布的场景(如 tenant_id + device_id)。 |
| Clustering Columns | 主键中 Partition Key 之后的部分,用于分区内排序和高效范围查询。 | 查询时可使用 =、IN、>、<、BETWEEN 等操作符。 |
| 主键定义语法 | PRIMARY KEY (partition_col) 或 PRIMARY KEY ((p1, p2), c1, c2) | 分区键必须完整指定才能查询;Clustering 列可部分使用。 |
| 分区大小控制 | 单个分区建议不超过 10 万行或 100 MB | 过大分区会导致读写延迟升高、OOM 风险。 |
| 查询模式驱动设计 | 先明确查询需求,再反推主键结构 | Cassandra 是”查询驱动建模”,非”实体驱动”。 |
5.2 时间序列数据建模
| 建模策略 | 说明 | 示例表结构 | 注意事项 |
|---|
| 按时间桶分片 | 将 Partition Key 加入时间粒度(如 day、hour),避免单一分区过大 | PRIMARY KEY ((device_id, date), timestamp) | date 可为 '2025-02-02',timestamp 为 clustering 列。 |
| 时间倒序存储 | 使用 WITH CLUSTERING ORDER BY (timestamp DESC) | CREATE TABLE events (...) WITH CLUSTERING ORDER BY (ts DESC); | 便于快速获取最新 N 条记录。 |
| TTL 自动过期 | 写入时设置 USING TTL,自动清理历史数据 | INSERT INTO logs (...) VALUES (...) USING TTL 2592000; — 30天 | 避免手动删除,减少 tombstone。 |
| 避免全时间扫描 | 不使用 ALLOW FILTERING 查询跨多日数据 | 应由应用层循环查询每日分区 | 跨分区聚合需外部系统(如 Spark)。 |
| 高基数设备处理 | 若 device_id 基数极高,可加入 hash 桶:((bucket, date), device_id, ts) | bucket = hash(device_id) % 16 | 防止热点分区。 |
5.3 反范式化与宽行设计
| 设计原则 | 说明 | 注意事项 |
|---|
| 反范式化 | 为每个查询模式单独建表,冗余存储数据 | 牺牲存储换查询性能;更新需写多表(应用层保证最终一致)。 |
| 宽行(Wide Row) | 单个 Partition Key 对应大量 Clustering Rows(如用户所有订单) | 利用分区内排序实现高效范围扫描;避免无限增长。 |
| 查询覆盖 | 表结构包含查询所需全部字段,避免额外 lookup | 例如:消息表直接存 sender_name 而非 sender_id。 |
| 写多读少 | 接受多次写入(如更新用户资料需改 3 张表),换取 O(1) 读取 | 适合读远多于写的场景(如社交 feed)。 |
| 数据生命周期管理 | 结合 TTL 或定期归档,防止宽行无限膨胀 | 可按时间窗口滚动创建新表(如 monthly_events_202502)。 |
5.4 使用集合类型(SET, LIST, MAP)
| 类型名称 | 语法示例 | 用途 | 限制与注意事项 |
|---|
| SET | tags SET<TEXT> | 存储无序、唯一元素集合 | 最大 65535 个元素;全量替换,不支持部分更新。 |
| LIST | comments LIST<TEXT> | 存储有序、可重复列表 | 插入/删除需重写整个 list;性能随长度下降。 |
| MAP | metadata MAP<TEXT, TEXT> | 存储键值对 | key 必须唯一;同样全量更新。 |
| 更新语法 | UPDATE table SET tags = tags + {'new'} WHERE ... | 添加元素(SET/MAP)或追加(LIST) | 不能直接修改某一项(如 map['key'] = 'val' 仅限 CQL 3.0+)。 |
| 删除元素 | DELETE tags['key'] FROM ... 或 UPDATE SET tags = tags - {'old'} | 移除特定元素 | 删除后生成 tombstone;频繁操作导致读放大。 |
| 存储开销 | 集合以单独 cell 存储,每个元素一个列名 | 大集合显著增加 SSTable 大小 | 建议元素数 < 1000;否则拆分为独立表。 |
5.5 用户自定义类型(UDT)
| 操作名称 | 语法 | 用途 | 示例 | 注意事项 |
|---|
| CREATE TYPE | CREATE TYPE address (street TEXT, city TEXT, zip INT); | 定义复合结构 | - | UDT 名在 keyspace 内唯一。 |
| 使用 UDT 列 | addr FROZEN<address> | 在表中引用 UDT | CREATE TABLE users (id UUID, home FROZEN<address>, PRIMARY KEY (id)); | 必须用 FROZEN 包裹,表示整体更新。 |
| 插入 UDT | INSERT INTO users (id, home) VALUES (uuid(), {street: 'Main', city: 'NYC', zip: 10001}); | 写入嵌套结构 | - | 字段名和类型必须匹配。 |
| 查询 UDT 字段 | SELECT home.city FROM users WHERE id = ... | 访问嵌套属性 | - | 不能对 UDT 内部字段建索引。 |
| ALTER TYPE ADD | ALTER TYPE address ADD country TEXT; | 扩展 UDT(添加字段) | - | 新字段对旧数据为 null;不能删除或重命名字段。 |
| UDT 限制 | 不支持嵌套 UDT(Cassandra 4.0+ 支持有限嵌套) | - | - | 避免深度嵌套;FROZEN 意味着无法部分更新。 |
5.6 二级索引与物化视图(Materialized Views)
| 功能名称 | 说明 | 语法示例 | 注意事项 |
|---|
| 二级索引(2i) | 对非主键列创建索引,支持等值查询 | CREATE INDEX ON users (email); | 仅适合低基数列(如 status);高基数列性能差;不支持范围查询。 |
| SASI 索引 | 支持 LIKE、范围、全文搜索的高级索引(需显式启用) | CREATE CUSTOM INDEX ON logs (message) USING 'org.apache.cassandra.index.sasi.SASIIndex'; | 需在 cassandra.yaml 启用;构建耗资源;不保证实时性。 |
| 物化视图(MV) | 自动维护基于原表的另一种主键视图 | CREATE MATERIALIZED VIEW users_by_email AS SELECT * FROM users WHERE email IS NOT NULL PRIMARY KEY (email, id); | 写放大(原表写 → MV 写);Cassandra 4.0 起默认禁用,需显式开启。 |
| MV 限制 | 视图主键必须包含原表主键所有列 | - | 不支持聚合、JOIN、UDT;一致性弱于原表。 |
| 替代方案 | 应用层维护”反向索引表” | 如单独建表 email_to_user_id | 更可控、性能更优;推荐生产环境使用。 |
| 索引性能警告 | 2i 和 MV 均在每个节点本地构建,跨节点查询效率低 | - | 避免在高频写入表上使用;优先考虑主键设计满足查询。 |
第六章:性能调优与运维
6.1 读写路径与 Memtable / SSTable
| 概念名称 | 说明 | 注意事项 |
|---|
| 写入路径 | 客户端写请求 → CommitLog(持久化) + Memtable(内存) → 后台刷盘为 SSTable | CommitLog 保证宕机不丢数据;Memtable 按 Partition Key 排序。 |
| Memtable | 内存中的可变数据结构,每个表一个,默认 512MB 或 10% 堆内存触发 flush | 可通过 memtable_heap_space_in_mb 调整大小。 |
| SSTable | 不可变的磁盘文件,包含数据、索引、摘要(Summary/Bloom Filter)等组件 | 多个 SSTable 需通过 Compaction 合并。 |
| 读取路径 | 先查 Memtable → 再查 Row Cache(可选)→ 最后查 SSTable(通过 Bloom Filter 快速过滤) | Bloom Filter 可能误报(false positive),但不会漏报。 |
| Row Cache | 缓存热数据行(默认关闭),适用于读多写少且数据集小的场景 | 开启会增加 GC 压力;建议仅缓存 < 1% 总数据量。 |
| Key Cache | 缓存 SSTable 的分区键偏移位置(默认开启) | 减少磁盘 seek;命中率可通过 nodetool info 查看。 |
6.2 Compaction 策略(SizeTiered, Leveled, TimeWindow)
| 策略名称 | 说明 | 适用场景 | 配置示例 | 注意事项 |
|---|
| SizeTieredCompactionStrategy (STCS) | 将大小相近的 SSTable 合并(默认策略) | 写密集、无 TTL、随机更新场景 | compaction = {'class': 'SizeTieredCompactionStrategy', 'bucket_high': '1.5', 'bucket_low': '0.5'} | 可能产生大文件;读放大较高;空间放大明显(最多 2x)。 |
| LeveledCompactionStrategy (LCS) | 将 SSTable 分层(L0~Ln),每层大小固定,合并更频繁 | 读密集、更新频繁、需低延迟读取 | compaction = {'class': 'LeveledCompactionStrategy', 'sstable_size_in_mb': '160'} | 磁盘写放大高(约 10x);初始 L0 层无排序;适合 SSD。 |
| TimeWindowCompactionStrategy (TWCS) | 按时间窗口(如 1 天)分组 SSTable,窗口内用 STCS 合并 | 时间序列数据、带 TTL、按时间查询 | compaction = {'class': 'TimeWindowCompactionStrategy', 'compaction_window_unit': 'DAYS', 'compaction_window_size': '1'} | 窗口结束后不再合并;避免跨窗口查询。 |
| 通用参数 | max_threshold(最大合并文件数)、min_threshold(最小触发数) | - | - | 修改策略需谨慎,可能触发全量重写。 |
6.3 节点扩容与数据再平衡
| 操作步骤名称 | 操作细节 | 注意事项 |
|---|
| 添加新节点 | 1. 配置 cassandra.yaml(cluster_name、seeds、listen_address 等) 2. 启动 Cassandra 服务 | 新节点自动加入 ring;无需手动分配 token(vnodes 默认启用)。 |
| 自动数据迁移 | 新节点启动后,通过 Gossip 发现集群状态,自动从副本节点流式拉取数据 | 流式传输(Streaming)期间不影响线上读写。 |
| 手动触发再平衡 | 使用 nodetool rebuild(从其他 DC 拉数据)或 nodetool repair(修复不一致) | 通常无需手动干预;扩容后建议运行 repair。 |
| 监控数据分布 | 使用 nodetool status 查看各节点 ownership 和 load | 理想情况下 ownership 应接近 1/N(N=节点数)。 |
| 避免热点 | 确保 partition key 设计均匀;vnodes 默认 256 个虚拟 token | 若关闭 vnodes,需手动计算 token 并设置 initial_token。 |
| 下线节点 | 执行 nodetool decommission(优雅下线)或 nodetool removenode(强制移除) | decommission 会将数据迁出;removenode 用于已宕机节点。 |
6.4 快照备份与恢复
| 操作名称 | 操作细节 | 命令示例 | 注意事项 |
|---|
| 创建快照 | 对指定 keyspace 或表生成硬链接(几乎瞬时完成) | nodetool snapshot myks 或 nodetool snapshot -t backup_20250202 myks | 快照存储在 data_dir/keyspace/table/snapshots/ 目录下。 |
| 列出快照 | 查看现有快照 | nodetool listsnapshots | 显示快照名、表、占用空间。 |
| 清理快照 | 删除指定快照释放磁盘空间 | nodetool clearsnapshot myks 或 nodetool clearsnapshot -t backup_20250202 | 仅删除硬链接,不影响原 SSTable。 |
| 从快照恢复 | 1. 停止 Cassandra 2. 替换对应表目录下的 SSTable 文件 3. 启动并执行 repair | cp -r snapshots/backup_20250202/* table_dir/ | 需确保 schema 一致;恢复后必须运行 nodetool repair。 |
| 自动快照 | 在 cassandra.yaml 中设置 auto_snapshot: true(默认开启) | - | DROP TABLE/KEYSPACE 前自动创建快照。 |
| 备份到远程 | 快照本身是本地文件,需配合 rsync、s3cmd 等工具上传 | rsync -av /var/lib/cassandra/data/myks/table/snapshots/ user@backup:/backup/ | 快照不含 commitlog,不能用于 PITR(时间点恢复)。 |
| 命令名称 | 语法 | 用途 | 示例 | 注意事项 |
|---|
| nodetool status | nodetool status [keyspace] | 查看集群节点状态(UN=Up Normal) | nodetool status myks | 显示 ownership、load、token range。 |
| nodetool info | nodetool info | 查看本节点基本信息(ID、DC、Rack、版本) | nodetool info | 包含 heap、cache 命中率等。 |
| nodetool tpstats | nodetool tpstats | 查看线程池状态(pending、active 任务数) | nodetool tpstats | 用于诊断写入积压或读延迟。 |
| nodetool compactionstats | nodetool compactionstats | 查看当前 compaction 进度 | nodetool compactionstats | 显示任务数、剩余字节。 |
| nodetool flush | nodetool flush [keyspace] [table] | 强制 Memtable 刷盘为 SSTable | nodetool flush myks users | 用于备份前确保数据落盘。 |
| nodetool repair | nodetool repair [keyspace] [table] | 触发反熵修复,同步副本间不一致数据 | nodetool repair -pr myks | -pr 表示只修复主 token 范围,避免重复。 |
| nodetool cleanup | nodetool cleanup [keyspace] | 删除当前节点不再负责的 stale 数据 | nodetool cleanup myks | 扩容后必须在旧节点执行。 |
| nodetool drain | nodetool drain | 停止接受写入,flush 所有 Memtable | nodetool drain | 用于安全停机前准备。 |
| nodetool describecluster | nodetool describecluster | 显示集群名称、partitioner、snitch 等 | nodetool describecluster | 用于确认集群拓扑一致性。 |
| nodetool getendpoints | nodetool getendpoints ks table key | 查询某 partition key 存储在哪些节点 | nodetool getendpoints myks users abc123 | 用于调试数据分布。 |
第七章:安全与监控
7.1 认证与授权(PasswordAuthenticator)
| 配置项/操作名称 | 配置方式与说明 | 示例值 / 命令 | 注意事项 |
|---|
| 启用认证 | 修改 cassandra.yaml:
authenticator: PasswordAuthenticator
authorizer: CassandraAuthorizer | authenticator: PasswordAuthenticator
authorizer: CassandraAuthorizer | 需重启 Cassandra 生效;默认 superuser 为 cassandra/cassandra。 |
| 创建用户 | 使用 CQL:
CREATE ROLE username WITH PASSWORD = 'pwd' AND LOGIN = true; | CREATE ROLE admin WITH PASSWORD = 'secure123' AND LOGIN = true AND SUPERUSER = true; | 首次登录需用 cassandra 用户创建新管理员。 |
| 授权操作 | GRANT SELECT ON TABLE ks.table TO user;
GRANT MODIFY ON KEYSPACE ks TO user; | GRANT ALL PERMISSIONS ON KEYSPACE myks TO app_user; | 权限包括:SELECT, MODIFY, AUTHORIZE, ALTER, DROP, DESCRIBE。 |
| 角色管理 | 支持角色继承:
CREATE ROLE dev;
GRANT dev TO alice; | - | 推荐基于角色分配权限,而非直接赋权给用户。 |
| 默认 system_auth keyspace | 存储用户和权限信息,需设置 replication_factor ≥ 3 | ALTER KEYSPACE system_auth WITH replication = {'class': 'NetworkTopologyStrategy', 'DC1': 3}; | 若 RF=1,节点故障将导致无法登录。 |
| 密码策略 | 通过自定义 IRoleManager 实现(社区版无内置密码复杂度/过期策略) | - | 企业版(DataStax)支持高级策略。 |
7.2 客户端加密(SSL/TLS)
| 配置项 | 配置方式与说明 | 文件位置 / 参数示例 | 注意事项 |
|---|
| 启用 client-to-node SSL | 在 cassandra.yaml 中配置:
client_encryption_options:
enabled: true
keystore: ...
keystore_password: ... | keystore: /etc/cassandra/conf/.keystore
truststore: /etc/cassandra/conf/.truststore | 需提前生成 Java keystore 和 truststore。 |
| 生成证书 | 使用 keytool:
keytool -genkeypair -alias cassandra -keyalg RSA -keystore .keystore | - | 所有节点应使用相同或互信的证书。 |
| cqlsh 连接加密 | 在 ~/.cassandra/cqlshrc 中配置:
[connection]
ssl = true
ssl_validate = false | 若关闭验证(ssl_validate=false),则无需 truststore | 生产环境应启用证书验证(ssl_validate=true)。 |
| 驱动连接 SSL | Java Driver 示例:
Cluster.builder().withSSL().build(); | Python driver: cluster = Cluster(..., ssl_context=ctx) | 需在客户端配置对应 truststore 或 CA。 |
| 双向 TLS(mTLS) | 设置 require_client_auth: true(在 client_encryption_options 下) | - | 强制客户端提供有效证书,增强安全性。 |
| 性能影响 | SSL 加解密增加 CPU 开销(约 10%~20%) | - | 建议在专用网络中评估是否必要。 |
7.3 日志配置与审计
| 日志类型 | 配置文件与说明 | 关键参数 | 注意事项 |
|---|
| 主日志(system.log) | 由 logback.xml 控制,位于 $CASSANDRA_HOME/conf/logback.xml | /var/log/cassandra/system.log,200MB | 默认 INFO 级别;可调至 DEBUG 用于排错。 |
| 审计日志(Audit Log) | Cassandra 4.0+ 支持,需在 cassandra.yaml 启用:
audit_logging_options:
enabled: true | audit_logs_dir: /var/log/cassandra/audit
included_keyspaces: [myks] | 记录 DDL/DML 操作(如 CREATE、INSERT);对性能有影响,谨慎开启。 |
| GC 日志 | 通过 jvm.options 启用:
-Xloggc:/var/log/cassandra/gc.log
-XX:+PrintGCDetails | - | 用于分析 JVM 停顿问题。 |
| 提交日志(CommitLog) | 非文本日志,但可通过 commitlog_sync 控制持久化策略 | commitlog_sync: periodic
commitlog_sync_period_in_ms: 10000 | 影响写入延迟与持久性权衡。 |
| 日志轮转 | logback.xml 中配置 RollingFileAppender | - | 避免磁盘被日志占满;建议保留 7~30 天。 |
| 敏感信息过滤 | 默认不记录密码、token 等 | - | 自定义 UDF/触发器可能泄露数据,需审查。 |
7.4 监控指标(JMX, Prometheus + Grafana)
| 监控方式 | 说明 | 关键指标 / 工具 | 注意事项 |
|---|
| JMX(内置) | Cassandra 通过 JMX 暴露数百个 MBean 指标 | org.apache.cassandra.metrics:type=Storage, name=Load
type=Table, scope=myks.users, name=ReadLatency | 默认监听 localhost:7199;生产环境需配置远程访问和认证。 |
| nodetool 查看指标 | 如 nodetool tablestats、nodetool proxyhistograms | Read/Write Latency, SSTable count, Space used | 适合临时诊断。 |
| Prometheus 抓取 | 使用 cassandra-exporter 或 jmx_exporter 暴露 JMX 指标为 HTTP metrics endpoint | jmx_exporter 配置 yaml 映射 MBean 到 Prometheus 指标 | 需额外部署 exporter 进程。 |
| Grafana 仪表盘 | 导入社区模板(如 DataStax 或 Instaclustr 提供的 dashboard) | 节点状态、读写吞吐、延迟分布、compaction 队列、cache 命中率 | 建议监控:Pending Compactions、Dropped Mutations、GC Pauses。 |
| 关键业务指标 | - Coordinator Read/Write Latency - Dropped Messages(因超时) - CAS Read/Write Failures | - | Dropped Mutations 表示写入失败,需立即告警。 |
| 告警规则建议 | - Heap usage > 80% - Pending compactions > 100 - Unavailable nodes > 0 | - | 结合 PagerDuty、Alertmanager 实现通知。 |
注:Cassandra 4.0+ 原生支持 OpenTelemetry,未来可替代 JMX。
第八章:开发集成
8.1 Java Driver 使用示例
| 方法/操作名称 | 语法与用途 | 代码示例(Java) | 注意事项 |
|---|
| 创建集群连接 | 使用 CqlSession.builder() 构建会话 | CqlSession session = CqlSession.builder().addContactPoint(new InetSocketAddress("127.0.0.1", 9042)).withLocalDatacenter("datacenter1").build(); | 必须指定 localDatacenter;支持多个 contact point。 |
| 执行同步查询 | session.execute(Statement) 返回 ResultSet | ResultSet rs = session.execute("SELECT name FROM users WHERE id = ?", uuid);
Row row = rs.one();
if (row != null) { String name = row.getString("name"); } | 查询参数使用 ? 占位符;避免拼接 SQL 防注入。 |
| 异步执行 | session.executeAsync(Statement) 返回 CompletableFuture | CompletableFuture future = session.executeAsync("INSERT INTO logs (id, msg) VALUES (?, ?)", id, msg);
future.whenComplete((rs, err) -> { ... }); | 适用于高并发写入;注意异常处理。 |
| 绑定命名参数 | 使用 BoundStatement 与预编译语句 | PreparedStatement ps = session.prepare("UPDATE users SET email = ? WHERE id = ?");
BoundStatement bs = ps.bind("a@example.com", userId);
session.execute(bs); | 预编译提升性能;可复用 PreparedStatement。 |
| 关闭连接 | 调用 session.close() | session.close(); | 应在应用关闭时调用,释放连接和线程资源。 |
| 配置重试与一致性 | 通过 RequestOptions 设置 | SimpleStatement stmt = SimpleStatement.builder("SELECT ...").setConsistencyLevel(DefaultConsistencyLevel.QUORUM).build(); | 默认一致性为 LOCAL_ONE;生产建议 QUORUM。 |
8.2 Python(cassandra-driver)使用示例
| 方法/操作名称 | 语法与用途 | 代码示例(Python) | 注意事项 |
|---|
| 创建集群连接 | 使用 Cluster 类 | from cassandra.cluster import Cluster
cluster = Cluster(['127.0.0.1'], port=9042)
session = cluster.connect('myks') | 默认使用 RoundRobinPolicy 负载均衡。 |
| 执行查询 | session.execute(query, parameters) | rows = session.execute("SELECT name FROM users WHERE id = %s", [user_id])
for row in rows: print(row.name) | 参数使用 %s 占位;返回 Row 对象,属性按列名访问。 |
| 异步执行 | 使用 session.execute_async() | from cassandra.concurrent import execute_concurrent
execute_concurrent(session, , concurrency=10) | 原生不支持 asyncio;需用 concurrent 或线程池。 |
| 使用 UDT | 先注册 UDT 类 | from cassandra.user_types import UserType
class Address(UserType): street = str city = str
cluster.register_user_type('myks', 'address', Address) | 需继承 UserType;字段名必须匹配。 |
| 批量写入 | 使用 BatchStatement | batch = BatchStatement()
batch.add("INSERT INTO t (a,b) VALUES (?, ?)", (1, 'x'))
batch.add("UPDATE t SET b=? WHERE a=?", ('y', 1))
session.execute(batch) | 默认为 LOGGED batch(有开销);非原子。 |
| 关闭连接 | 调用 cluster.shutdown() | cluster.shutdown() | 释放连接池和后台线程。 |
8.3 Spring Data Cassandra 集成
| 配置/操作名称 | 说明 | 代码/配置示例 | 注意事项 |
|---|
| 添加依赖 | Maven 引入 spring-boot-starter-data-cassandra | org.springframework.boot
spring-boot-starter-data-cassandra | 需匹配 Spring Boot 与 Cassandra 版本。 |
| 配置连接 | application.yml 中设置 | spring: cassandra: contact-points: 127.0.0.1 port: 9042 keyspace-name: myks local-datacenter: datacenter1 | 必须指定 local-datacenter。 |
| 定义实体类 | 使用 @Table、@PrimaryKey 等注解 | @Table
public class User { @PrimaryKey private UUID id; private String name;
} | 支持 @Column、@CassandraType 等。 |
| Repository 接口 | 继承 CassandraRepository | public interface UserRepository extends CassandraRepository<User, UUID> { List findByName(String name);
} | 方法名自动解析为 CQL(有限支持);复杂查询需 @Query。 |
| 自定义查询 | 使用 @Query 注解 | @Query("SELECT * FROM users WHERE email = ? ALLOW FILTERING")
List findByEmailCustom(String email); | 避免 ALLOW FILTERING;仅用于调试。 |
| 手动模板操作 | 注入 CassandraTemplate | @Autowired
CassandraTemplate template;
template.select("SELECT ...", User.class); | 适用于动态 CQL 构建。 |
8.4 与 Kafka / Spark 集成场景
| 集成组件 | 场景说明 | 实现方式与工具 | 注意事项 |
|---|
| Kafka → Cassandra | 实时消费消息并写入 Cassandra(如日志、事件流) | 使用 Kafka Connect + Cassandra Sink Connector(如 DataStax 提供) 或自定义 Flink/Spark Streaming 作业 | 控制写入速率避免压垮 Cassandra;使用批量写入。 |
| Spark → Cassandra | 批量分析 Cassandra 数据或写入结果 | 使用 spark-cassandra-connector:
df = spark.read.format("org.apache.spark.sql.cassandra").options(table="users", keyspace="myks").load() | 需 colocate Spark 与 Cassandra 节点(same rack)以优化本地读取。 |
| Cassandra → Kafka | 监听 Cassandra 变更并推送至 Kafka(CDC) | 使用 Cassandra CommitLog 解析工具(如 Cassandra Kafka Sink、Debezium 实验性支持)或应用层双写 | 原生无 CDC;双写需处理一致性。 |
| 流处理架构 | Kafka(缓冲)→ Flink/Spark Streaming(处理)→ Cassandra(存储) | Flink 使用 CassandraSink:
stream.addSink(new CassandraSink<>(...)); | 设置合理的并行度与背压策略。 |
| 数据导出 | 将 Cassandra 表全量导出到 HDFS/S3 供 Spark 分析 | 使用 DSBulk 或 sstable2json + Spark 读取 SSTable | DSBulk 支持高效导出/导入。 |
| 一致性保障 | 在 Kafka-Cassandra 链路中保证 exactly-once | Kafka 事务 + 幂等写入(Cassandra upsert 天然幂等) | 需 Kafka >= 0.11 且启用事务。 |