BigTable
BigTable是存储领域的基石级论文,其设计思想至今还在被广泛使用。在实际数据湖的开发优化中,对于列式存储、稀疏存储有了些体会和零散的理解,重读两遍BigTable将知识串起来。
原文行文比较散,而且内容彼此交织,不便于理解。本文将BigTable拆分为三个子概念理解:逻辑存储格式、物理存储格式、分布式控制。
前言
BigTable的出发点是希望存储PB级别数据并且提供灵活性和高性能。常见的数据库关系模型能够高效处理结构化数据,但是无法高效服务半结构化数据和非结构化数据。WebPage属于半结构化数据,Google Earth包含大量非结构化的视频图片数据,强行使用关系模型得不偿失。
| 维度 | 结构化数据 | 半结构化数据 | 非结构化数据 |
|---|---|---|---|
| 核心特征 | 固定架构,所有数据字段一致 | 有标签/标记,字段可缺失或扩展 | 无预定义模式,自由格式 |
| 模式 | 严格、预先定义 | 动态可变、自描述(同一类型有不同的标识,例如Json对象有不同的key) | 无 |
| 组织方式 | 二维表(行×列) | 树/图结构,键值对、嵌套对象 | 文件流、二进制/文本块 |
| 典型格式 | 关系数据库表、定长CSV、Excel固定表 | JSON、XML、HTML、日志 | 图片、视频、音频、Word/PDF、纯文本 |
| 存储方式 | 关系型数据库(MySQL、PostgreSQL、Oracle) | NoSQL(MongoDB、Redis)、数据湖 | 对象存储(OSS/S3)、分布式文件系统(HDFS) |
| 查询效率 | 高 | 中 | 低,需复杂预处理 |
| 灵活性 | 低,改表结构成本高 | 高,字段可随时增减 | 最高,但难以直接分析 |
BigTable将数据存储为半结构化模式(逻辑存储),从而获得了灵活性,引入LSM-Tree(物理存储)获取高写入性能,将分布式协调交给Chubby(类似Etcd,提供强一致存储和分布式锁),将文件存储和容错交给GFS,使用最小组件来构筑系统,显著降低了分布式复杂度。
可以说是大师之作,值得反复学习。
逻辑存储格式
论文开篇点题,给出了提纲挈领的总结语句:BigTable 是一个稀疏的、分布式的、可持久化的多维有序 map 。map 通过rowKey,columnKey 和时间戳进行索引,map 中的每个值都是一个连续的字节数组。
数据存储格式表示为(row:string, column:string, time:int64) → string 。
以网页存储为例,引出列族和限定符概念:
列族是预先定义的数据属性,限定符则是列族具体的属性,限定符可以为空;列族描述数据的共性(contents + anchor),限定符是具体的值,随存储对象动态变化。网页的内容是HTML,锚点则是引用该网页的网站。显然,引用网页的网站是无法被预先统计和预测的,因此使用动态属性标识能够给予极大的灵活性。
上述示例可以理解为类似Json对象的格式,只是行键和列族名是必备的schema,不能缺省:
1 | { |
行键标识一个对象,列族对应某个属性或某些属性的组合,限定符是某个属性的具体值,使用时间戳区别版本;不同列族下的数据对应不同的对象,因此采取不同的存储模式、访问模式、额度限制、压缩策略、垃圾回收策略(保留最近N个版本,或者足够新的版本)。BigTable是行级存储,提供行键事务,一行对外表现为整体,内部根据列族进行多维展开,展开后被称为单元格。BigTable按照行键字典序来维护数据有序性,行键通常大小是10-100字节,最大支持64KB的行键。行键的选取会影响数据排列,在WebTable示例中,行键是域名翻转的结果,因为网页域名是由小到大排列(.com在最后),翻转后才能让相同根域名下的数据相邻存放从而获得更高的局部性。
一张表可能承载PB级数据,存放到一台机器不合理,因此数据表会在达到一定阈值后(100-200MB)被拆分(选择某个行键一分为二),拆分后连续的多行数据被称为tablet, tablet是分发、负载均衡的最小单位。数据表的格式在用户创建时指定。
元数据存储
BigTable元数据分为协调元数据和表状态元数据:
- 协调元数据:系统全局数据。存放在
Chubby(可以理解为强一致性KV,对外提供KV存储(支持命名空间)和分布式锁);master lock: 确保系统只存在一个主节点root tablet position: 表状态元数据的根tablet(只有一个,永不拆分)tablet server lock: 每台数据存储服务器持有一个分布式锁并定期续约来保活table schema: 表的列族、存储模式、配额信息acl: 访问控制信息
- 状态元数据:记录
tablet的运行信息,用于索引、分裂,迁移,恢复。存储在BigTable自身中,具有固定的schema。- 行键: (
table id,End Row Key) - 列族:
location: 承载tablet的服务器标识log: 每条提交日志的恢复起点 + 组成tablet的SST元数据event:分裂、合并、迁移、恢复信息
- 行键: (
每行状态元数据约1KB,系统使用三级索引,当tablet容量上限设置为128MB时,大约能索引$2^{34}$个tablet,即$2^{61}$字节数据(2EB,2000PB)
物理存储格式
建立在GFS之上的LSM-Tree,只是WAL被称为commit log。
读操作先查内存表再查SST;写操作先写日志,再提交到内存表。
压缩分为三类:
- minor: memtable -> sst
- merge: memtable + some ssts -> sst
- major: all ssts -> sst

BigTable在上述架构基础上采取了诸多优化来获得稳定性能: - Locality groups: 对列族分组,同一个
tablet的不同locality group对应不同的SST;物理隔离不会同时访问的数据、区分不同访问频率和模式的数据。频繁访问的小数据局部组可以声明为In Memory,全量加载到内存后不访问磁盘。 - 数据压缩:Google 采用两级压缩:先用
Bentley-McIlroy的长公共字符串压缩,再用16 KB窗口的快速算法。在WebTable上可以达到10:1的压缩比。 - 两级缓存:
scan cache缓存KV,Table Cache缓存数据块 - 布隆过滤器: 跳过不包含目标的数据块
- 合并提交日志:如果每个
tablet持有独立的日志,那么会产生多个并发小文件写入(一台tablet server管理10~1000tablet,写入性能不佳,将其合并,一台tablet server上的所有tablet写入数据到一个日志文件。tablet server崩溃迁移数据时,为了减少开销,会先对group commit log进行分块并行排序,迁移目标tablet只读取对应部分日志即可。 - 加速恢复:迁移前和停止服务前,进行一次
minor compaction来避免日志重放 - 利用
SST的不可变性
分布式控制
延续了控制面和数据面分离的思路:
- 关键控制信息依赖
Chubby Master单点调度协调(无状态,崩溃不丢数据)- 客户端缓存位置,丢失/过期最多需要三次网络往返(
Chubby -> Root MetaData -> MetaData)获取tablet server location。读写请求直接下发tablet server,通过权限校验后执行
控制面正常运行高度依赖Chubby,只要Chubby不崩溃,不超载,系统就能平稳运行。
崩溃恢复如下图所示:
master崩溃恢复四步走:抢占chubby master lock-> 轮询存活tablet server-> 查找metadata table,定位需要分配、迁移的tablet-> 重建状态tablet server崩溃后无法续约分布式锁,被master监听到锁过期事件后,抢占并删除对应的分布式锁,随后标记tablet server对应的tablet为未分配,执行迁移操作
测试与数据分析
Bigtable 被用于差异极大的场景:
Crawl表压缩比极高(11%),列族和locality group数量多,主要用于批量处理;Google Earth服务表虽然只有0.5TB,但33%数据驻留内存,需要低延迟;Personalized Search拥有93个列族,展示了多租户共享同一张表的能力。
BigTable的读写性能受益于服务器数量,增加服务器数量虽然不能线性scale,但是可以减少单次开销,提高聚合吞吐量。
总结
有序、稀疏、列式存储的典范之作,列族设计常看常新。
不支持跨行事务为客户端使用带来了巨量心智负担。