首页/目录/全部文章

全部文章

八个专题的源码、算法与协议笔记都在这里。

笔记列表

RocksDB在Ceph中的调用与实现分析

RocksDB在Ceph中的调用与实现分析

目录

  1. 文档目标
  2. 整体定位
  3. 核心代码位置
  4. 抽象层KeyValueDB
  5. RocksDBStore适配层
  6. BlueStore如何接入RocksDB
  7. 数据库打开与初始化流程
  8. 写路径分析
  9. 读路径与迭代器分析
  10. 列族与分片机制
  11. MergeOperator机制
  12. BlueFS与Env的关系
  13. 缓存、统计与compact
  14. reshard重分片流程
  15. 典型调用链总结
  16. 关键结论

文档目标

这份文档不是泛泛介绍 upstream RocksDB 的 LSM 树原理,而是专门分析:

  • Ceph 里为什么要用 RocksDB
  • BlueStore 如何通过 KeyValueDB 抽象使用 RocksDB
  • RocksDBStore 如何把 Ceph 的 prefix/key 模型映射到底层 RocksDB
  • 事务、读写、迭代、列族、分片、merge、compact、重分片这些机制在 Ceph 中是怎么落地的

如果你正在阅读这些文件:

  • cephMain/src/kv/KeyValueDB.h
  • cephMain/src/kv/KeyValueDB.cc
  • cephMain/src/kv/RocksDBStore.h
  • cephMain/src/kv/RocksDBStore.cc
  • cephMain/src/os/bluestore/BlueStore.cc

那么这份文档的目标就是把它们串成一条完整主线。


整体定位

在 Ceph 的 BlueStore 架构中,对象数据 主要直接落到块设备上,而 元数据 则主要存放在 RocksDB 中。

可以把它粗略理解为:

  • BlueStore 负责对象语义、空间管理、缓存、事务组织
  • RocksDB 负责元数据的 KV 持久化
  • BlueFS 负责为 RocksDB 提供底层文件存储环境

它们之间的分工如下:

层次 角色 主要职责
OSD/BlueStore 上层调用者 发起对象元数据读写、事务提交、迭代扫描
KeyValueDB 抽象层 定义统一 KV 接口,屏蔽具体后端
RocksDBStore 适配层 把 Ceph 接口映射为 RocksDB API
RocksDB 元数据引擎 管理 memtable、WAL、SST、compaction、迭代器
BlueFS / Env 存储环境 提供 RocksDB 文件读写运行环境

因此,从 Ceph 代码视角看,真正值得重点理解的,不是 RocksDB 全部内部实现,而是:

  1. KeyValueDB 抽象语义
  2. RocksDBStore 如何实现这个抽象
  3. BlueStore 如何调用 KeyValueDB

核心代码位置

1. 抽象接口层

  • cephMain/src/kv/KeyValueDB.h
  • cephMain/src/kv/KeyValueDB.cc

这层定义了通用 KV 数据库接口,包括:

  • 事务对象 TransactionImpl
  • 读取接口 get()
  • 迭代器接口 IteratorImpl / WholeSpaceIteratorImpl
  • merge operator 抽象
  • compact、缓存、统计等可选能力

2. RocksDB适配层

  • cephMain/src/kv/RocksDBStore.h
  • cephMain/src/kv/RocksDBStore.cc

这层把上面的抽象映射为 RocksDB 的:

  • rocksdb::DB
  • rocksdb::WriteBatch
  • rocksdb::Iterator
  • rocksdb::ColumnFamilyHandle
  • rocksdb::Options
  • rocksdb::Env

3. BlueStore接入层

  • cephMain/src/os/bluestore/BlueStore.cc

BlueStore 通过 KeyValueDB::create() 创建数据库实例,通过 db->get_transaction()db->submit_transaction()db->get()db->get_iterator() 等接口完成元数据操作。


抽象层KeyValueDB

KeyValueDB 的核心价值是:把 Ceph 上层的元数据访问模型统一成 prefix/key 语义

1. prefix + key 模型

Ceph 不是简单把所有元数据塞进一个平面的 key 空间,而是按逻辑类别分前缀管理。例如:

  • 某些前缀用于对象元数据
  • 某些前缀用于统计信息
  • 某些前缀用于 omap
  • 某些前缀用于 allocator / freelist 相关元数据

所以 KeyValueDB 大部分接口都长这样:

1
2
3
4
get(prefix, key, value)
set(prefix, key, value)
rmkey(prefix, key)
get_iterator(prefix)

这意味着调用者思考的是“逻辑前缀 + 逻辑 key”,而不是 RocksDB 原始键编码细节。

2. 事务抽象

TransactionImpl 表示一批待提交的写操作,其职责是:

  • 累积 set
  • 累积 rmkey
  • 累积 rm_range_keys
  • 累积 merge
  • 暴露操作数与字节规模

这层并不规定底层必须怎么实现,只规定“提交前是一批内存中的写集合,提交后一次性生效”。

RocksDBStore 中,这个抽象最终会映射成 rocksdb::WriteBatch

3. 迭代器抽象

KeyValueDB 定义了两类迭代器:

  • IteratorImpl
    面向 prefix 过滤后的逻辑视图
  • WholeSpaceIteratorImpl
    面向底层实现暴露的原始全键空间

这种两层设计非常重要,因为:

  • 上层通常只关心某个 prefix 下的 key
  • 实现层有时需要直接处理原始编码键,例如 "prefix\0key"
  • 当存在多个 column family / shard 时,实现层还可能需要做跨迭代器合并

4. MergeOperator抽象

KeyValueDB::MergeOperator 是对 RocksDB associative merge 的抽象包装。

它定义:

  • merge_nonexistent()
  • merge()
  • name()

其意义是:

  • Ceph 上层只注册逻辑 merge 规则
  • 具体怎么挂接到 RocksDB 默认 CF 或独立 CF,由 RocksDBStore 负责

RocksDBStore适配层

RocksDBStore 是这套链路的核心。

它解决的问题是:

Ceph 的 prefix/key 抽象,如何安全且高效地映射到 RocksDB 的 DB / CF / Iterator / WriteBatch / Options / Env 体系。

1. 它管理的关键对象

RocksDBStore 内部最重要的成员包括:

  • rocksdb::DB *db
    真正的 RocksDB 数据库实例
  • rocksdb::Env *env
    底层运行环境,可能是默认文件系统,也可能是 BlueRocksEnv
  • rocksdb::BlockBasedTableOptions bbt_opts
    block cache、filter、index 等表选项
  • default_cf
    默认列族句柄
  • cf_handles
    逻辑 prefix 到实际列族分片句柄的路由表
  • cf_ids_to_prefix
    用于把 RocksDB column family id 反查回 prefix
  • dbstats
    RocksDB 统计对象

2. 它完成的核心映射

映射一:事务抽象 -> WriteBatch

KeyValueDB::TransactionImpl

映射成:

RocksDBStore::RocksDBTransactionImpl

其中核心成员就是:

  • rocksdb::WriteBatch bat

映射二:逻辑key -> RocksDB原始key

有两种模式:

  1. 没有独立列族时
    原始 key = "prefix\0key"

  2. 有独立列族时
    prefix 直接决定列族,列族内只保存裸 key

映射三:prefix -> ColumnFamily

若 prefix 被配置为 column family,则:

  • prefix 直接映射到某个 CF
  • 如果该 prefix 进一步分片,则再根据 key 的 hash 路由到某个 shard CF

映射四:全空间迭代 -> prefix迭代

底层先提供 whole-space iterator,再按 prefix 做过滤包装;
如果实现层知道某个 prefix 已经独占某个 CF,则可以直接用更高效的单 CF iterator。


BlueStore如何接入RocksDB

BlueStore 自己并不直接 new rocksdb::DB,而是通过 KeyValueDB 工厂来创建。

1. 创建入口

KeyValueDB.cc 的工厂非常简单:

1
2
3
4
5
6
7
8
9
10
KeyValueDB *KeyValueDB::create(CephContext *cct, const string& type,
const string& dir,
map<string,string> options,
void *p)
{
if (type == "rocksdb") {
return new RocksDBStore(cct, dir, options, p);
}
return NULL;
}

因此 BlueStore 只需要提供:

  • 类型字符串 "rocksdb"
  • 路径
  • KV 选项
  • 可选 Env 指针

就能拿到 RocksDBStore 实例。

2. BlueStore中的典型初始化动作

BlueStore 在初始化数据库时,通常会做三件事:

  1. 调用 KeyValueDB::create()
  2. 调用 db->set_merge_operator() 注册某些 prefix 的 merge 规则
  3. 调用 db->open()db->create_and_open()

这意味着:

  • KeyValueDB 负责统一入口
  • RocksDBStore 负责真实数据库打开
  • merge operator 必须在 open 前就注册完成

3. 运行期调用方式

BlueStore 后续主要通过以下接口使用 RocksDB:

  • get_transaction()
  • submit_transaction()
  • submit_transaction_sync()
  • get()
  • get_iterator()
  • compact_*()
  • estimate_prefix_size()

也就是说,BlueStore 使用的是“数据库抽象”,而非 RocksDB 原生 API。


数据库打开与初始化流程

RocksDBStore 的打开主流程核心在 do_open()

1. init阶段

init() 做的事情很克制:

  • 保存 options 字符串
  • 试解析一遍配置,尽早发现错误

不真正打开数据库

这样 BlueStore 在 mount 前就能先验证 RocksDB 参数是否合法。

2. load_rocksdb_options阶段

load_rocksdb_options() 是选项构造中心,负责:

  • 解析文本配置串
  • 安装 Ceph 日志器
  • 绑定 Env
  • 配置 block cache / row cache
  • 配置 bloom filter / index type / table factory
  • 安装默认 CF 的 merge router
  • 配置 RocksDB statistics

可以理解为:这里把 Ceph 配置世界转换成 RocksDB Options 世界。

3. create_and_open路径

如果数据库尚不存在,则:

  1. 创建目录
  2. 打开 RocksDB
  3. 应用 column family / sharding 布局
  4. 持久化 sharding/def

4. open已有数据库路径

如果是打开已有数据库,则:

  1. 读取持久化的 sharding 定义
  2. 校验 RocksDB 实际列族与目标布局是否一致
  3. 打开已有列族
  4. 如处于重分片恢复阶段,则补建缺失列族
  5. 建立运行期 cf_handles 路由表

这里最关键的一点是:

Ceph 在 reopen 时不是“盲开 RocksDB”,而是先验证逻辑分片布局与物理列族布局是否匹配。


写路径分析

1. 从BlueStore到事务对象

BlueStore 需要写元数据时,先拿一个事务:

1
2
3
4
BlueStore
-> db->get_transaction()
-> RocksDBStore::RocksDBTransactionImpl
-> rocksdb::WriteBatch bat

此时所有写操作只是追加到 WriteBatch 里。

2. set路径

调用 set(prefix, key, value) 时,RocksDBTransactionImpl::set() 会先判断:

  • 这个 prefix 是否有独立 CF
  • 如果有,进一步判断是否要按 shard 分片

然后有两种写法:

情况一:prefix 落在 default CF

把 key 编码成:

prefix + '\0' + key

然后写入 default_cf

情况二:prefix 落在独立CF

直接把裸 key 写入目标 CF

3. 删除路径

删除有三种典型方式:

  • rmkey()
  • rmkeys_by_prefix()
  • rm_range_keys()

其中 rmkeys_by_prefix()rm_range_keys() 里有一个非常重要的优化:

  • 如果删除键数较少,则逐 key Delete
  • 如果删除范围很大,则改用 DeleteRange

这样做的目的是在:

  • 精确性
  • 写放大
  • tombstone 数量

之间做平衡。

4. merge路径

merge() 不会立即读出旧值再合并,而是把 merge 操作写入 RocksDB,由 merge operator 在后续读取或 compaction 时解释。

这适合:

  • 计数器
  • 统计项
  • 增量聚合类元数据

5. 提交路径

所有事务最后都进入:

  • submit_transaction()
  • submit_transaction_sync()

它们最终都会调用:

  • submit_common()

submit_common() 负责:

  1. 设置 rocksdb::WriteOptions
  2. 根据配置决定是否禁用 WAL
  3. 打印事务内容用于调试
  4. 调用 db->Write(woptions, &bat)
  5. 采集 RocksDB perf context 数据

因此,Ceph 事务和 RocksDB WriteBatch 的关系可以概括为:

一个 Ceph KV 事务,最终就是一次 RocksDB db->Write()


读路径与迭代器分析

1. get路径

RocksDBStore::get() 的主要工作依然是路由:

  1. 判断 prefix 是否映射独立 CF
  2. 若需要,则进一步选择 shard
  3. 调用 RocksDB Get

同样分两种情况:

  • default CF:读 "prefix\0key"
  • 独立 CF:读裸 key

2. split_key的意义

默认 CF 中,原始键是组合键:

prefix\0key

因此 split_key() 是一个很基础但非常重要的辅助函数。

它的作用是把 RocksDB 原始 key 拆回:

  • prefix
  • 逻辑 key

很多迭代器逻辑都依赖它。

3. WholeSpaceIterator

WholeSpaceIteratorImpl 表示底层数据库实现直接暴露的原始迭代器。

在 RocksDB 场景下,它通常对应:

  • 某个 column family 上的 rocksdb::Iterator

它能返回:

  • key()
  • raw_key()
  • value()
  • value_as_sv()

这里的 raw_key() 特别重要,因为它可能是未过滤 prefix 的真实编码键。

4. PrefixIterator

KeyValueDB 默认提供一个 PrefixIteratorImpl 包装器,它的作用是:

  • 底层先拿 whole-space iterator
  • 外层再按 prefix 过滤

这保证即使某个后端没有实现更高级的原生 prefix iterator,也能对上层暴露统一行为。

5. RocksDBStore中的优化

RocksDBStore 不只是使用通用 PrefixIteratorImpl,它还会根据:

  • prefix 是否独占某个列族
  • 迭代边界是否落在单个 shard

来决定是否直接创建更高效的单列族 iterator。

这就是为什么 RocksDBStore 里会有:

  • new_shard_iterator()
  • check_cf_handle_bounds()
  • 跨 shard 的 merge iterator

列族与分片机制

这是 Ceph 对 RocksDB 使用方式里最关键、也最“非通用 RocksDB 教材式”的部分。

1. 为什么需要列族

如果所有 prefix 都放在 default CF 中,会遇到几个问题:

  • 不同类型元数据混在一起
  • compaction 相互影响
  • 缓存难以差异化
  • 热点前缀容易形成瓶颈

因此 Ceph 支持把某些 prefix 提升为独立 column family。

2. 为什么还需要分片

即使一个 prefix 独占一个 CF,仍可能有问题:

  • 某个前缀的数据量太大
  • 某个前缀写入太热
  • 单 CF compaction 压力太集中

因此 Ceph 允许一个逻辑列族继续拆成多个 shard:

1
O -> O-0, O-1, O-2, O-3 ...

3. 分片定义字符串

Ceph 使用文本字符串描述 sharding 布局,例如:

1
O(4,0-10)=write_buffer_size=...

它描述的是:

  • 逻辑列族名
  • shard 数量
  • 用 key 的哪一段参与 hash
  • 该列族附加的 RocksDB 选项

parse_sharding_def() 就是这个语法的解析器。

4. 路由规则

路由过程是:

  1. 通过 prefix 找到逻辑列族
  2. 如果只有一个 handle,直接落该 CF
  3. 如果有多个 shard,则取 key 的某个子串做 hash
  4. hash % shard_count 选出最终 shard

5. 运行期持久化

Ceph 会把当前 sharding 定义写到:

  • sharding/def

这样下次重启时可以:

  • 校验 RocksDB 中真实存在的列族
  • 恢复重分片状态
  • 拒绝错误布局的数据库

BlueStore各Prefix的Key/Value格式总表

下面这张表总结的是 BlueStore 传给 KeyValueDB 的逻辑格式

真正落到 RocksDB 时,RocksDBStore 还会再做一层物理映射:

  • 若该 prefix 落在 default CF,则物理 key = prefix + '\0' + logical_key
  • 若该 prefix 有独立 CF,则物理 key = logical_key

因此表中重点关注的是 BlueStore 自己定义的 logical_keyvalue 编码规则。

Prefix 主要用途 logical key 格式 value 格式 说明
S super 元数据 通常是可读字段名,如 nid_maxblobid_maxondisk_formatmin_alloc_size 对应字段的 encode() 后 bufferlist 保存 BlueStore 全局配置、版本、计数器、高水位等
T 统计信息 可读字段名,或 pool_id 的二进制编码 key 统计结构编码后的 bufferlist,常通过 merge 更新 BLUESTORE_GLOBAL_STATFS_KEY 也存放在这里
C collection 元数据 stringify(cid) cnode_t 的编码 保存 collection 的 bits 等元数据
O onode / 对象元数据 ghobject_t 的排序友好二进制编码:shard + pool + hash + namespace + key/name + snap + generation + 'o' bluestore_onode_t + spanning blobs + inline extents 编码后的 bufferlist BlueStore 最核心的一类记录
M 普通 omap u64 nid + '.' + user_key;边界键分别用 '-''~' header/tail 特殊值,普通项为用户 omap value 默认对象 omap
P pgmeta omap u64 nid + '.' + user_key;不带 pool/hash 前缀 M 相同 用于 meta collection 的 omap
m per-pool omap u64 pool + u64 nid + '.' + user_key M 相同 把 pool 维度编码进 key 前缀
p per-pg omap u64 pool + u32 hash + u64 nid + '.' + user_key M 相同 比 per-pool 再多一层 hash 维度
L deferred 事务 u64 seq deferred_transaction_t 的编码 启动回放 deferred write 就扫这类键
X shared blob 元数据 u64 sbid bluestore_shared_blob_t 的编码 保存 shared blob 引用关系等元数据
B freelist 空间管理 u64 offset u64 length 经典 freelist 形式,表示某段空闲区间
b bitmap freelist 空间管理 BitmapFreelistManager 定义的位图键 位图相关结构编码 属于 bitmap allocator 的内部格式

1. S:super元数据

这类记录最直观,key 往往就是字段名,value 是对应字段的编码结果。例如:

  • nid_max
  • blobid_max
  • freelist_type
  • ondisk_format
  • min_compat_ondisk_format
  • min_alloc_size
  • per_pool_omap

它们的共同特点是:

  • key 可读
  • value 一般是定长标量或简单结构的 encode() 结果
  • 用于 BlueStore 启动时恢复全局状态

2. T:统计信息

T 类记录主要用于:

  • 全局 statfs
  • per-pool 统计

这里有两类常见 key:

  1. 可读字符串 key,例如:
    • bluestore_statfs
  2. 二进制 pool 统计 key:
    • u64(pool_id)

这类 value 常不是“整块重写”,而是通过 merge() 增量更新。

3. C:collection元数据

C 类记录较简单:

  • key = stringify(cid)
  • value = cnode_t

典型作用是保存:

  • collection 的 bits
  • collection 的基础拓扑元信息

4. O:对象onode元数据

O 是最重要的一类。

key 格式

对象 key 是按排序需求定制的二进制编码,核心顺序是:

  1. shard
  2. pool
  3. hash
  4. namespace
  5. key/object name
  6. snap
  7. generation
  8. 后缀字符 'o'

它的目标不是人可读,而是:

  • 能按对象逻辑顺序扫描
  • 能支持范围列举
  • 能支持 lower_bound/upper_bound

value 格式

O 类 value 由三部分拼接而成:

  1. bluestore_onode_t
  2. spanning blob 信息
  3. inline extent map 数据

如果 extent map 已经分片,则分片内容不一定都 inline 放在这个 value 里。

5. M/P/m/p:omap家族

omap 的关键点是:底层 key 不是直接等于用户 omap key

它会在用户 key 前面加上对象/池/PG 相关前缀,以便:

  • 把同一个对象的 omap 放在连续范围里
  • 构造 header/key/tail 三种边界

三类边界字符是:

  • '-':header
  • '.':实际用户 key
  • '~':tail

所以同一个对象的 omap 区间通常长这样:

1
2
3
... + nid + '-'
... + nid + '.' + user_key
... + nid + '~'

不同变种区别在于 nid 前面是否再编码:

  • M:只有 nid
  • P:pgmeta 场景,也只有 nid
  • mpool + nid
  • ppool + hash + nid

6. L:deferred事务

L 类记录表示尚未完全回放/清理的 deferred 写事务。

  • key = u64 seq
  • value = deferred_transaction_t

启动恢复时,BlueStore 会扫描 PREFIX_DEFERRED,逐条回放。

7. X:shared blob

X 类记录保存 shared blob 元数据。

  • key = u64 sbid
  • value = bluestore_shared_blob_t

value 里最关键的是:

  • shared blob 标识
  • ref_map
  • 引用计数相关状态

8. B:传统freelist

B 类记录表示空闲空间区间:

  • key = u64 offset
  • value = u64 length

它描述的语义很直接:

offset 开始,有一段长度为 length 的空闲物理空间。

9. b:bitmap freelist

b 类记录属于 bitmap allocator 的内部格式。

它和 B 的区别是:

  • B 偏向“区间列表”
  • b 偏向“位图块”

这里的 key/value 具体布局不在 BlueStore.cc 里定义,而由 BitmapFreelistManager 负责。
如果你继续往下看 allocator 路径,重点就该转到 bitmap freelist 的实现文件。


MergeOperator机制

1. 为什么要做两层适配

RocksDB 每个 column family 只能配置一个 merge operator。

但 Ceph 的逻辑模型是:

  • 不同 prefix 可能需要不同 merge 语义

如果所有这些 prefix 都落在 default CF,就会出现冲突:

  • default CF 里混着多个 prefix
  • 每个 prefix 却想用不同 merge 规则

因此 RocksDBStore 设计了两种适配器:

2. MergeOperatorRouter

用于 default CF。

它的逻辑是:

  1. 从原始 key 中识别 prefix
  2. 在已注册的 merge operator 列表中找到对应 prefix
  3. 把 merge 委托给对应 Ceph MergeOperator

3. MergeOperatorLinker

用于非 default CF。

因为独立 CF 对应单一逻辑 prefix,所以不需要“路由”,只需要直接绑定。

4. 为什么要构造稳定名称

RocksDB 在打开数据库时会校验 merge operator 名称是否一致。

因此 RocksDBStore 会把:

  • prefix
  • operator name

组合成稳定的整体名字,以保证:

  • 注册顺序变化不会影响最终名称
  • open 时能正确做一致性检查

BlueFS与Env的关系

1. Env是什么

rocksdb::Env 是 RocksDB 对底层文件环境的抽象。

RocksDB 的很多操作并不直接调用 POSIX 文件 API,而是通过 Env

  • 创建文件
  • 读写 WAL
  • 读写 SST
  • 创建目录
  • 后台线程

2. Ceph为什么要定制Env

在 BlueStore 场景里,RocksDB 文件并不一定落在普通文件系统目录中,而可能落在:

  • BlueFS 管理的设备空间

因此 BlueStore 会在合适条件下构造:

  • BlueRocksEnv

再把它作为 void* p 传给 KeyValueDB::create()

然后 RocksDBStore 在构造或 open 时把它转回:

  • rocksdb::Env *env

3. 结果是什么

RocksDB 仍以为自己在使用标准 Env
但实际底层文件已经由 BlueFS 承载。

这使得 Ceph 实现了:

  • RocksDB 逻辑不大改
  • BlueFS 接管底层存储
  • WAL/SST 与 BlueStore 设备布局协同工作

缓存、统计与compact

1. 缓存

RocksDBStore::load_rocksdb_options() 里会统一构造缓存体系:

  • block cache
  • row cache
  • bloom filter
  • index type

Ceph 还支持:

  • 全局 RocksDB cache 大小配置
  • 某些列族单独 block cache 配置
  • 优先级缓存接口 PriorityCache

这说明 Ceph 并不只是“用一下 RocksDB”,而是在缓存资源分配上做了比较深的整合。

2. 统计

RocksDBStore 支持:

  • RocksDB statistics
  • RocksDB perf context
  • Ceph perf counters

因此既能看:

  • RocksDB 内部写 WAL / memtable / delay 的耗时

也能看:

  • Ceph 视角的 get latency / submit latency / compact 指标

3. compact

Ceph 暴露了多种 compact 入口:

  • 全量 compact
  • prefix compact
  • range compact
  • async compact

RocksDBStore 还维护了一个后台 compact 队列,用来:

  • 合并相邻 compact 范围
  • 避免前台线程直接承担长时间 compact

reshard重分片流程

reshard 是 RocksDBStore 里最复杂的一类维护逻辑之一。

它的本质是:

把已有数据库中的 key,从旧的列族布局搬迁到新的列族布局。

1. 为什么需要reshard

因为随着系统运行,最初的列族布局可能不再适合:

  • 数据量变大
  • 某类前缀成为热点
  • compaction 压力集中
  • 需要新的缓存策略

2. reshard的核心步骤

  1. 解析目标 sharding 定义
  2. 写入重分片锁
  3. 列出现有列族
  4. 打开所有旧列族
  5. 计算缺失的新列族并创建
  6. 构造新的 cf_handles
  7. 遍历旧列族全部 key
  8. 根据新规则重算目标 CF
  9. 批量执行 Delete + Put
  10. 删除旧布局中不再需要的空列族
  11. 写回新的 sharding/def

3. 为什么实现里有很多“失败注入”开关

resharding_ctrl 里有:

  • unittest_fail_after_first_batch
  • unittest_fail_after_processing_column
  • unittest_fail_after_successful_processing

这是因为重分片是一个跨多步的维护过程,最怕中途失败导致状态不一致。

所以实现必须支持:

  • 故障注入
  • 启动后恢复
  • 继续完成中断的重分片

典型调用链总结

1. 初始化调用链

1
2
3
4
5
6
7
BlueStore
-> KeyValueDB::create("rocksdb", ...)
-> new RocksDBStore(...)
-> set_merge_operator(...)
-> open() / create_and_open()
-> do_open()
-> rocksdb::DB::Open(...)

2. 写事务调用链

1
2
3
4
5
6
7
BlueStore
-> db->get_transaction()
-> tx->set()/rmkey()/merge()
-> RocksDBTransactionImpl::bat 累积写操作
-> db->submit_transaction()
-> RocksDBStore::submit_common()
-> rocksdb::DB::Write()

3. 单key读取调用链

1
2
3
4
5
BlueStore
-> db->get(prefix, key, &value)
-> RocksDBStore::get()
-> 判断 prefix 对应的 CF / shard
-> rocksdb::DB::Get()

4. 迭代器调用链

1
2
3
4
5
BlueStore
-> db->get_iterator(prefix)
-> RocksDBStore::get_iterator()
-> 选择 whole-space / shard iterator / merge iterator
-> rocksdb::Iterator

5. merge调用链

1
2
3
4
5
6
7
BlueStore
-> db->set_merge_operator(prefix, op)
-> RocksDBStore open 前安装 merge operator
-> tx->merge(prefix, key, value)
-> RocksDB WriteBatch::Merge()
-> Router/Linker
-> 具体 Ceph MergeOperator

关键结论

1. KeyValueDB是Ceph真正依赖的接口

Ceph 上层并不直接依赖 RocksDB 原生 API,而是依赖 KeyValueDB 抽象。

2. RocksDBStore是关键桥梁

RocksDBStore 负责把:

  • Ceph 的 prefix/key 模型
  • 事务模型
  • 迭代器模型
  • merge 模型

映射到 RocksDB 的对象体系中。

3. Ceph对RocksDB的使用并不“浅”

Ceph 并不是简单把 RocksDB 当成一个黑盒 KV 库,而是深度利用了:

  • column family
  • merge operator
  • Env
  • WriteBatch
  • iterator bounds
  • perf context
  • block cache
  • DeleteRange
  • 重分片

4. BlueFS是RocksDB接入BlueStore设备体系的关键

没有 BlueFS / BlueRocksEnv,RocksDB 只能把文件落到普通文件系统;
有了它,Ceph 才能把 RocksDB WAL/SST 纳入 BlueStore 的设备管理体系。

5. 读这部分代码的最佳顺序

建议按下面顺序阅读:

  1. KeyValueDB.h
  2. KeyValueDB.cc
  3. RocksDBStore.h
  4. RocksDBStore.cc
  5. BlueStore.cc 中数据库初始化、读写事务、iterator 调用部分

如果按这个顺序看,理解成本会明显低很多。

BlueStore对象存储与元数据管理分析

BlueStore对象存储与元数据管理分析

1. 概述

BlueStore是Ceph的默认对象存储后端,采用了一种创新的存储架构:

  • 数据直接存储在块设备上:避免传统文件系统的开销
  • 元数据存储在RocksDB中:提供高性能的元数据管理
  • 分离式设计:数据路径和元数据路径完全分离

2. 对象存储架构

2.1 存储层次结构

BlueStore采用多层次的存储抽象,从逻辑对象到物理存储的映射关系如下:

1
2
3
4
5
6
7
8
9
Object (Onode)

Extent (逻辑extent)

Blob (数据块)

PExtent (物理extent)

Block Device (块设备)

2.1.1 Onode (Object Node)

定义位置: BlueStore.h:1328

Onode是对象的内存表示,包含:

  • 对象标识符 (ghobject_t oid): 唯一标识一个对象
  • 对象元数据 (bluestore_onode_t onode): 存储在RocksDB中的元数据
  • ExtentMap (extent_map): 逻辑extent到Blob的映射
  • BufferSpace (bc): 缓存的数据缓冲区
1
2
3
4
5
6
7
struct Onode {
Collection *c; // 所属集合
ghobject_t oid; // 对象ID
bluestore_onode_t onode; // 元数据
ExtentMap extent_map; // extent映射
BufferSpace bc; // 缓冲区空间
};

元数据存储

  • Onode的元数据存储在RocksDB中,键为 PREFIX_OBJ + key
  • 包含对象大小、修改时间、扩展属性等信息

2.1.2 Extent (逻辑Extent)

定义位置: BlueStore.h:829

Extent将对象的逻辑偏移映射到Blob:

1
2
3
4
5
6
struct Extent {
uint32_t logical_offset; // 对象内的逻辑偏移
uint32_t blob_offset; // Blob内的偏移
uint32_t length; // 长度
BlobRef blob; // 指向的Blob
};

特点

  • Extent按logical_offset排序,存储在boost::intrusive::set
  • 一个对象可以有多个Extent,支持稀疏对象
  • Extent可以跨越多个Blob(通过spanning blob机制)

2.1.3 Blob (数据块)

定义位置: BlueStore.h:651bluestore_types.h:498

Blob是数据的逻辑容器,包含:

1
2
3
4
5
struct Blob {
bluestore_blob_t blob; // Blob元数据
bluestore_blob_use_tracker_t used_in_blob; // 引用计数跟踪
SharedBlobRef shared_blob; // 共享Blob引用(如果共享)
};

Blob元数据 (bluestore_blob_t):

1
2
3
4
5
6
7
8
struct bluestore_blob_t {
PExtentVector extents; // 物理extent列表
uint32_t logical_length; // 逻辑长度(压缩前)
uint32_t compressed_length; // 压缩后长度
uint32_t flags; // 标志位(压缩、校验和等)
uint8_t csum_type; // 校验和类型
buffer::ptr csum_data; // 校验和数据
};

Blob特性

  • 压缩支持: 可以存储压缩数据,通过FLAG_COMPRESSED标志
  • 校验和: 支持数据完整性校验,通过FLAG_CSUM标志
  • 共享Blob: 多个对象可以共享同一个Blob(写时复制)
  • 引用计数: 通过bluestore_blob_use_tracker_t跟踪使用情况

2.1.4 PExtent (物理Extent)

定义位置: bluestore_types.h:98

PExtent表示块设备上的实际物理位置:

1
2
3
4
struct bluestore_pextent_t {
uint64_t offset; // 块设备上的偏移
uint32_t length; // 长度
};

特点

  • 一个Blob可以包含多个PExtent(支持不连续存储)
  • PExtent由Allocator分配和管理
  • 支持对齐到块大小边界

2.2 数据写入流程

2.2.1 写入路径

  1. 接收写入请求

    • 客户端请求写入对象数据
    • BlueStore创建或获取Onode
  2. 分配空间

    • 通过Allocator分配物理空间(PExtent)
    • 创建或扩展Blob来容纳数据
  3. 数据准备

    • 可选:压缩数据(如果启用)
    • 可选:计算校验和
    • 准备写入缓冲区
  4. 写入块设备

    • 直接写入到块设备(绕过文件系统)
    • 使用异步I/O提高性能
  5. 更新元数据

    • 创建或更新Extent映射
    • 更新Blob元数据
    • 更新Onode元数据
  6. 提交事务

    • 将元数据变更写入RocksDB事务
    • 提交事务,确保一致性

2.2.2 关键代码路径

写入入口: BlueStore::_do_write()

  • 创建WriteContext管理写入操作
  • 调用Writer::do_write()执行实际写入

Writer类: Writer.hWriter.cc

  • 管理写入的数据块(blob_data_t)
  • 处理压缩、对齐等操作
  • 协调空间分配和数据写入

2.3 数据读取流程

  1. 查找Onode

    • 从RocksDB读取对象元数据
    • 加载到内存缓存
  2. 解析ExtentMap

    • 根据逻辑偏移查找对应的Extent
    • 确定需要读取的Blob和偏移
  3. 读取数据

    • 从块设备读取物理数据
    • 可选:验证校验和
    • 可选:解压缩数据
  4. 返回数据

    • 将数据返回给客户端
    • 可能缓存到BufferSpace中

3. 元数据管理

3.1 元数据存储架构

BlueStore使用RocksDB作为元数据存储引擎:

1
KeyValueDB *db = nullptr;  // RocksDB实例

3.1.1 元数据类型

1. 对象元数据 (Onode)

  • 键格式: PREFIX_OBJ + key
  • 值内容: 序列化的bluestore_onode_t
  • 包含信息:
    • 对象大小
    • 修改时间
    • 扩展属性
    • ExtentMap的编码(如果较小则内联)

2. Blob元数据

  • 存储位置: 内联在Onode中(小对象)或单独存储(大对象)
  • SharedBlob: 共享Blob的元数据单独存储
    • 键格式: PREFIX_SHARED_BLOB + sbid
    • 包含引用计数、物理extent等信息

3. 集合元数据 (Collection)

  • 键格式: PREFIX_COLL + cid
  • 值内容: 集合的元数据(如bits等)

4. OMap数据

  • 键格式: omap_prefix + onode_key + user_key
  • 值内容: 用户定义的键值对

3.2 元数据组织

3.2.1 ExtentMap编码

ExtentMap可以以两种方式存储:

1. 内联存储 (小对象)

  • ExtentMap直接编码在Onode中
  • 适合extent数量少的对象

2. 分片存储 (大对象)

  • ExtentMap分成多个shard
  • 每个shard独立编码和存储
  • 支持按需加载,减少内存占用

编码格式:

1
2
3
4
5
// Extent编码包含:
// - logical_offset (逻辑偏移)
// - blob_offset (Blob偏移)
// - length (长度)
// - blob_id (Blob ID或SharedBlob ID)

3.2.2 SharedBlob管理

目的: 支持写时复制(Copy-on-Write),多个对象可以共享数据

结构:

1
2
3
4
5
struct SharedBlob {
uint64_t sbid; // SharedBlob ID
bluestore_extent_ref_map_t ref_map; // 引用计数映射
bluestore_blob_t blob; // Blob元数据
};

引用计数:

  • 使用bluestore_extent_ref_map_t跟踪每个extent的引用数
  • 当引用计数为0时,可以释放物理空间

3.3 事务管理

3.3.1 TransContext

定义位置: BlueStore.h:276

每个事务由TransContext管理:

1
2
3
4
struct TransContext {
KeyValueDB::Transaction t; // RocksDB事务
// ... 其他事务相关数据
};

3.3.2 事务流程

  1. 开始事务

    • 创建TransContext
    • 创建RocksDB事务
  2. 记录变更

    • 修改Onode元数据
    • 更新ExtentMap
    • 记录Blob变更
  3. 提交事务

    • 将元数据变更写入RocksDB
    • 确保原子性
  4. 同步

    • 可选:同步RocksDB(fsync)
    • 确保数据持久化

3.4 缓存管理

3.4.1 Onode缓存

结构: OnodeCacheShard

  • 缓存Onode对象
  • 使用LRU策略管理
  • 支持分片(sharding)提高并发性

3.4.2 Buffer缓存

结构: BufferCacheShard

  • 缓存对象数据缓冲区
  • 管理Buffer对象的状态(EMPTY/CLEAN/WRITING)
  • 支持延迟写入(deferred write)

3.4.3 缓存策略

  • LRU: 最近最少使用
  • TwoQ: 两级队列(热数据和冷数据)
  • 优先级缓存: 根据访问模式调整优先级

4. 空间管理

4.1 Allocator (分配器)

BlueStore使用Allocator管理块设备上的空闲空间:

支持的分配器类型:

  • StupidAllocator: 简单的interval_set实现
  • BitmapAllocator: 基于位图的分配器
  • AvlAllocator: 基于AVL树的分配器
  • BtreeAllocator: 基于B树的分配器
  • Btree2Allocator: B树分配器的改进版
  • HybridAllocator: 混合分配器(主分配器+位图分配器)

分配流程:

  1. 根据请求大小和提示(hint)选择分配策略
  2. 从空闲空间中找到合适的extent
  3. 返回分配的PExtent列表

4.2 FreelistManager (空闲列表管理器)

作用: 持久化空闲空间信息到RocksDB

实现: BitmapFreelistManager

  • 使用位图跟踪空闲块
  • 支持合并操作符(XOR)进行增量更新
  • 在RocksDB中持久化空闲空间状态

5. 高级特性

5.1 压缩

支持: 通过Compression.hCompression.cc实现

压缩流程:

  1. Estimator: 估算压缩收益
  2. Scanner: 扫描可压缩的数据
  3. 压缩执行: 使用配置的压缩算法
  4. 元数据更新: 更新Blob的压缩标志和长度

压缩标志: bluestore_blob_t::FLAG_COMPRESSED

5.2 校验和

支持: 通过Checksummer实现

校验和类型:

  • CRC32C
  • XXHASH32
  • XXHASH64

存储: 校验和数据存储在bluestore_blob_t::csum_data

5.3 共享Blob (写时复制)

场景:

  • 对象克隆
  • 快照
  • 去重

机制:

  • 多个对象共享同一个Blob
  • 写入时创建新的Blob(写时复制)
  • 通过引用计数管理生命周期

6. 性能优化

6.1 延迟写入 (Deferred Write)

目的: 合并小写入,提高性能

机制:

  • 小写入先缓存到DeferredBatch
  • 延迟到合适时机批量写入
  • 减少I/O次数

6.2 对齐优化

块对齐:

  • 数据对齐到块大小边界
  • 提高I/O效率

分配对齐:

  • 根据分配单元对齐
  • 减少碎片

6.3 缓存优化

多级缓存:

  • Onode缓存(元数据)
  • Buffer缓存(数据)
  • 块设备缓存(操作系统级)

缓存策略:

  • 根据访问模式调整
  • 支持预热和预取

7. 数据一致性

7.1 事务保证

  • 原子性: 通过RocksDB事务保证
  • 持久性: 通过fsync保证
  • 一致性: 通过写前日志(WAL)保证

7.2 崩溃恢复

  • RocksDB恢复: 自动恢复未提交的事务
  • 空间一致性: 通过FreelistManager重建空闲空间
  • 元数据校验: 启动时验证元数据完整性

8. 总结

BlueStore的存储架构具有以下特点:

  1. 分离式设计: 数据和元数据分离存储,各自优化
  2. 直接I/O: 数据直接写入块设备,避免文件系统开销
  3. 灵活映射: 多层次的映射关系,支持稀疏对象和共享数据
  4. 高性能: 通过缓存、延迟写入、对齐等优化提高性能
  5. 可扩展: 支持压缩、校验和、共享Blob等高级特性

这种设计使得BlueStore能够提供高性能、高可靠性的对象存储服务,特别适合大规模分布式存储场景。

9. 关键数据结构关系图

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
┌─────────────────┐
│ Object (Onode)│
│ - oid │
│ - onode (meta) │
│ - extent_map │
│ - buffer_space │
└────────┬────────┘

│ contains

┌─────────────────┐
│ Extent │
│ - logical_off │
│ - blob_offset │
│ - length │
│ - blob_ref │
└────────┬────────┘

│ references

┌─────────────────┐
│ Blob │
│ - blob (meta) │
│ - use_tracker │
│ - shared_blob │
└────────┬────────┘

│ contains

┌─────────────────┐
│ PExtent │
│ - offset │
│ - length │
└────────┬────────┘

│ maps to

┌─────────────────┐
│ Block Device │
└─────────────────┘

Metadata Storage (RocksDB):
┌─────────────────┐
│ Onode Metadata │───┐
│ (PREFIX_OBJ) │ │
└─────────────────┘ │

┌─────────────────┐ │
│ SharedBlob Meta │ │
│ (PREFIX_SHARED) │ │
└─────────────────┘ │

┌─────────────────┐ │
│ Collection Meta │ │
│ (PREFIX_COLL) │ │
└─────────────────┘ │

┌─────────────────┐ │
│ OMap Data │ │
│ (omap_prefix) │ │
└─────────────────┘ │


┌─────────────────┐
│ RocksDB │
└─────────────────┘

10. 参考代码位置

  • 核心类定义: cephMain/src/os/bluestore/BlueStore.h
  • 核心实现: cephMain/src/os/bluestore/BlueStore.cc
  • 类型定义: cephMain/src/os/bluestore/bluestore_types.h
  • 写入器: cephMain/src/os/bluestore/Writer.h/cc
  • 压缩: cephMain/src/os/bluestore/Compression.h/cc
  • 分配器: cephMain/src/os/bluestore/Allocator*.h/cc
  • 空闲列表: cephMain/src/os/bluestore/FreelistManager.h/cc

BlueStore 流程图

下面用 Mermaid 画两张图:一张偏 模块架构,一张偏 写事务状态机 + 读路径流程(内容对应 BlueStore.cc 里的主要协作关系,而非逐函数罗列)。


1. BlueStore 架构图(组件与依赖)

flowchart TB
  subgraph OSD["OSD / ObjectStore 接口"]
    OS[queue_transactions / read / collection ops]
  end

  subgraph BS["BlueStore 核心(BlueStore.cc)"]
    COLL[Collection + OnodeSpace]
    ONODE[Onode / ExtentMap / Blob / SharedBlob]
    BUF[BufferSpace + BufferCacheShard]
    TXC[TransContext 状态机]
    TH[Throttle + OpSequencer]
    DEF[Deferred 队列与回放]
  end

  subgraph ALLOC["空间管理"]
    AL[Allocator]
    FM[FreelistManager]
  end

  subgraph PERSIST["持久化"]
    KV[(KeyValueDB / RocksDB)]
    BF[BlueFS]
    BDEV[(BlockDevice: block / block.db / block.wal)]
  end

  subgraph BG["后台与回调"]
    KVTH[kv_sync_thread / kv_finalize_thread]
    FIN[Finisher]
    MPT[mempool_thread]
  end

  OS --> COLL
  OS --> TXC
  COLL --> ONODE
  read_path[read → _do_read] --> ONODE
  read_path --> BUF
  TXC --> ONODE
  TXC --> BUF
  TXC --> TH
  TXC --> AL
  TXC --> FM
  TXC --> KV
  ONODE --> KV
  FM --> KV
  TXC --> BDEV
  BUF --> BDEV
  KV --> BF
  BF --> BDEV
  KVTH --> TXC
  TXC --> FIN
  DEF --> BDEV
  DEF --> KV
  MPT -.-> BS

要点(和 BlueStore.cc 一致):

  • 元数据 / omap / deferred / 分配元数据KV;启用 BlueFS 时 RocksDB 文件BlueFS 里,再落到 块设备
  • 对象字节 直接 Allocator → pextent → BlockDevice
  • TransContext 串起 AIO、KV 提交、deferred、on_commit

2. 流程图:写路径(queue_transactions + _txc_state_proc

stateDiagram-v2
  direction LR

  [*] --> PREPARE: queue_transactions\n_txc_add_transaction
  PREPARE --> AIO_WAIT: 有块设备 AIO\n_txc_aio_submit
  PREPARE --> IO_DONE: 无 AIO

  AIO_WAIT --> IO_DONE: AIO 完成\n_txc_finish_io\n(同 OSR 顺序推进)

  IO_DONE --> KV_QUEUED: 入 kv_queue\n通知 kv_sync_thread

  state KV_QUEUED {
      [*] --> Note
      Note: 可选 bluestore_sync_submit_transaction\n同线程 _txc_apply_kv
  }

  KV_QUEUED --> KV_SUBMITTED: db->submit_transaction\n_txc_apply_kv
  KV_SUBMITTED --> KV_DONE: sync / finalize\n_txc_committed_kv\n(on_commit → finisher)

  KV_DONE --> DEFERRED_QUEUED: deferred_txn 非空\n_deferred_queue
  KV_DONE --> FINISHING: 无 deferred

  DEFERRED_QUEUED --> DEFERRED_CLEANUP: 延迟写完成
  DEFERRED_CLEANUP --> FINISHING

  FINISHING --> DONE: _txc_finish\n释放/回收/出队
  DONE --> [*]

补充一条线性泳道(便于和代码对照):

写路径主链解析 Transaction改 Onode/Extent/Blob\n分配/释放 extent提交块写 AIO编码 onode/shared_blob\n写入 txc->tdeferred 记入 KV\n_finalize_kv_txc_state_procKV 持久化deferred?deferred 数据写_txc_finish

3. 流程图:读路径(read_do_read

无对象readCollection 锁 + get_onode返回 -ENOENT_do_read_read_cache\n划分 ready vs 待读BufferSpace 命中\n拼 bl按 blob 合并读请求\nAIO 读 BlockDevice校验和 / 解压bufferlist 返回

4. 流程图:挂载时与存储栈的打开顺序(_open_db_and_around

mount → _mount_open_db_and_aroundpath + fsid + lock_open_bdev_open_db 只读_open_super_meta_open_fm + _init_alloc关 DB 再 _open_db 最终模式_post_init_alloc 可写时_upgrade_super + _open_collections_kv_start + _deferred_replay

若需要 可导出 PNG/SVG:把上述 Mermaid 粘进支持 Mermaid 的渲染器(如部分 Markdown 预览、Mermaid Live Editor),即可导出。若你希望 单张合并大图(架构 + 写状态机一页),可以说一下更偏向「部署视角」还是「代码函数视角」,我可以再收紧节点命名以便贴进设计文档。

Ceph OSD EC(纠删码)实现过程与原理分析

Ceph OSD EC(纠删码)实现过程与原理分析

1. 概述

EC(Erasure Code,纠删码)是 Ceph 中一种重要的数据冗余方式,相比传统的多副本复制,EC 可以提供更高的存储效率。例如,使用 k=4, m=2 的 EC 配置(即 4+2),可以将 6 个数据块编码为 4 个数据块和 2 个校验块,在保证可以容忍 2 个块丢失的同时,存储效率为 66.7%(相比 3 副本的 33.3%)。

1.1 核心概念

  • k(数据块数):原始数据被分割成的数据块数量
  • m(校验块数):通过编码生成的校验块数量
  • 条带(Stripe):数据被分割和编码的基本单位
  • 分片(Shard):每个条带块存储在不同的 OSD 上,称为一个分片
  • 条带宽度(Stripe Width):一个条带的总大小 = k * chunk_size

1.2 架构设计

EC 实现采用了分层架构:

1
2
3
4
5
6
7
8
9
PrimaryLogPG (上层)

PGBackend (接口层)

ECBackend (EC 实现层)

ECCommon (通用功能层)

ErasureCodeInterface (编码算法接口)

2. 核心组件

2.1 PGBackend 接口

PGBackend 是后端抽象接口,定义了统一的接口规范:

  • Listener:由上层(PrimaryLogPG)实现,提供回调接口
  • RecoveryHandle:恢复操作句柄
  • 核心方法
    • submit_transaction():提交事务
    • objects_read_sync():同步读取对象
    • objects_read_and_reconstruct():异步读取并重构对象
    • recover_object():恢复对象

2.2 ECBackend

ECBackend 是 EC 实现的核心类,继承自 ECCommon

1
2
3
4
5
6
7
8
9
10
11
12
class ECBackend : public ECCommon {
// 核心组件
ReadPipeline read_pipeline; // 读取管道
RMWPipeline rmw_pipeline; // 读-修改-写管道
ECRecoveryBackend recovery_backend; // 恢复后端
ErasureCodeInterfaceRef ec_impl; // 编码算法实现

// 核心方法
void submit_transaction(...); // 提交事务
void objects_read_and_reconstruct(...); // 读取并重构
int recover_object(...); // 恢复对象
};

2.3 ECCommon

ECCommon 提供 EC 操作的通用接口和数据结构:

  • ec_extent_t:扩展结构,包含错误码、扩展映射和分片扩展映射
  • read_request_t:读取请求结构
  • read_result_t:读取结果结构
  • shard_read_t:分片读取结构

2.4 ECUtil

ECUtil 提供 EC 相关的工具函数和数据结构:

  • stripe_info_t:条带信息,包含 k、m、chunk_size 等
  • shard_extent_set_t:分片扩展集合
  • shard_extent_map_t:分片扩展映射
  • 对齐操作:EC_ALIGN_SIZE = 4KB,所有操作必须按 4KB 对齐

3. 写入流程

3.1 整体流程

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
客户端写入请求

PrimaryLogPG::do_op()

ECBackend::submit_transaction()

ECTransaction::generate_transactions() // 生成写入计划

RMWPipeline::start() // 读-修改-写管道

读取现有数据(如果需要)

编码数据

发送到各分片 OSD

等待所有分片确认

提交事务

3.2 写入计划(WritePlan)

ECTransaction::WritePlan 负责规划写入操作:

1
2
3
4
5
6
7
8
9
10
11
12
struct WritePlan {
bool want_read; // 是否需要读取
std::list<WritePlanObj> plans; // 写入计划列表
};

class WritePlanObj {
const hobject_t hoid; // 对象 ID
std::optional<shard_extent_set_t> to_read; // 需要读取的分片
shard_extent_set_t will_write; // 将要写入的分片
bool invalidates_cache; // 是否使缓存失效
bool do_parity_delta_write; // 是否进行校验增量写入
};

写入计划生成逻辑

  1. 分析写入范围:确定哪些条带需要更新
  2. 确定读取需求
    • 部分写入(Partial Write):需要读取现有数据
    • 全条带写入:不需要读取
  3. 确定写入分片
    • 数据分片(0 到 k-1):写入更新的数据块
    • 校验分片(k 到 k+m-1):写入重新计算的校验块

3.3 读-修改-写(RMW)流程

对于部分写入,需要执行 RMW 操作:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
1. 读取现有数据
- 确定需要读取的分片(通常是所有 k+m 个分片)
- 发送 ECSubRead 消息到各分片

2. 等待读取完成
- 收集各分片的响应
- 如果某些分片不可用,使用可用分片解码

3. 修改数据
- 将新数据覆盖到读取的数据中
- 重新计算受影响的条带

4. 编码数据
- 使用 ErasureCodeInterface::encode() 编码
- 生成 k 个数据块和 m 个校验块

5. 写入数据
- 发送 ECSubWrite 消息到各分片
- 等待所有分片确认

3.4 编码过程

编码使用 ErasureCodeInterface::encode()

1
2
3
4
5
6
7
8
9
// 伪代码
int encode(
const bufferlist &in, // 输入数据
map<int, bufferlist> &encoded // 输出:分片ID -> 编码后的数据块
) {
// 1. 将输入数据按 chunk_size 分割成 k 个数据块
// 2. 使用编码算法(如 Reed-Solomon)计算 m 个校验块
// 3. 返回 k+m 个编码块
}

3.5 分片写入

编码完成后,将数据块发送到对应的分片 OSD:

1
2
3
4
5
6
7
8
9
10
11
12
// ECBackend::handle_sub_write()
void handle_sub_write(
pg_shard_t from,
OpRequestRef msg,
ECSubWrite &op,
const ZTracer::Trace &trace,
ECListener &eclistener
) {
// 1. 验证写入请求
// 2. 应用事务到本地存储
// 3. 发送确认消息
}

4. 读取流程

4.1 整体流程

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
客户端读取请求

PrimaryLogPG::do_op()

ECBackend::objects_read_and_reconstruct()

ReadPipeline::start() // 读取管道

确定需要读取的分片

发送 ECSubRead 消息

等待响应(可能需要解码)

重构原始数据

返回给客户端

4.2 读取策略

EC 读取有两种策略:

  1. 快速读取(Fast Read)

    • 只从 k 个数据分片读取
    • 不需要解码,直接拼接数据
    • 适用于所有数据分片可用的情况
  2. 降级读取(Degraded Read)

    • 某些数据分片不可用
    • 需要从 k 个可用分片(可能包括校验分片)读取
    • 使用 ErasureCodeInterface::decode() 解码重构数据

4.3 解码过程

解码使用 ErasureCodeInterface::decode()

1
2
3
4
5
6
7
8
9
// 伪代码
int decode(
const map<int, bufferlist> &chunks, // 输入:可用分片的数据
map<int, bufferlist> &decoded // 输出:所有分片的数据
) {
// 1. 确定需要重构的分片
// 2. 使用解码算法重构缺失的分片
// 3. 返回完整的数据(k 个数据块)
}

4.4 读取管道(ReadPipeline)

ReadPipeline 管理异步读取流程:

1
2
3
4
5
6
7
class ReadPipeline {
// 1. 确定读取分片集合
// 2. 发送读取请求
// 3. 收集响应
// 4. 如果数据不完整,触发解码
// 5. 调用完成回调
};

5. 恢复流程

5.1 恢复触发

恢复在以下情况触发:

  1. Peering 完成后:发现某些分片缺失数据
  2. OSD 上线:需要回填数据
  3. 手动触发:管理员命令

5.2 恢复流程

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
检测到缺失对象

ECBackend::recover_object()

ECRecoveryBackend::recover_object()

确定恢复源(从哪些分片读取)

读取数据(使用 ReadPipeline)

解码重构数据

编码数据(如果需要写入到新分片)

推送到目标分片(PushOp)

等待确认

更新恢复状态

5.3 恢复源选择

恢复源选择遵循以下原则:

  1. 优先使用数据分片:如果 k 个数据分片都可用,直接读取
  2. 使用校验分片:如果某些数据分片不可用,使用校验分片解码
  3. 最小化网络传输:优先选择本地或网络距离近的分片

5.4 恢复后端(ECRecoveryBackend)

ECRecoveryBackend 专门处理恢复操作:

1
2
3
4
5
6
7
8
9
10
class ECRecoveryBackend : public RecoveryBackend {
// 恢复操作句柄
ECRecoveryHandle *open_recovery_op();

// 运行恢复操作
void run_recovery_op(ECRecoveryHandle &h, int priority);

// 处理恢复推送
void handle_recovery_push(const PushOp &op, ...);
};

6. 条带对齐与扩展管理

6.1 对齐要求

EC 操作必须按 4KB(EC_ALIGN_SIZE)对齐:

  • 原因:编码算法要求数据块大小固定
  • 实现:所有读写操作都对齐到 4KB 边界
  • 影响:部分写入可能需要读取整个条带

6.2 扩展(Extent)管理

EC 使用扩展(Extent)来管理数据范围:

1
2
3
4
5
6
7
8
9
10
11
// 扩展集合:表示连续的偏移量区间
using extent_set = interval_set<uint64_t, ...>;

// 扩展映射:将偏移量区间映射到缓冲区
using extent_map = interval_map<uint64_t, bufferlist, ...>;

// 分片扩展集合:每个分片的扩展集合
using shard_extent_set_t = shard_id_map<extent_set>;

// 分片扩展映射:每个分片的扩展映射
using shard_extent_map_t = shard_id_map<extent_map>;

6.3 扩展缓存(ExtentCache)

ECExtentCache 缓存已读取的扩展,避免重复读取:

1
2
3
4
5
6
7
8
9
10
11
12
13
class ECExtentCache {
// LRU 缓存
LRU &cache;

// 查询缓存
bool get(const hobject_t &hoid,
const extent_set &extents,
shard_extent_map_t *out);

// 更新缓存
void put(const hobject_t &hoid,
const shard_extent_map_t &data);
};

7. 部分写入优化

7.1 问题

传统 RMW 流程需要:

  1. 读取整个条带(k+m 个分片)
  2. 修改数据
  3. 重新编码
  4. 写入所有分片(k+m 个)

这会产生大量网络 I/O。

7.2 优化策略

Ceph 实现了多种优化:

  1. 校验增量写入(Parity Delta Write)

    • 只读取受影响的数据分片
    • 计算校验增量(delta)
    • 只更新校验分片
  2. 子块(Sub-chunk)优化

    • 对于大对象,可以按子块处理
    • 减少需要读取的数据量
  3. 零去重(Zero Deduplication)

    • 检测零缓冲区
    • 用共享的零缓冲区替换,减少内存使用

8. 消息类型

EC 使用专门的消息类型进行分片间通信:

8.1 ECSubWrite

分片写入消息:

1
2
3
4
5
6
7
struct ECSubWrite {
pg_shard_t from; // 来源分片
hobject_t soid; // 对象 ID
eversion_t version; // 版本
ObjectStore::Transaction txn; // 事务
// ...
};

8.2 ECSubRead

分片读取消息:

1
2
3
4
5
6
struct ECSubRead {
pg_shard_t from; // 来源分片
hobject_t soid; // 对象 ID
shard_read_t reads; // 读取请求
// ...
};

8.3 ECSubWriteReply / ECSubReadReply

对应的回复消息。

9. 关键数据结构

9.1 stripe_info_t

条带信息:

1
2
3
4
5
6
7
8
9
10
11
12
class stripe_info_t {
uint64_t stripe_width; // 条带宽度 = k * chunk_size
uint64_t chunk_size; // 块大小
uint32_t k; // 数据块数
uint32_t m; // 校验块数

// 计算条带号
uint64_t logical_to_stripe(uint64_t logical) const;

// 计算分片内的偏移量
uint64_t logical_to_prev_chunk_start(uint64_t logical) const;
};

9.2 shard_id_t

分片 ID 类型,用于标识不同的分片(0 到 k+m-1)。

9.3 ec_align_t

对齐区域:

1
2
3
4
5
struct ec_align_t {
uint64_t offset; // 偏移量(对齐到 4KB)
uint64_t length; // 长度(对齐到 4KB)
uint32_t flags; // 标志位
};

10. 性能优化

10.1 并行处理

  • 并行读取:同时从多个分片读取
  • 并行写入:同时写入多个分片
  • 异步操作:使用回调机制,避免阻塞

10.2 缓存策略

  • 扩展缓存:缓存已读取的扩展
  • 对象上下文缓存:缓存对象元数据
  • 零缓冲区共享:共享零缓冲区,减少内存

10.3 网络优化

  • 批量消息:合并多个小消息
  • 压缩:对大数据进行压缩
  • 优先级队列:区分恢复和客户端 I/O 的优先级

11. 错误处理

11.1 分片故障

当某些分片不可用时:

  1. 读取:使用可用分片解码重构数据
  2. 写入:如果写入分片不足,等待或使用临时分片
  3. 恢复:触发恢复流程,修复缺失数据

11.2 数据不一致

检测到数据不一致时:

  1. Scrub:定期检查数据完整性
  2. Repair:修复不一致的数据
  3. 日志记录:记录错误信息

12. 总结

EC 实现的核心要点:

  1. 分层架构:清晰的接口和实现分离
  2. 异步处理:使用管道和回调机制
  3. 对齐要求:所有操作按 4KB 对齐
  4. 优化策略:多种优化减少 I/O
  5. 容错能力:可以容忍 m 个分片故障

EC 相比多副本复制的优势:

  • 存储效率高:例如 4+2 配置效率为 66.7%,而 3 副本为 33.3%
  • 可配置性强:可以根据需求选择不同的 k 和 m
  • 适合大对象:对于大对象,EC 的优势更明显

EC 的劣势:

  • 计算开销:编码和解码需要 CPU 计算
  • 部分写入性能:部分写入需要 RMW,性能较差
  • 恢复复杂度:恢复流程比多副本复杂

13. 相关文件

核心文件

  • ECBackend.h/cc:EC 后端主实现
  • ECCommon.h/cc:EC 通用功能
  • ECUtil.h/cc:EC 工具函数
  • ECTransaction.h/cc:EC 事务处理
  • ECExtentCache.h/cc:扩展缓存

消息类型

  • ECMsgTypes.h/cc:EC 消息类型定义
  • MOSDECSubOpWrite.h:分片写入消息
  • MOSDECSubOpRead.h:分片读取消息

编码算法

  • erasure-code/ErasureCodeInterface.h:编码算法接口
  • 具体实现:jerasure、isa、shec 等插件

14. EC 与三副本(Replicated)详细对比

14.1 基本概念对比

特性 EC(纠删码) 三副本(Replicated)
数据组织 数据被编码成 k 个数据块 + m 个校验块 数据完整复制 3 份
存储位置 每个块存储在不同的 OSD(分片) 每个副本存储在不同的 OSD
容错能力 可以容忍 m 个分片丢失 可以容忍 2 个副本丢失
存储效率 高(例如 4+2 为 66.7%) 低(33.3%)
典型配置 k=4, m=2(4+2)或 k=8, m=3(8+3) size=3

14.2 架构实现对比

EC 架构

1
2
3
4
5
6
7
PrimaryLogPG

ECBackend (实现 PGBackend 接口)
├── ReadPipeline (读取管道)
├── RMWPipeline (读-修改-写管道)
├── ECRecoveryBackend (恢复后端)
└── ErasureCodeInterface (编码算法)

三副本架构

1
2
3
4
5
PrimaryLogPG

ReplicatedBackend (实现 PGBackend 接口)
├── Push/Pull 机制
└── 直接复制操作

14.3 写入流程对比

EC 写入流程

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
1. 客户端写入请求
2. 生成写入计划(WritePlan)
- 分析需要更新的条带
- 确定是否需要读取现有数据(RMW)
3. 如果需要 RMW:
- 读取现有数据(k+m 个分片)
- 解码重构原始数据
- 合并新数据
4. 编码数据
- 使用 ErasureCodeInterface::encode()
- 生成 k 个数据块 + m 个校验块
5. 并行写入所有分片(k+m 个)
- 发送 ECSubWrite 消息
6. 等待所有分片确认
7. 提交事务

特点

  • 需要编码计算(CPU 开销)
  • 必须写入所有 k+m 个分片
  • 部分写入需要 RMW,性能较差
  • 所有操作必须 4KB 对齐

三副本写入流程

1
2
3
4
5
6
7
1. 客户端写入请求
2. 主副本接收请求
3. 并行写入 3 个副本
- 主副本本地写入
- 发送 RepOp 消息到 2 个副本
4. 等待所有副本确认
5. 提交事务

特点

  • 无需编码计算
  • 只需写入 3 个副本
  • 部分写入直接覆盖,性能好
  • 无对齐要求

14.4 读取流程对比

EC 读取流程

1
2
3
4
5
6
7
8
9
10
1. 客户端读取请求
2. 确定读取策略:
- 快速读取:从 k 个数据分片读取(如果都可用)
- 降级读取:从 k 个可用分片读取(可能包括校验分片)
3. 发送 ECSubRead 消息
4. 收集响应
5. 如果需要解码:
- 使用 ErasureCodeInterface::decode()
- 重构原始数据
6. 返回给客户端

特点

  • 快速读取:只需读取 k 个分片,无需解码
  • 降级读取:需要解码计算
  • 最小读取分片数:k 个

三副本读取流程

1
2
3
1. 客户端读取请求
2. 主副本直接读取本地数据
3. 返回给客户端

特点

  • 直接读取,无需解码
  • 只需读取 1 个副本
  • 性能最优

14.5 恢复流程对比

EC 恢复流程

1
2
3
4
5
6
7
8
9
10
1. 检测到缺失分片
2. 确定恢复源(从哪些分片读取)
- 优先使用数据分片
- 如果数据分片不足,使用校验分片
3. 读取数据(使用 ReadPipeline)
- 从 k 个可用分片读取
4. 解码重构数据
5. 编码数据(如果需要写入到新分片)
6. 推送到目标分片(PushOp)
7. 等待确认

特点

  • 需要编码/解码计算
  • 恢复源选择复杂
  • 恢复流程复杂

三副本恢复流程

1
2
3
4
5
1. 检测到缺失副本
2. 从其他副本 Pull 数据
- 选择可用的副本作为源
3. 直接复制数据
4. 等待确认

特点

  • 无需编码/解码
  • 恢复源选择简单(任意可用副本)
  • 恢复流程简单

14.6 性能对比

操作类型 EC 三副本 说明
全条带写入 中等 EC 需要编码,但可以并行写入
部分写入 EC 需要 RMW,三副本直接覆盖
顺序读取 最快 EC 快速读取只需 k 个分片
随机读取 中等 EC 需要对齐,可能读取更多数据
降级读取 EC 需要解码计算
恢复速度 EC 需要编码/解码
CPU 开销 EC 需要编码/解码计算
网络开销 中等 EC 写入 k+m 个分片,三副本写入 3 个副本

14.7 存储效率对比

EC 存储效率计算

1
2
3
4
5
6
7
配置:k=4, m=2 (4+2)
存储效率 = k / (k + m) = 4 / 6 = 66.7%
容错能力:可以容忍 2 个分片丢失

配置:k=8, m=3 (8+3)
存储效率 = k / (k + m) = 8 / 11 = 72.7%
容错能力:可以容忍 3 个分片丢失

三副本存储效率

1
2
3
配置:size=3
存储效率 = 1 / 3 = 33.3%
容错能力:可以容忍 2 个副本丢失

14.8 适用场景对比

EC 适用场景

适合的场景

  • 大对象存储:对象越大,EC 优势越明显
  • 冷数据/归档数据:访问频率低,对性能要求不高
  • 存储成本敏感:需要更高的存储效率
  • 顺序读写为主:避免部分写入的 RMW 开销
  • 数据量巨大:存储效率带来的成本节省显著

不适合的场景

  • 小对象频繁写入:RMW 开销大
  • 随机写入为主:对齐要求导致额外 I/O
  • 对延迟敏感:编码/解码增加延迟
  • CPU 资源受限:编码/解码需要 CPU 计算

三副本适用场景

适合的场景

  • 高性能要求:需要低延迟、高吞吐
  • 小对象频繁写入:直接覆盖,无额外开销
  • 随机访问:无对齐要求,性能好
  • 热数据:访问频率高,性能优先
  • 简单运维:恢复流程简单

不适合的场景

  • 存储成本敏感:存储效率低
  • 大容量冷数据:存储成本高
  • 存储空间受限:需要更高的存储效率

14.9 代码实现对比

EC 关键代码路径

1
2
3
4
5
6
7
8
9
10
11
// 写入
ECBackend::submit_transaction()
→ ECTransaction::generate_transactions() // 生成写入计划
→ RMWPipeline::start() // RMW 流程
→ ErasureCodeInterface::encode() // 编码
handle_sub_write() // 分片写入

// 读取
ECBackend::objects_read_and_reconstruct()
→ ReadPipeline::start() // 读取管道
→ ErasureCodeInterface::decode() // 解码(如需要)

三副本关键代码路径

1
2
3
4
5
6
7
8
9
// 写入
ReplicatedBackend::submit_transaction()
→ 直接写入主副本
→ 发送 RepOp 到副本
→ 等待确认

// 读取
ReplicatedBackend::objects_read_sync()
→ 直接读取主副本

14.10 消息类型对比

EC 消息类型

  • ECSubWrite:分片写入消息
  • ECSubRead:分片读取消息
  • ECSubWriteReply:分片写入确认
  • ECSubReadReply:分片读取响应
  • PushOp:恢复推送操作

三副本消息类型

  • RepOp:复制操作消息
  • RepOpReply:复制操作确认
  • PushOp:恢复推送操作
  • PullOp:恢复拉取操作

14.11 数据一致性对比

特性 EC 三副本
写入一致性 需要所有 k+m 个分片确认 需要所有 3 个副本确认
读取一致性 从 k 个分片读取,可能解码 从主副本读取
部分写入 需要 RMW,可能影响一致性 直接覆盖,一致性简单
对齐要求 必须 4KB 对齐 无对齐要求

14.12 运维复杂度对比

方面 EC 三副本
配置复杂度 需要选择 k 和 m 只需设置 size=3
恢复复杂度 复杂(需要编码/解码) 简单(直接复制)
故障排查 复杂(涉及分片和解码) 简单(直接查看副本)
性能调优 复杂(涉及对齐、RMW 优化) 相对简单

14.13 成本分析

存储成本(以 100TB 原始数据为例)

EC (4+2)

  • 实际存储:100TB / 0.667 ≈ 150TB
  • 存储成本:150TB × 单价

三副本

  • 实际存储:100TB × 3 = 300TB
  • 存储成本:300TB × 单价

节省:EC 可以节省约 50% 的存储空间

计算成本

EC

  • CPU 开销:编码/解码需要 CPU 计算
  • 网络开销:写入 k+m 个分片

三副本

  • CPU 开销:几乎无额外计算
  • 网络开销:写入 3 个副本

14.14 总结对比表

维度 EC 三副本 胜者
存储效率 66.7% (4+2) 33.3% ✅ EC
写入性能(全条带) 中等 ✅ 三副本
写入性能(部分) ✅ 三副本
读取性能 最快 ✅ 三副本
CPU 开销 ✅ 三副本
恢复速度 ✅ 三副本
运维复杂度 复杂 简单 ✅ 三副本
适用对象大小 大对象 任意 -
成本 ✅ EC

14.15 选择建议

选择 EC 的情况

  1. 存储成本是主要考虑因素
  2. 数据主要是大对象
  3. 以顺序读写为主
  4. 对延迟不敏感
  5. 有足够的 CPU 资源

选择三副本的情况

  1. 性能是主要考虑因素
  2. 数据主要是小对象
  3. 随机访问较多
  4. 对延迟敏感
  5. CPU 资源受限
  6. 需要简单的运维

混合使用

  • 热数据使用三副本(性能优先)
  • 冷数据使用 EC(成本优先)
  • 通过 Ceph 的 tiering 功能自动迁移

OSD.cc 快速参考卡片

OSD.cc 快速参考卡片

🚀 核心入口函数

函数 行号 功能
OSD::OSD() 2741 OSD 构造函数
OSD::pre_init() 2941 预初始化
OSD::init() 4050 完整初始化(24步)
OSD::shutdown() 5084 关闭流程
OSD::ms_dispatch() 8158 消息分发入口

📨 消息处理

消息类型 处理函数 行号 说明
MOSDOp handle_op() ~8500 客户端 I/O 操作
MOSDPGLog handle_pg_log() ~9000 PG 日志同步
MOSDPGNotify handle_pg_notify() ~9100 PG 通知
MOSDPGCreate2 handle_pg_create() ~5500 PG 创建
MCommand handle_command() 8007 管理命令

🗂️ PG 管理

函数 行号 功能
OSD::load_pgs() ~5500 从存储加载 PG
OSD::create_pg() ~5600 创建新 PG
OSD::lookup_pg() ~5700 查找 PG
OSDService::queue_for_pg_delete() 2238 PG 删除

🔄 恢复和回填

函数 行号 功能
OSDService::queue_recovery_context() 2004 恢复上下文入队
OSDService::_queue_for_recovery() 2358 恢复操作入队
OSDService::kick_recovery_queue() ~2400 触发恢复队列

🧹 Scrub

函数 行号 功能
OSDService::queue_for_scrub() 2149 Scrub 入队
OSD::scrub_purged_snaps() 8041 清理快照映射

⏰ 定时器

函数 行号 功能
OSD::tick() 6951 持有锁的 tick
OSD::tick_without_osd_lock() 7006 不持有锁的 tick

💓 心跳

函数 行号 功能
OSD::_add_heartbeat_peer() 6098 添加心跳对等节点
OSD::maybe_update_heartbeat_peers() 6159 更新心跳对等节点
OSD::heartbeat_check() ~6800 心跳检查
OSD::heartbeat_dispatch() 8136 心跳消息分发

🗺️ OSDMap 管理

函数 行号 功能
OSDService::activate_map() 646 激活新 Map
OSDService::get_nextmap_reserved() 663 保留 Map
OSDService::await_reserved_maps() 705 等待 Map 释放
OSDService::publish_map() ~650 发布 Map

🎯 操作调度

函数 行号 功能
OSDService::enqueue_back() 1969 操作入队(尾部)
OSDService::enqueue_front() 1979 操作入队(头部)

🔧 工具函数

函数 行号 功能
OSD::write_superblock() 2418 写入 Superblock
OSD::mkfs() 2448 创建文件系统
OSD::get_osd_compat_set() 263 获取兼容性特性

📊 关键数据结构

类型 位置 说明
OSDService OSD.h 服务层,提供共享服务
OSDShard OSD.h 分片,管理 PG 和调度器
ShardedOpWQ OSD.h 分片工作队列
HeartbeatInfo OSD.h 心跳信息

🔑 关键代码段

初始化流程(24步)

1
2
3
4
5
6
7
8
9
行 4050-4583: OSD::init()
├─ 行 4059-4063: 初始化基础组件
├─ 行 4086-4090: 挂载存储
├─ 行 4141-4173: 读取 Superblock
├─ 行 4177-4185: 加载 OSDMap
├─ 行 4274: 加载 PG
├─ 行 4376: 启动操作线程池
├─ 行 4393-4408: Monitor 认证
└─ 行 4566: 启动引导

消息处理流程

1
2
3
4
5
行 8158: ms_dispatch() 入口

行 8500+: _dispatch() 路由

根据消息类型分发到相应处理函数

PG 操作流程

1
2
3
4
5
6
7
客户端操作 → ms_dispatch()

handle_op() → lookup_pg()

enqueue_back() → 操作队列

ShardedOpWQ → PG::do_request()

📍 重要行号速查

  • 文件头部注释: 行 17-29
  • OSD 构造函数: 行 2741
  • OSD 初始化: 行 4050
  • 消息分发: 行 8158
  • 定时器 tick: 行 6951, 7006
  • 心跳检查: 行 ~6800
  • PG 创建: 行 ~5500
  • 恢复调度: 行 2004, 2358

🎓 学习路径

  1. 第一天:阅读 OSD::init() (行 4050-4583)
  2. 第二天:阅读 OSD::ms_dispatch() (行 8158+) 和消息路由
  3. 第三天:阅读 PG 管理相关函数
  4. 第四天:阅读恢复和回填机制
  5. 第五天:阅读心跳和定时器

💡 调试技巧

  1. 设置断点

    • OSD::init() - 观察初始化流程
    • OSD::ms_dispatch() - 观察消息分发
    • OSD::tick() - 观察定时任务
  2. 日志输出

    • 使用 dout() 宏查看调试信息
    • 关注 derr 错误输出
  3. 符号查找

    • 使用 IDE 的 “Go to Definition”
    • 使用 grep 查找函数调用

提示:配合 OSD.h 一起阅读,理解类定义和接口

OSD.cc 快速阅读指南

OSD.cc 快速阅读指南

📋 文件概览

OSD.cc 是 Ceph OSD 的核心实现文件,约 12515 行代码。它实现了 OSD 守护进程的主要功能,包括:

  • OSD 生命周期管理(启动/停止)
  • 消息接收与分发
  • PG 管理(创建/加载/删除)
  • 恢复、回填、scrub 调度
  • 心跳监控
  • OSDMap 管理

🎯 快速理解路径(推荐阅读顺序)

第一阶段:理解整体架构(1-2小时)

1. 文件头部和基础函数(行 1-300)

  • 目的:了解文件结构和基础工具函数
  • 关键内容
    • 文件说明注释(行 17-29)
    • 兼容性特性集合函数(行 231-268)
    • OSDService 构造函数(行 278-340)

2. OSD 构造函数和初始化(行 2720-2950)

  • 目的:理解 OSD 如何启动
  • 关键函数
    1
    2
    3
    OSD::OSD(...)           // 行 2741-2884:构造函数,初始化所有组件
    OSD::pre_init() // 行 2941+:预初始化
    OSD::init() // 行 4050-4583:完整初始化流程(24个步骤)
  • 阅读重点
    • 组件初始化顺序
    • 存储后端挂载
    • Superblock 读取
    • OSDMap 加载
    • PG 加载
    • 认证和引导

3. OSD 关闭流程(行 5084-5300)

  • 目的:理解 OSD 如何优雅关闭
  • 关键函数
    1
    2
    OSD::shutdown()         // 行 5084+:关闭流程
    OSDService::shutdown() // 行 605+:服务层关闭

第二阶段:理解消息处理(2-3小时)

4. 消息分发入口(行 8158-8500)

  • 目的:理解消息如何被路由和处理
  • 关键函数
    1
    2
    OSD::ms_dispatch()      // 行 8158+:消息分发入口
    OSD::_dispatch() // 行 8500+:实际消息路由
  • 阅读重点
    • 消息类型识别
    • 路由到 PG 或服务层
    • 锁的使用

5. 客户端操作处理(行 8500-9500)

  • 目的:理解客户端 I/O 请求处理流程
  • 关键消息类型
    • MOSDOp:客户端读写操作
    • MOSDOpReply:操作回复
  • 阅读重点
    • 操作如何路由到 PG
    • 错误处理机制

6. OSD 间消息处理(行 9500-10500)

  • 目的:理解 OSD 之间的通信
  • 关键消息类型
    • MOSDPGLog:PG 日志同步
    • MOSDPGNotify:PG 通知
    • MOSDPGInfo:PG 信息交换
    • MOSDPGCreate2:PG 创建请求

第三阶段:理解 PG 管理(3-4小时)

7. PG 创建和加载(行 5500-6500)

  • 目的:理解 PG 如何被创建和管理
  • 关键函数
    1
    2
    3
    OSD::handle_pg_create()    // PG 创建处理
    OSD::load_pgs() // 从存储加载 PG
    OSD::create_pg() // 创建新 PG
  • 阅读重点
    • PG 创建触发条件
    • PG 状态初始化
    • PG 分裂和合并处理

8. PG 删除和清理(行 2238-2400)

  • 目的:理解 PG 如何被删除
  • 关键函数
    1
    OSDService::queue_for_pg_delete()  // 行 2238+

第四阶段:理解核心功能(4-5小时)

9. 恢复和回填调度(行 2004-2400)

  • 目的:理解数据恢复机制
  • 关键函数
    1
    2
    OSDService::queue_recovery_context()  // 行 2004+
    OSDService::_queue_for_recovery() // 行 2358+
  • 阅读重点
    • 恢复优先级
    • 恢复节流
    • 回填机制

10. Scrub 调度(行 2044-2238)

  • 目的:理解数据一致性检查
  • 关键函数
    1
    2
    OSDService::queue_for_scrub()         // 行 2149+
    OSDService::queue_scrub_event_msg() // 行 2089+

11. Agent 和缓存层(行 782-872)

  • 目的:理解缓存层对象提升/降级
  • 关键函数
    1
    2
    OSDService::agent_entry()              // 行 782+:Agent 主循环
    OSDService::promote_throttle_recalibrate() // 行 872+:节流校准

第五阶段:理解系统维护(2-3小时)

12. 定时器任务(行 6951-7065)

  • 目的:理解定期执行的任务
  • 关键函数
    1
    2
    OSD::tick()                    // 行 6951+:持有锁的 tick
    OSD::tick_without_osd_lock() // 行 7006+:不持有锁的 tick
  • 阅读重点
    • 心跳检查
    • Monitor 报告
    • Scrub 启动
    • PG 创建恢复

13. 心跳机制(行 6098-6950)

  • 目的:理解 OSD 健康监控
  • 关键函数
    1
    2
    3
    4
    OSD::_add_heartbeat_peer()      // 行 6098+
    OSD::maybe_update_heartbeat_peers() // 行 6159+
    OSD::heartbeat_check() // 行 6800+
    OSD::heartbeat_dispatch() // 行 8136+

14. OSDMap 管理(行 646-761)

  • 目的:理解集群映射管理
  • 关键函数
    1
    2
    3
    OSDService::activate_map()      // 行 646+:激活新 Map
    OSDService::get_nextmap_reserved() // 行 663+:保留 Map
    OSDService::await_reserved_maps() // 行 705+:等待 Map 释放

🔍 关键代码段速查

初始化流程(24个步骤)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
OSD::init() (行 4050-4583)
├─ 1. 初始化基础组件(追踪器、定时器)
├─ 2. 挂载对象存储
├─ 3. 验证对象名称长度支持
├─ 4. 读取和验证 Superblock
├─ 5. 加载 OSDMap
├─ 6. 检查遗留的 PG 删除
├─ 7. 处理兼容性特性升级
├─ 8. 初始化元数据对象
├─ 9. 加载 OSD 类
├─ 10. 设置 Epoch 和清理
├─ 11. 加载 PG
├─ 12. 初始化统计信息
├─ 13. 配置消息传递器和认证
├─ 14. 初始化 Manager 客户端
├─ 15. 注册消息分发器
├─ 16. 初始化服务并发布信息
├─ 17. 处理 PG 分裂和合并
├─ 18. 启动工作线程
├─ 19. 与 Monitor 认证
├─ 20. 更新 CRUSH 信息
├─ 21. 最终初始化
├─ 22. 订阅 Monitor 消息
├─ 23. 启动引导过程
└─ 24. QoS 相关配置覆盖

消息处理流程

1
2
3
4
5
6
7
8
9
10
11
Messenger 接收消息

OSD::ms_dispatch() (行 8158+)

OSD::_dispatch() (行 8500+)

根据消息类型路由:
├─ 客户端操作 → PG::do_request()
├─ OSD 间消息 → 相应处理函数
├─ Monitor 消息 → 相应处理函数
└─ 管理命令 → handle_command()

PG 操作流程

1
2
3
4
5
6
7
8
9
10
11
12
13
客户端发送 MOSDOp

OSD::ms_dispatch()

OSD::_dispatch() → handle_op()

查找 PG → lookup_pg()

OSDService::enqueue_back() → 加入操作队列

ShardedOpWQ 处理 → PG::do_request()

执行操作 → 回复客户端

📚 辅助阅读材料

相关文件

  1. OSD.h:类定义和接口
  2. PG.cc:PG 具体实现
  3. PeeringState.cc:PG 对等状态机
  4. PGBackend.cc:PG 后端实现

关键数据结构

  • OSDService:服务层,提供共享服务
  • OSDShard:分片,管理 PG 和调度器
  • ShardedOpWQ:分片工作队列
  • HeartbeatInfo:心跳信息

💡 阅读技巧

1. 使用 IDE 的符号导航

  • 跳转到函数定义
  • 查找函数调用
  • 查看类继承关系

2. 关注注释

  • 文件头部注释(行 17-29)
  • 函数注释(已添加中文注释)
  • 关键代码段注释

3. 理解调用链

  • 从入口函数开始(如 init()ms_dispatch()
  • 跟踪函数调用
  • 理解数据流

4. 分模块阅读

  • 不要试图一次性理解所有代码
  • 按功能模块阅读(如:初始化、消息处理、PG 管理)
  • 每个模块理解后再看下一个

5. 使用调试器

  • 设置断点
  • 单步执行
  • 观察变量值

🎓 学习路径建议

初学者(1-2周)

  1. 阅读文件头部注释
  2. 理解 OSD::init() 流程
  3. 理解 OSD::ms_dispatch() 消息分发
  4. 理解一个简单的消息处理流程(如 MOSDOp

中级(2-4周)

  1. 理解 PG 创建和加载
  2. 理解恢复和回填机制
  3. 理解心跳机制
  4. 理解 OSDMap 管理

高级(1-2个月)

  1. 深入理解所有消息类型处理
  2. 理解 PG 分裂和合并
  3. 理解 Agent 和缓存层
  4. 理解性能优化机制

⚠️ 注意事项

  1. 文件很大:不要试图一次性读完,按模块阅读
  2. 依赖复杂:理解 OSD 需要了解 PG、PeeringState 等
  3. 并发处理:注意锁的使用和并发安全
  4. 状态机:PG 状态机是核心,需要重点理解

🔗 快速跳转

  • 初始化:行 2741(构造函数)、行 4050(init)
  • 消息处理:行 8158(ms_dispatch)
  • PG 管理:行 5500+(PG 创建)
  • 恢复调度:行 2004+(恢复队列)
  • 心跳:行 6098+(心跳对等节点)
  • 定时器:行 6951(tick)

最后更新:2024年
文件版本:基于 Ceph 主分支

OSD 操作线程池(osd_op_tp)工作机制分析

OSD 操作线程池(osd_op_tp)工作机制分析

1. 概述

osd_op_tp 是 OSD 中用于处理客户端操作和内部操作的核心线程池。它采用分片(Sharded)架构,将操作按照 PG 进行分片,每个分片有独立的调度队列,多个工作线程并发处理不同分片的操作。

2. 架构设计

2.1 核心组件

1
2
3
4
5
6
7
8
osd_op_tp (ShardedThreadPool)
├── 多个工作线程(根据配置创建)
└── op_shardedwq (ShardedOpWQ)
├── 多个分片(OSDShard)
│ ├── scheduler (操作调度器)
│ ├── pg_slots (PG 槽位映射)
│ └── context_queue (上下文队列)
└── 操作入队/出队逻辑

2.2 线程数量计算

1
2
3
4
5
6
7
8
9
10
11
12
int OSD::get_num_op_threads()
{
// 如果明确配置了每个分片的线程数
if (cct->_conf->osd_op_num_threads_per_shard)
return get_num_op_shards() * cct->_conf->osd_op_num_threads_per_shard;

// 根据存储类型选择不同的线程数
if (store_is_rotational) // HDD
return get_num_op_shards() * cct->_conf->osd_op_num_threads_per_shard_hdd;
else // SSD
return get_num_op_shards() * cct->_conf->osd_op_num_threads_per_shard_ssd;
}

线程数 = 分片数 × 每个分片的线程数

3. 工作流程

3.1 操作入队(Enqueue)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
void OSD::ShardedOpWQ::_enqueue(OpSchedulerItem&& item)
{
// 1. 根据 PG ID 计算分片索引(哈希)
uint32_t shard_index =
item.get_ordering_token().hash_to_shard(osd->shards.size());

// 2. 获取对应的分片
OSDShard* sdata = osd->shards[shard_index];

// 3. 将操作加入调度器队列
{
std::lock_guard l{sdata->shard_lock};
sdata->scheduler->enqueue(std::move(item));
}

// 4. 通知等待的线程
sdata->sdata_cond.notify_all();
}

关键点:

  • 使用 PG ID 的哈希值确定分片,确保同一 PG 的操作总是路由到同一分片
  • 这保证了同一 PG 的操作顺序性

3.2 操作处理(Process)

每个工作线程执行 _process 方法,这是核心处理循环:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
void OSD::ShardedOpWQ::_process(uint32_t thread_index, 
ceph::heartbeat_handle_d *hb)
{
// 1. 确定处理哪个分片
uint32_t shard_index = thread_index % osd->num_shards;
OSDShard* sdata = osd->shards[shard_index];

// 2. 检查队列是否为空
sdata->shard_lock.lock();
if (sdata->scheduler->empty()) {
// 等待新操作到达
sdata->sdata_cond.wait(...);
}

// 3. 从调度器出队一个操作
WorkItem work_item = sdata->scheduler->dequeue();
auto item = std::move(std::get<OpSchedulerItem>(work_item));

// 4. 获取或创建 PG 槽位
const auto token = item.get_ordering_token(); // PG ID
auto r = sdata->pg_slots.emplace(token, nullptr);
if (r.second) {
r.first->second = make_unique<OSDShardPGSlot>();
}
OSDShardPGSlot *slot = r.first->second.get();

// 5. 将操作加入 PG 槽位的待处理队列
slot->to_process.push_back(std::move(item));

// 6. 获取 PG 引用并加锁
PGRef pg = slot->pg;
if (pg) {
pg->lock(); // 获取 PG 锁
}

// 7. 从待处理队列取出操作
auto qi = std::move(slot->to_process.front());
slot->to_process.pop_front();

// 8. 执行操作
qi.run(osd, sdata, pg, tp_handle);

// 9. 释放 PG 锁
if (pg) {
pg->unlock();
}
}

3.3 操作执行流程

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
操作入队

根据 PG ID 哈希到分片

加入调度器队列

工作线程从调度器出队

加入 PG 槽位的 to_process 队列

获取 PG 锁

从 to_process 取出操作

执行操作(qi.run())
├── 客户端操作 → PG::do_op()
├── Peering 事件 → PG::do_peering_event()
└── 其他操作 → 相应处理函数

释放 PG 锁

4. 关键数据结构

4.1 OSDShard(分片)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
class OSDShard {
ceph::mutex shard_lock; // 分片锁
OSDMapRef shard_osdmap; // 分片的 OSDMap 引用
unique_ptr<OpScheduler> scheduler; // 操作调度器
map<spg_t, unique_ptr<OSDShardPGSlot>> pg_slots; // PG 槽位映射

// 等待机制
ceph::mutex sdata_wait_lock;
ceph::condition_variable sdata_cond;
uint32_t waiting_threads; // 等待的线程数
bool stop_waiting; // 停止等待标志

list<Context*> context_queue; // 上下文队列(用于 oncommit 回调)
};

4.2 OSDShardPGSlot(PG 槽位)

1
2
3
4
5
6
7
8
9
10
class OSDShardPGSlot {
PGRef pg; // PG 引用
list<OpSchedulerItem> to_process; // 待处理操作队列
list<OpSchedulerItem> waiting; // 等待操作队列
list<OpSchedulerItem> waiting_peering; // 等待 Peering 的操作

uint64_t requeue_seq; // 重新入队序列号
uint32_t num_running; // 正在运行的操作数
set<spg_t> waiting_for_split; // 等待分裂的 PG
};

4.3 OpSchedulerItem(操作调度项)

包含:

  • ordering_token:排序令牌(通常是 PG ID),用于确定分片
  • map_epoch:操作对应的 OSDMap epoch
  • 操作类型:客户端操作、Peering 事件等
  • 操作数据:OpRequestRef、PGPeeringEventRef 等

5. 调度策略

5.1 调度器类型

OSD 支持多种调度器:

  • mClockScheduler:基于 mClock 算法的 QoS 调度器
  • WeightedPriorityQueue:加权优先级队列
  • SimpleQueue:简单队列(用于测试)

5.2 调度优先级

操作根据以下因素确定优先级:

  1. 操作类型:Peering 事件通常优先级较高
  2. 客户端优先级:客户端请求的优先级
  3. 成本:操作的预期成本(IO 数量、延迟等)
  4. 时间戳:操作的时间戳

6. 并发控制

6.1 分片级别的并发

  • 每个分片有独立的锁(shard_lock
  • 不同分片的操作可以并发处理
  • 同一分片内的操作串行处理(通过调度器)

6.2 PG 级别的并发

  • 每个 PG 有自己的锁(PG::lock()
  • 同一 PG 的操作必须串行处理
  • 不同 PG 的操作可以并发处理

6.3 线程分配策略

1
uint32_t shard_index = thread_index % osd->num_shards;
  • 线程通过 thread_index % num_shards 分配到分片
  • 多个线程可以处理同一分片(提高并发度)
  • 但同一 PG 的操作仍然串行(通过 PG 锁保证)

7. 等待和唤醒机制

7.1 等待条件

线程在以下情况会等待:

  1. 调度器队列为空:等待新操作到达
  2. 操作被调度到未来:等待调度时间到达
  3. PG 不存在:等待 PG 创建或加载
  4. Map epoch 不匹配:等待 OSDMap 更新

7.2 唤醒机制

1
2
3
4
5
// 入队时唤醒
sdata->sdata_cond.notify_all(); // 或 notify_one()

// 等待条件
sdata->sdata_cond.wait(lock);

8. 特殊处理场景

8.1 PG 不存在

1
2
3
4
5
6
7
8
9
10
while (!pg) {
// 检查是否需要创建 PG
if (create_info) {
pg = osd->handle_pg_create_info(osdmap, create_info);
}
// 或者等待 PG 加载
else {
_add_slot_waiter(token, slot, std::move(qi));
}
}

8.2 Map Epoch 不匹配

1
2
3
4
if (qi.get_map_epoch() > osdmap->get_epoch()) {
// 操作需要更新的 OSDMap,等待 Map 更新
_add_slot_waiter(token, slot, std::move(qi));
}

8.3 PG 分裂

1
2
3
4
if (!slot->waiting_for_split.empty()) {
// PG 正在分裂,等待分裂完成
_add_slot_waiter(token, slot, std::move(qi));
}

9. 性能优化

9.1 分片设计

  • 减少锁竞争:不同分片使用不同的锁
  • 提高缓存局部性:同一 PG 的操作总是在同一分片处理
  • 负载均衡:通过哈希均匀分布 PG 到分片

9.2 调度优化

  • 优先级调度:重要操作优先处理
  • 成本感知:考虑操作的 IO 成本
  • QoS 保证:通过 mClock 算法保证不同客户端的服务质量

9.3 批量处理

  • 上下文队列:批量处理 oncommit 回调
  • 操作合并:某些操作可以合并处理

10. 配置参数

参数 说明 默认值
osd_op_num_threads_per_shard 每个分片的线程数 -
osd_op_num_threads_per_shard_hdd HDD 每个分片的线程数 2
osd_op_num_threads_per_shard_ssd SSD 每个分片的线程数 3
osd_op_thread_timeout 线程超时时间 60s
osd_op_thread_suicide_timeout 线程自杀超时 300s

11. 总结

osd_op_tp 通过以下机制实现高效的操作处理:

  1. 分片架构:将操作按 PG 分片,减少锁竞争
  2. 多线程并发:每个分片多个线程,提高吞吐量
  3. 智能调度:根据优先级、成本等因素调度操作
  4. 顺序保证:同一 PG 的操作串行处理,保证一致性
  5. 等待机制:合理等待 PG 创建、Map 更新等条件

这种设计在保证操作顺序性的同时,最大化了并发性能。

OSD.cc 整体业务流程和设计原理详细分析

OSD.cc 整体业务流程和设计原理详细分析

一、整体架构设计

1.1 核心组件架构

OSD.cc 实现了 Ceph 分布式存储系统的核心组件 OSD(Object Storage Daemon),采用分层架构设计:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
┌─────────────────────────────────────────────────────────┐
│ OSD 主类 (OSD) │
│ - 生命周期管理 (init/shutdown) │
│ - 消息分发与路由 │
│ - PG 管理与调度 │
│ - 心跳与健康检查 │
└─────────────────────────────────────────────────────────┘

┌───────────────┼───────────────┐
│ │ │
┌───────▼──────┐ ┌──────▼──────┐ ┌─────▼──────┐
│ OSDService │ │ ObjectStore │ │ Messenger │
│ - PG 调度 │ │ - 数据持久化 │ │ - 消息传递 │
│ - 恢复管理 │ │ - 元数据管理 │ │ - 连接管理 │
│ - Scrub 调度 │ │ │ │ │
│ - 状态管理 │ │ │ │ │
└─────────────┘ └──────────────┘ └────────────┘
│ │ │
└───────────────┼───────────────┘

┌───────▼───────┐
│ PG (Placement │
│ Group) │
│ - 数据一致性 │
│ - Peering │
│ - 恢复/回填 │
└───────────────┘

1.2 设计原则

  1. 分层解耦:OSD 作为协调层,具体操作下沉到 PG、PGBackend 等组件
  2. 异步非阻塞:大量使用异步消息处理和回调机制
  3. 状态机驱动:PG 的 peering、恢复等通过状态机管理
  4. 资源池化:使用线程池、连接池等提高性能
  5. 容错设计:心跳检测、故障报告、自动恢复

二、主要业务流程

2.1 OSD 启动流程

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
pre_init()

init()
├─> 挂载 ObjectStore
├─> 读取 superblock
├─> 验证兼容性特性
├─> 加载 OSDMap
├─> 初始化服务组件
│ ├─> OSDService::init()
│ ├─> 启动心跳线程
│ ├─> 启动 agent 线程
│ └─> 初始化调度器
├─> 加载 PG (load_pgs)
└─> final_init()
├─> 启动 tick 定时器
├─> 启动操作处理线程池
└─> 开始 boot 流程
├─> _preboot()
├─> _get_purged_snaps()
├─> start_waiting_for_healthy()
└─> _send_boot()
└─> 等待 Monitor 响应

关键代码位置:

  • OSD::init() (4051行):主初始化函数
  • OSD::final_init() (4463行):最终初始化
  • OSD::start_boot() (7236行):启动引导流程

2.2 OSDMap 更新流程

OSDMap 是集群拓扑的核心,其更新流程如下:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
handle_osd_map(MOSDMap *m)

1. 等待 PG 消费完旧 map(防止缓存溢出)

2. 验证消息来源和 FSID

3. 处理完整 map 和增量 map
├─> 解码 map 数据
├─> 应用增量更新
├─> 验证 CRC
└─> 存储到 ObjectStore

4. 更新 superblock
├─> 记录 epoch 范围
├─> 更新 pg_num_history
└─> 记录 purged_snaps

5. _committed_osd_maps()
├─> 更新各 shard 的 map
├─> 通知 PG map 变更
└─> 触发 PG 处理

6. consume_map() / activate_map()
├─> PG 分裂/合并检测
├─> PG 创建/删除
└─> 推进 PG 状态

关键代码位置:

  • OSD::handle_osd_map() (8561行):处理 OSDMap 消息
  • OSD::_committed_osd_maps() (8968行):提交 map 变更
  • OSD::consume_map() (9664行):PG 消费 map
  • OSD::activate_map() (9724行):激活新 map

2.3 客户端操作处理流程

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
客户端发送 MOSDOp

ms_fast_dispatch()
├─> 创建 OpRequest
├─> 路由到 PG
│ ├─> 新客户端:等待 map 更新
│ └─> 旧客户端:直接路由
└─> enqueue_op()

op_shardedwq (分片工作队列)

_process() (工作线程)
├─> _lookup_lock_pg()
├─> pg->do_request()
│ ├─> 权限检查
│ ├─> 路由到主/副本
│ ├─> 执行操作
│ └─> 发送回复
└─> pg->unlock()

关键代码位置:

  • OSD::ms_fast_dispatch() (8045行):快速消息分发
  • OSD::enqueue_op():操作入队
  • OSD::_process():操作处理(在 OSD.h 中定义)

2.4 PG 生命周期管理

2.4.1 PG 创建流程

1
2
3
4
5
6
7
8
9
10
11
12
13
handle_fast_pg_create(MOSDPGCreate2 *m)

1. 验证消息来源和权限

2. 检查 PG 是否已存在

3. 创建 PG 实例
├─> PrimaryLogPG::create()
├─> 初始化 PG 状态
└─> 注册到 OSD

4. 触发 Peering
└─> PG::handle_peering_event()

2.4.2 PG 分裂流程

1
2
3
4
5
6
7
8
9
10
consume_map() / activate_map()

track_pools_and_pg_num_changes()
├─> 检测 pg_num 变化
└─> identify_splits_and_merges()

split_pgs()
├─> 创建子 PG
├─> 复制父 PG 数据
└─> 删除父 PG

关键代码位置:

  • OSD::handle_fast_pg_create() (9817行):处理 PG 创建
  • OSD::split_pgs() (9781行):PG 分裂
  • OSDService::identify_splits_and_merges() (359行):识别分裂/合并

2.5 恢复与回填流程

1
2
3
4
5
6
7
8
9
10
11
12
13
14
恢复触发
├─> PG Peering 完成
├─> OSDMap 变更
└─> 手动触发

queue_for_recovery()
├─> 计算恢复成本
└─> 加入调度队列

PGRecovery 处理
├─> 预留恢复槽位
├─> 选择恢复对象
├─> 执行恢复操作
└─> 更新 PG 状态

关键代码位置:

  • OSDService::queue_recovery_context() (1828行):恢复上下文入队
  • OSDService::_queue_for_recovery() (2158行):恢复任务入队

2.6 Scrub 流程

1
2
3
4
5
6
7
8
9
10
11
定时触发 / 手动触发

queue_for_scrub()
├─> 创建 scrub 事件消息
└─> 加入调度队列

PGScrub 处理
├─> 检查对象完整性
├─> 比较主副本数据
├─> 修复不一致
└─> 更新统计信息

关键代码位置:

  • OSDService::queue_for_scrub() (2018行):scrub 入队
  • OSD::handle_fast_scrub() (8227行):处理 scrub 消息

2.7 心跳与故障检测流程

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
heartbeat_entry() (心跳线程)

heartbeat()
├─> 更新统计信息
├─> 检查满载状态
├─> 向对等节点发送 PING
└─> 检查超时

handle_osd_ping()
├─> 处理 PING/PING_REPLY
├─> 更新时间戳
└─> 共享 OSDMap

heartbeat_check()
├─> 检查对等节点响应
└─> 标记故障节点

send_failures()
└─> 向 Monitor 报告故障

关键代码位置:

  • OSD::heartbeat_entry() (6532行):心跳线程入口
  • OSD::heartbeat() (6604行):执行心跳
  • OSD::handle_osd_ping() (6184行):处理心跳消息
  • OSD::heartbeat_check() (6555行):检查心跳超时

三、关键设计原理

3.1 消息分发机制

OSD 采用两级消息分发:

  1. 快速路径 (Fast Dispatch)

    • 无需 OSDMap 锁的消息直接处理
    • 包括:PING、PG 创建/通知、Scrub 等
    • 函数:ms_fast_dispatch()
  2. 标准路径 (Standard Dispatch)

    • 需要 OSDMap 锁的消息
    • 包括:OSDMap 更新、命令等
    • 函数:_dispatch()

设计优势:

  • 减少锁竞争
  • 提高并发性能
  • 降低延迟

3.2 PG 分片架构

OSD 使用分片(Shard)机制管理 PG:

1
2
3
4
OSD
├─> Shard[0] -> PG[0, N, 2N, ...]
├─> Shard[1] -> PG[1, N+1, 2N+1, ...]
└─> Shard[N-1] -> PG[N-1, 2N-1, ...]

设计优势:

  • 减少锁竞争(每个 shard 独立锁)
  • 提高并行度
  • 更好的 NUMA 亲和性

3.3 操作调度机制

OSD 使用多种调度器:

  1. mClock Scheduler(默认):

    • 基于多级反馈队列
    • 支持 QoS 保证
    • 动态成本计算
  2. WeightedPriorityQueue(传统):

    • 基于优先级和权重
    • 固定成本模型

调度队列类型:

  • op_shardedwq:客户端操作队列
  • 恢复队列:queue_recovery_context()
  • Scrub 队列:queue_for_scrub()
  • 快照修剪队列:queue_for_snap_trim()

3.4 OSDMap 缓存与引用计数

OSDMap 采用引用计数管理:

1
2
3
4
5
6
7
8
9
get_nextmap_reserved()
├─> 增加引用计数
└─> 返回 map 引用

PG 使用 map

release_map()
├─> 减少引用计数
└─> 引用为 0 时可释放

设计优势:

  • 防止 map 在使用中被释放
  • 支持 map 缓存
  • 自动内存管理

3.5 满载状态管理

OSD 维护多级满载状态:

1
NONE < NEARFULL < BACKFILLFULL < FULL < FAILSAFE

状态计算:

  • recalc_full_state():根据使用率计算状态
  • check_full_status():更新当前状态
  • need_fullness_update():检查是否需要上报

设计优势:

  • 渐进式节流
  • 防止数据丢失
  • 支持测试注入

3.6 心跳机制设计

心跳采用双通道设计:

1
2
3
4
前端通道 (Front)         后端通道 (Back)
│ │
├─> 客户端网络 ├─> 集群网络
└─> 可配置 └─> 可配置

心跳消息类型:

  • PING:发送心跳请求
  • PING_REPLY:心跳响应
  • YOU_DIED:通知对方已下线

时间戳同步:

  • 使用单调时钟避免时钟漂移
  • HeartbeatStamps 维护时间戳
  • 计算网络延迟和时钟偏差

3.7 故障检测与恢复

故障检测:

  1. 心跳超时检测
  2. 连接断开检测
  3. 操作失败检测

故障处理:

  1. 标记故障节点
  2. 向 Monitor 报告
  3. 触发 PG Peering
  4. 启动恢复流程

四、关键数据结构

4.1 OSD 核心数据结构

1
2
3
4
5
6
7
8
9
class OSD {
OSDService service; // 服务层
ObjectStore *store; // 存储后端
OSDMapRef osdmap; // 当前 OSDMap
map<spg_t, PGRef> pg_map; // PG 映射表
OpShardedWQ op_shardedwq; // 操作队列
HeartbeatThread heartbeat_thread; // 心跳线程
// ...
};

4.2 OSDService 核心数据结构

1
2
3
4
5
6
7
8
class OSDService {
Objecter *objecter; // 对象器
OpScheduler op_scheduler; // 操作调度器
map<epoch_t, OSDMapRef> map_cache; // Map 缓存
RecoveryReserver recovery_reserver; // 恢复资源预留
ScrubReserver scrub_reserver; // Scrub 资源预留
// ...
};

五、性能优化设计

5.1 异步处理

  • 大量使用异步回调和 Future
  • 避免阻塞主线程
  • 提高并发性能

5.2 批处理优化

  • OSDMap 批量更新
  • 操作批量提交
  • 减少系统调用

5.3 NUMA 优化

  • set_numa_affinity():设置 CPU 亲和性
  • 根据存储和网络 NUMA 节点优化
  • 减少跨节点访问

5.4 缓存策略

  • OSDMap 缓存(LRU)
  • 连接缓存
  • 对象元数据缓存

六、容错与可靠性

6.1 数据一致性

  • PG Peering 保证一致性
  • Scrub 检测和修复不一致
  • 事务保证原子性

6.2 故障恢复

  • 自动故障检测
  • 自动恢复和回填
  • 降级模式支持

6.3 优雅关闭

  • shutdown():完整关闭流程
  • fast_shutdown():快速关闭模式
  • 确保数据持久化

七、总结

OSD.cc 是 Ceph 存储系统的核心实现,采用:

  1. 分层架构:清晰的职责划分
  2. 异步设计:高并发性能
  3. 状态机驱动:可靠的状态管理
  4. 资源池化:高效的资源利用
  5. 容错设计:高可用性保障

整个设计体现了分布式系统的最佳实践,在性能、可靠性和可维护性之间取得了良好的平衡。

PG.cc 业务功能与原理详细分析

PG.cc 业务功能与原理详细分析

一、PG 概述

1.1 PG 的定义与作用

Placement Group (PG) 是 Ceph 分布式存储系统的核心抽象,负责:

  • 数据分片管理:将对象映射到特定的 PG,实现数据分片
  • 一致性保证:通过 Peering 机制保证副本间一致性
  • 故障恢复:自动检测和恢复数据不一致
  • 负载均衡:通过 PG 分布实现数据在 OSD 间的均衡

1.2 PG 在 Ceph 架构中的位置

1
2
3
4
5
6
7
8
9
客户端请求

OSD (OSD.cc)
├─> 路由到 PG
└─> PG (PG.cc)
├─> PeeringState (状态机)
├─> PGBackend (数据操作)
├─> PGLog (日志管理)
└─> Scrubber (数据校验)

二、核心业务功能

2.1 PG 生命周期管理

2.1.1 PG 初始化

1
PG::PG(OSDService *o, OSDMapRef curmap, const PGPool &_pool, spg_t p)

初始化流程:

  1. 设置 PG ID 和集合(Collection)
  2. 初始化 SnapMapper(快照映射器)
  3. 创建 PeeringState(peering 状态机)
  4. 初始化统计信息

关键代码位置:

  • PG::PG() (199行):构造函数
  • PG::init() (854行):初始化 PG 状态
  • PG::read_state() (1086行):从磁盘读取 PG 状态

2.1.2 PG 状态读取

1
void PG::read_state(ObjectStore *store)

流程:

  1. 从 ObjectStore 读取 PG 元数据
  2. 读取 PG 信息(info)和 PastIntervals
  3. 读取 PGLog(操作日志)
  4. 初始化 PeeringState
  5. 触发 Peering 流程

数据结构:

  • pg_info_t:PG 基本信息(版本、统计等)
  • PastIntervals:历史区间信息
  • PGLog:操作日志

2.2 Peering 机制

2.2.1 Peering 概述

Peering 是 PG 的核心机制,用于:

  • 确定权威数据:在多个副本间确定哪个版本是权威的
  • 同步状态:确保所有副本对 PG 状态达成一致
  • 处理分歧:解决副本间的数据分歧

2.2.2 Peering 事件处理

1
void PG::do_peering_event(PGPeeringEventRef evt, PeeringCtx &rctx)

Peering 事件类型:

  • NullEvt:空事件
  • Activate:激活 PG
  • AdvMap:OSDMap 推进
  • Query:查询对等节点状态
  • Notify:通知对等节点
  • Info:信息交换

处理流程:

1
2
3
4
5
6
7
8
9
10
11
收到 Peering 事件

验证事件有效性(epoch 检查)

调用 recovery_state.handle_event()

状态机转换

执行相应操作(PeeringCtx)

写入事务(write_if_dirty)

关键代码位置:

  • PG::do_peering_event() (2128行):处理 peering 事件
  • PG::queue_peering_event() (2142行):将事件加入队列
  • PG::handle_advance_map() (2194行):处理 map 推进

2.2.3 OSDMap 推进处理

1
2
3
4
5
void PG::handle_advance_map(
OSDMapRef osdmap, OSDMapRef lastmap,
vector<int>& newup, int up_primary,
vector<int>& newacting, int acting_primary,
PeeringCtx &rctx)

处理步骤:

  1. 检测 acting set 变化
  2. 检测 up set 变化
  3. 触发 Peering 状态机转换
  4. 处理 PG 分裂/合并
  5. 更新 PG 状态

2.3 客户端操作处理

2.3.1 操作路由

PG 接收来自 OSD 的操作请求,根据 PG 状态决定处理方式:

操作状态检查:

  • can_discard_op():检查操作是否可以丢弃
  • can_discard_request():检查请求是否可以丢弃
  • old_peering_msg():检查是否为过期的 peering 消息

关键代码位置:

  • PG::can_discard_op() (1970行):检查操作是否过期
  • PG::can_discard_request() (2070行):检查请求是否过期

2.3.2 操作等待机制

当 PG 处于非活跃状态时,操作会被暂存:

1
2
3
void PG::requeue_op(OpRequestRef op)
void PG::requeue_ops(list<OpRequestRef> &ls)
void PG::requeue_map_waiters()

等待队列类型:

  • waiting_for_map:等待 OSDMap 更新
  • waiting_for_peered:等待 Peering 完成
  • waiting_for_flush:等待刷新完成
  • waiting_for_clean_to_primary_repair:等待清理完成

2.4 恢复(Recovery)机制

2.4.1 恢复触发

1
void PG::queue_recovery()

触发条件:

  • PG 是主副本(is_primary())
  • PG 已完成 Peering(is_peered())
  • 存在缺失对象(missing objects)

恢复流程:

1
2
3
4
5
6
7
8
9
10
11
12
13
queue_recovery()

计算恢复成本(平均对象大小)

加入恢复队列(osd->queue_for_recovery)

恢复调度器处理

选择恢复对象

执行恢复操作

finish_recovery_op()

关键代码位置:

  • PG::queue_recovery() (425行):将恢复加入队列
  • PG::start_recovery_op() (493行):开始恢复操作
  • PG::finish_recovery_op() (508行):完成恢复操作
  • PG::finish_recovery() (452行):完成所有恢复

2.4.2 恢复状态管理

1
2
void PG::clear_recovery_state()
void PG::cancel_recovery()

恢复状态:

  • recovery_queued:是否已加入恢复队列
  • recovery_ops_active:活跃的恢复操作数
  • recovering_oids:正在恢复的对象集合

2.5 Scrub 机制

2.5.1 Scrub 概述

Scrub 用于检测和修复数据不一致:

Scrub 类型:

  • Shallow Scrub:检查对象元数据和校验和
  • Deep Scrub:读取并验证对象数据

2.5.2 Scrub 调度

1
2
3
Scrub::schedule_result_t PG::start_scrubbing(
const Scrub::SchedEntry& candidate,
Scrub::OSDRestrictions osd_restrictions)

Scrub 条件检查:

  • PG 必须是活跃的(is_active())
  • 不能有 NOSCRUB 标志
  • 满足 scrub 间隔要求

关键代码位置:

  • PG::start_scrubbing() (1278行):启动 scrub
  • PG::scrub_requested() (1328行):scrub 请求处理
  • PG::on_scrub_schedule_input_change() (1314行):scrub 调度变化处理

2.6 快照管理

2.6.1 快照映射

1
2
3
4
void PG::update_object_snap_mapping(
ObjectStore *t, const hobject_t &soid, const set<snapid_t> &snaps)
void PG::clear_object_snap_mapping(
ObjectStore *t, const hobject_t &soid)

快照映射器(SnapMapper)作用:

  • 维护对象到快照的映射关系
  • 支持快照克隆和删除
  • 管理快照元数据

2.6.2 快照修剪

1
2
void PG::queue_snap_retrim(snapid_t snap)
void PG::on_active_actmap()

快照修剪流程:

  1. 检测需要修剪的快照
  2. 加入修剪队列(snap_trimq)
  3. 执行修剪操作
  4. 更新快照映射

关键代码位置:

  • PG::queue_snap_retrim() (1529行):加入快照修剪队列
  • PG::on_active_actmap() (1550行):激活时处理快照修剪

2.7 PG 分裂与合并

2.7.1 PG 分裂

1
void PG::split_into(pg_t child_pgid, PG *child, unsigned split_bits)

分裂流程:

  1. 创建子 PG
  2. 复制父 PG 状态到子 PG
  3. 更新 SnapMapper 的 split_bits
  4. 复制快照修剪队列
  5. 执行数据分裂(_split_into)

关键代码位置:

  • PG::split_into() (528行):PG 分裂
  • PG::start_split_stats() (545行):开始分裂统计
  • PG::finish_split_stats() (550行):完成分裂统计

2.7.2 PG 合并

1
2
3
void PG::merge_from(map<spg_t,PGRef>& sources, PeeringCtx &rctx,
unsigned split_bits,
const pg_merge_meta_t& last_pg_merge_meta)

合并流程:

  1. 合并源 PG 的状态
  2. 合并 PGLog
  3. 合并集合(merge_collection)
  4. 更新 SnapMapper
  5. 删除源 PG 元数据

关键代码位置:

  • PG::merge_from() (555行):PG 合并

2.8 Backoff 机制

2.8.1 Backoff 概述

Backoff 用于在 PG 状态不稳定时阻止客户端操作:

1
2
3
4
void PG::add_backoff(const ceph::ref_t<Session>& s, 
const hobject_t& begin,
const hobject_t& end)
void PG::release_backoffs(const hobject_t& begin, const hobject_t& end)

Backoff 状态:

  • STATE_NEW:新建
  • STATE_ACKED:已确认
  • STATE_DELETING:删除中

使用场景:

  • Peering 过程中
  • 恢复过程中
  • PG 分裂/合并时

关键代码位置:

  • PG::add_backoff() (582行):添加 backoff
  • PG::release_backoffs() (608行):释放 backoff
  • PG::clear_backoffs() (670行):清除所有 backoff

2.9 资源预留机制

2.9.1 恢复资源预留

1
2
3
4
void PG::request_remote_recovery_reservation(
unsigned priority,
PGPeeringEventURef on_grant,
PGPeeringEventURef on_preempt)

资源类型:

  • 本地资源:本地 OSD 的 IO 资源
  • 远程资源:远程 OSD 的恢复资源
  • Scrub 资源:scrub 操作的资源

关键代码位置:

  • PG::request_local_background_io_reservation() (1385行)
  • PG::request_remote_recovery_reservation() (1410行)
  • PG::cancel_remote_recovery_reservation() (1423行)

三、核心设计原理

3.1 状态机驱动

PG 使用状态机(PeeringState)管理 PG 生命周期:

1
2
3
4
5
Initial → Reset → Started → Peering → Active

Incomplete

Down

状态转换触发:

  • OSDMap 变更
  • Peering 事件
  • 恢复完成
  • 错误处理

3.2 版本控制机制

3.2.1 版本类型

  • last_update:最后更新版本
  • last_complete:最后完整版本
  • log_tail:日志尾部版本

3.2.2 版本比较

通过版本比较确定:

  • 哪个副本数据最新
  • 是否需要恢复
  • 数据是否一致

3.3 PGLog 机制

3.3.1 PGLog 作用

  • 操作记录:记录所有写操作
  • 版本追踪:追踪对象版本变化
  • 恢复依据:用于恢复缺失对象

3.3.2 PGLog 结构

1
2
3
4
5
class PGLog {
IndexedLog log; // 操作日志
map<hobject_t, pg_missing_t> missing; // 缺失对象
// ...
};

3.4 锁机制

3.4.1 PG 锁

1
2
3
void PG::lock(bool no_lockdep) const
void PG::unlock() const
bool PG::is_locked() const

锁的作用:

  • 保护 PG 状态一致性
  • 防止并发修改
  • 确保操作原子性

锁的持有者:

  • 操作处理线程
  • Peering 线程
  • 恢复线程

3.5 事务机制

3.5.1 事务使用

1
2
3
4
5
6
7
8
9
void PG::prepare_write(
pg_info_t &info,
pg_info_t &last_written_info,
PastIntervals &past_intervals,
PGLog &pglog,
bool dirty_info,
bool dirty_big_info,
bool need_write_epoch,
ObjectStore::Transaction &t)

事务包含:

  • PG 信息更新
  • PGLog 写入
  • 对象操作
  • 元数据更新

3.6 引用计数管理

3.6.1 引用计数

1
2
void PG::get(const char* tag)
void PG::put(const char* tag)

引用计数用途:

  • 防止 PG 在使用中被删除
  • 跟踪 PG 使用情况
  • 调试内存泄漏

四、关键数据结构

4.1 PG 核心数据结构

1
2
3
4
5
6
7
8
9
class PG {
spg_t pg_id; // PG ID
coll_t coll; // 集合
pg_info_t info; // PG 信息
PeeringState recovery_state; // Peering 状态机
PGLog projected_log; // 投影日志
SnapMapper snap_mapper; // 快照映射器
// ...
};

4.2 等待队列

1
2
3
4
map<entity_name_t, list<OpRequestRef>> waiting_for_map;
list<OpRequestRef> waiting_for_peered;
list<OpRequestRef> waiting_for_flush;
list<OpRequestRef> waiting_for_clean_to_primary_repair;

4.3 恢复相关

1
2
3
bool recovery_queued;              // 是否已加入恢复队列
int recovery_ops_active; // 活跃恢复操作数
set<hobject_t> recovering_oids; // 正在恢复的对象

五、业务流程示例

5.1 客户端写操作流程

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
客户端发送写请求

OSD 路由到 PG

PG::can_discard_op() 检查

PG 状态检查
├─> 非活跃 → 加入 waiting_for_peered
└─> 活跃 → 继续处理

权限检查 (op_has_sufficient_caps)

路由到主/副本
├─> 主副本 → 执行写操作
└─> 副本 → 接收复制

更新 PGLog

提交事务

发送回复

5.2 Peering 流程

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
OSDMap 变更

handle_advance_map()

检测 acting set 变化

触发 Peering 事件

do_peering_event()

recovery_state.handle_event()

状态机转换
├─> Reset → Started
├─> Started → Peering
└─> Peering → Active

查询对等节点状态

确定权威数据

同步状态

激活 PG (on_activate)

处理等待的操作

5.3 恢复流程

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
检测到缺失对象

queue_recovery()

加入恢复队列

恢复调度器选择 PG

start_recovery_op()

选择恢复对象

从对等节点拉取数据

写入本地

finish_recovery_op()

继续下一个对象

所有对象恢复完成

finish_recovery()

触发 scrub(可选)

六、性能优化设计

6.1 异步处理

  • Peering 事件异步处理
  • 恢复操作异步执行
  • 操作等待队列异步唤醒

6.2 批量操作

  • 批量恢复对象
  • 批量更新统计信息
  • 批量处理等待操作

6.3 资源预留

  • 恢复资源预留避免冲突
  • Scrub 资源预留保证优先级
  • 本地 IO 资源预留

七、容错与可靠性

7.1 状态一致性

  • 通过 Peering 保证副本一致性
  • 通过 PGLog 追踪操作历史
  • 通过版本控制检测分歧

7.2 故障恢复

  • 自动检测缺失对象
  • 自动从对等节点恢复
  • 自动处理副本故障

7.3 数据校验

  • Scrub 检测数据不一致
  • 自动修复(auto-repair)
  • 校验和验证

八、总结

PG.cc 实现了 Ceph 存储系统的核心逻辑:

  1. 状态管理:通过 PeeringState 状态机管理 PG 生命周期
  2. 一致性保证:通过 Peering 和 PGLog 保证数据一致性
  3. 故障恢复:自动检测和恢复数据不一致
  4. 操作处理:路由和处理客户端操作
  5. 资源管理:管理恢复、scrub 等后台任务

整个设计体现了分布式系统的核心原则:

  • 状态机驱动:清晰的状态转换
  • 版本控制:精确的版本追踪
  • 异步处理:高效的并发处理
  • 容错设计:可靠的故障处理

PG 作为 Ceph 的核心抽象,承担了数据分片、一致性保证、故障恢复等关键职责,是整个存储系统稳定运行的基础。

PG 存储调用与数据读取机制分析

PG 存储调用与数据读取机制分析

1. 概述

PG(Placement Group)是 Ceph OSD 中管理对象存储的核心组件。PG 通过 ObjectStore 接口与底层存储后端(如 BlueStore、FileStore)交互,实现对象的读写、元数据管理等功能。

1.1 核心组件

  • ObjectStore:存储后端抽象接口
  • CollectionHandle (ch):集合句柄,用于标识和管理 PG 的存储集合
  • PGTransaction:PG 层的事务抽象
  • ObjectStore::Transaction:存储层的事务实现

1.2 存储层次

1
2
3
4
5
6
7
PG (PrimaryLogPG)

PGBackend (ReplicatedBackend/ECBackend)

ObjectStore (BlueStore/FileStore)

底层存储设备

2. 存储初始化与集合管理

2.1 PG 存储初始化

PG 在创建时会初始化存储相关的组件:

1
2
3
4
5
6
7
8
9
10
11
12
13
// PG 构造函数
PG::PG(OSDService *o, OSDMapRef curmap, const PGPool &_pool, spg_t p) :
pg_whoami(o->whoami, p.shard),
pg_id(p),
coll(p), // 集合 ID,基于 PG ID
osd(o),
cct(o->cct),
osdriver(osd->store, coll_t(), OSD::make_snapmapper_oid()),
// ...
{
// coll 是 PG 的集合标识符
// 格式:coll_t(pgid),例如 coll_t(1.23)
}

关键成员变量

  • coll:集合 ID(coll_t),标识 PG 的存储集合
  • ch:集合句柄(ObjectStore::CollectionHandle),用于访问集合

2.2 打开集合(Collection)

PG 在加载时会打开对应的存储集合:

1
2
// OSD::load_pgs() 中
pg->ch = store->open_collection(pg->coll);

集合句柄的作用

  • 标识 PG 的存储集合
  • 用于事务提交和查询操作
  • 保证事务在同一个集合内按顺序执行

2.3 集合与 PG 的对应关系

1
2
3
4
5
6
// PG ID 到集合 ID 的转换
coll_t coll(pgid); // 直接使用 PG ID 作为集合 ID

// 例如:
// pgid = spg_t(1.23)
// coll = coll_t(1.23)

3. 存储读取机制

3.1 读取 PG 元数据

PG 启动时需要从存储中读取元数据:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
// PG::read_state() - 读取 PG 状态
void PG::read_state(ObjectStore *store)
{
// 1. 读取 PG 信息
int r = read_info(store, pg_id, coll, info_from_disk,
past_intervals_from_disk, info_struct_v);

// 2. 读取 PGLog
recovery_state.get_pg_log().read_log_and_missing(
store, coll, pgmeta_oid, info_from_disk,
past_intervals_from_disk, cct);

// 3. 初始化 PeeringState
recovery_state.init(...);
}

读取 PG 信息

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
int PG::read_info(ObjectStore *store, spg_t pgid, const coll_t &coll,
pg_info_t &info, PastIntervals &past_intervals,
__u8 &struct_v)
{
// 1. 打开集合
auto ch = store->open_collection(coll);
ceph_assert(ch);

// 2. 准备要读取的键
set<string> keys;
keys.insert(string(infover_key)); // 版本信息
keys.insert(string(info_key)); // PG 信息
keys.insert(string(biginfo_key)); // 大信息(PastIntervals)
keys.insert(string(fastinfo_key)); // 快速信息

// 3. 构建元数据对象 ID
ghobject_t pgmeta_oid(pgid.make_pgmeta_oid());

// 4. 从 OMAP 读取值
map<string,bufferlist> values;
int r = store->omap_get_values(ch, pgmeta_oid, keys, &values);

// 5. 解码数据
auto p = values[string(infover_key)].cbegin();
decode(struct_v, p);

p = values[string(info_key)].begin();
decode(info, p);

p = values[string(biginfo_key)].begin();
decode(past_intervals, p);

return 0;
}

3.2 读取对象数据

3.2.1 通过 PGBackend 读取

PG 通过 PGBackend 读取对象数据:

1
2
3
4
5
6
7
8
9
10
11
12
// PrimaryLogPG::do_op() - 处理客户端操作
void PrimaryLogPG::do_op(OpRequestRef& op)
{
// 1. 获取对象上下文
ObjectContextRef obc = get_object_context(head, false);

// 2. 创建操作上下文
OpContext *ctx = new OpContext(op, obc, this);

// 3. 执行操作(包括读取)
execute_ctx(ctx);
}

3.2.2 读取操作流程

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
// PrimaryLogPG::execute_ctx() - 执行操作上下文
void PrimaryLogPG::execute_ctx(OpContext *ctx)
{
// 1. 准备事务
int result = prepare_transaction(ctx);

// 2. 如果是只读操作
if (ctx->op_t->empty() || result < 0) {
// 完成读取
complete_read_ctx(result, ctx);
return;
}

// 3. 写入操作需要提交事务
// ...
}

3.2.3 直接存储读取

在某些场景下,PG 会直接调用存储接口读取数据:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
// 读取对象属性
int r = osd->store->getattr(
ch,
ghobject_t(oid, ghobject_t::NO_GEN, pg_whoami.shard),
OI_ATTR, // 对象信息属性
bv);

// 读取对象统计信息
struct stat st;
int r = osd->store->stat(
ch,
ghobject_t(oid, ghobject_t::NO_GEN, pg_whoami.shard),
&st);

// 读取 OMAP 值
int r = store->omap_get_values(
ch,
pgmeta_oid,
keys,
&values);

3.3 异步读取

对于 EC(纠删码)池,支持异步读取:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
// PrimaryLogPG::OpContext::start_async_reads()
void PrimaryLogPG::OpContext::start_async_reads(PrimaryLogPG *pg)
{
inflightreads = 1;

// 准备异步读取请求
list<boost::tuple<hobject_t, uint64_t, uint64_t, uint32_t>> in;
in.swap(pending_async_reads);

// 通过 PGBackend 异步读取
pg->pgbackend->objects_read_async(
in,
new OnReadComplete(pg, this),
pg->get_pool().fast_read);
}

4. 存储写入机制

4.1 事务构建

PG 通过事务(Transaction)来组织写入操作:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
// PrimaryLogPG::execute_ctx() - 构建事务
void PrimaryLogPG::execute_ctx(OpContext *ctx)
{
// 1. 创建 PG 事务
ctx->op_t.reset(new PGTransaction);

// 2. 准备事务(添加操作)
int result = prepare_transaction(ctx);

// 3. 将 PG 事务转换为存储事务
ObjectStore::Transaction t;
ctx->op_t->apply(&osdriver, &t);

// 4. 提交事务
int r = osd->store->queue_transaction(ch, std::move(t), NULL);
}

4.2 事务提交流程

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
// PrimaryLogPG::execute_ctx() - 提交事务
void PrimaryLogPG::execute_ctx(OpContext *ctx)
{
// 1. 准备存储事务
ObjectStore::Transaction t;

// 2. 应用 PG 事务到存储事务
ctx->op_t->apply(&osdriver, &t);

// 3. 注册提交回调
ctx->register_on_commit([...]() {
// 发送回复给客户端
osd->send_message_osd_client(reply, m->get_connection());
});

// 4. 通过 PGBackend 提交(会复制到副本)
issue_repop(repop, ctx);
}

4.3 通过 PGBackend 写入

PG 通过 PGBackend 处理写入操作:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
// ReplicatedBackend::submit_transaction()
void ReplicatedBackend::submit_transaction(
const hobject_t &hoid,
const eversion_t &version,
PGTransactionUPtr &&t,
const eversion_t &trim_to,
const eversion_t &roll_forward_to,
const vector<pg_log_entry_t> &log_entries,
boost::optional<pg_hit_set_history_t> &hset_history,
Context *on_all_commit,
Context *on_all_applied,
Context *on_local_applied_sync)
{
// 1. 转换 PG 事务为存储事务
ObjectStore::Transaction _t;
t->apply(&osdriver, &_t);

// 2. 提交本地事务
int r = store->queue_transaction(ch, std::move(_t), ...);

// 3. 发送到副本(如果是主副本)
if (is_primary()) {
send_transaction_to_replicas(...);
}
}

4.4 直接存储写入

在某些场景下,PG 会直接写入存储:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
// 写入 PG 状态
void PeeringState::write_if_dirty(ObjectStore::Transaction& t)
{
if (dirty_info) {
// 准备写入数据
map<string,bufferlist> km;
prepare_info_keymap(...);

// 添加到事务
t.omap_setkeys(coll, pgmeta_oid, km);
}

// 提交事务
int r = store->queue_transaction(ch, std::move(t), NULL);
}

5. 存储接口详解

5.1 ObjectStore 核心接口

5.1.1 集合操作

1
2
3
4
5
6
7
8
9
10
11
// 打开集合
ObjectStore::CollectionHandle ch = store->open_collection(coll);

// 创建新集合
ObjectStore::CollectionHandle ch = store->create_new_collection(coll);

// 刷新集合(等待所有事务完成)
ch->flush();

// 异步刷新提交
bool idle = ch->flush_commit(callback);

5.1.2 对象读取

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
// 读取对象数据
int store->read(
CollectionHandle &ch,
const ghobject_t &oid,
uint64_t offset,
size_t len,
bufferlist &bl);

// 读取对象属性
int store->getattr(
CollectionHandle &ch,
const ghobject_t &oid,
const char *name,
bufferlist &value);

// 获取对象统计信息
int store->stat(
CollectionHandle &ch,
const ghobject_t &oid,
struct stat *st);

5.1.3 OMAP 操作

1
2
3
4
5
6
7
8
9
10
11
12
// 读取 OMAP 值
int store->omap_get_values(
CollectionHandle &ch,
const ghobject_t &oid,
const set<string> &keys,
map<string,bufferlist> *out);

// 读取 OMAP 头部
int store->omap_get_header(
CollectionHandle &ch,
const ghobject_t &oid,
bufferlist *header);

5.1.4 事务操作

1
2
3
4
5
6
// 提交事务
int store->queue_transaction(
CollectionHandle &ch,
Transaction &&t,
TrackedOpRef op = TrackedOpRef(),
ThreadPool::TPHandle *handle = NULL);

5.2 Transaction 操作类型

存储事务支持多种操作:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
// 写入数据
t.write(oid, offset, len, data);

// 设置属性
t.setattr(oid, name, value);

// OMAP 操作
t.omap_setkeys(oid, keys);
t.omap_setheader(oid, header);
t.omap_rmkeys(oid, keys);

// 对象操作
t.touch(oid); // 创建对象
t.remove(oid); // 删除对象
t.clone(oid, noid); // 克隆对象
t.clone_range(oid, noid, offset, len, dest_offset);

6. 数据读取场景

6.1 客户端读取操作

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
// PrimaryLogPG::do_op() - 处理读取操作
void PrimaryLogPG::do_op(OpRequestRef& op)
{
// 1. 获取对象上下文
ObjectContextRef obc = get_object_context(head, false);

// 2. 创建操作上下文
OpContext *ctx = new OpContext(op, obc, this);

// 3. 执行操作
execute_ctx(ctx);

// 4. 在 prepare_transaction() 中:
// - 如果是读取操作,通过 do_osd_ops() 读取数据
// - 数据从 ObjectContext 或直接读取存储获取
}

6.2 恢复读取

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
// ReplicatedBackend::recover_object() - 恢复对象
int ReplicatedBackend::recover_object(
const hobject_t &hoid,
eversion_t v,
ObjectContextRef head,
ObjectContextRef obc,
RecoveryHandle *_h)
{
if (get_parent()->get_local_missing().is_missing(hoid)) {
// 本地缺失,从副本拉取
prepare_pull(v, hoid, head, h);
} else {
// 本地有,推送到副本
start_pushes(hoid, obc, h);
}
}

6.3 Scrub 读取

1
2
3
4
5
6
7
8
9
// 在 Scrub 过程中读取对象
// 1. 读取对象数据
store->read(ch, oid, 0, size, data);

// 2. 读取对象信息
store->getattr(ch, oid, OI_ATTR, oi_bl);

// 3. 验证数据完整性
verify_object_data(oid, data, oi);

7. 数据写入场景

7.1 客户端写入操作

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
// PrimaryLogPG::execute_ctx() - 执行写入
void PrimaryLogPG::execute_ctx(OpContext *ctx)
{
// 1. 准备事务
ctx->op_t.reset(new PGTransaction);

// 2. 执行 OSD 操作(构建事务)
int result = prepare_transaction(ctx);

// 3. 转换为存储事务
ObjectStore::Transaction t;
ctx->op_t->apply(&osdriver, &t);

// 4. 通过 PGBackend 提交(包含复制)
issue_repop(repop, ctx);
}

7.2 PG 状态写入

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
// PeeringState::write_if_dirty() - 写入 PG 状态
void PeeringState::write_if_dirty(ObjectStore::Transaction& t)
{
if (dirty_info) {
// 准备 PG 信息
map<string,bufferlist> km;
prepare_info_keymap(...);

// 添加到事务
t.omap_setkeys(coll, pgmeta_oid, km);
}

if (dirty_big_info) {
// 准备大信息(PastIntervals)
// ...
}

// PGLog 写入
pg_log.write_log_and_missing(t, pgmeta_oid, ...);
}

7.3 恢复写入

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
// ReplicatedBackend::run_recovery_op() - 执行恢复操作
void ReplicatedBackend::run_recovery_op(
PGBackend::RecoveryHandle *_h,
int priority)
{
RPGHandle *h = static_cast<RPGHandle *>(_h);

// 1. 发送推送(推送到副本)
send_pushes(priority, h->pushes);

// 2. 发送拉取(从副本拉取)
send_pulls(priority, h->pulls);

// 3. 发送删除
send_recovery_deletes(priority, h->deletes);
}

8. 存储调用路径

8.1 读取路径

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
客户端请求

OSD::dispatch_op()

PrimaryLogPG::do_op()

PrimaryLogPG::execute_ctx()

PrimaryLogPG::prepare_transaction()

PrimaryLogPG::do_osd_ops() // 读取操作

ObjectContext::get() // 从缓存获取
↓ (缓存未命中)
ObjectStore::read() // 直接读取存储

底层存储后端(BlueStore/FileStore)

8.2 写入路径

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
客户端请求

OSD::dispatch_op()

PrimaryLogPG::do_op()

PrimaryLogPG::execute_ctx()

PrimaryLogPG::prepare_transaction()

PGTransaction::add_*() // 添加操作到 PG 事务

PGTransaction::apply() // 转换为存储事务

PGBackend::submit_transaction()

ObjectStore::queue_transaction() // 提交到存储

底层存储后端(BlueStore/FileStore)

9. 关键数据结构

9.1 CollectionHandle

1
2
3
4
5
6
7
// 集合句柄
ObjectStore::CollectionHandle ch;

// 用途:
// 1. 标识 PG 的存储集合
// 2. 用于事务提交
// 3. 保证事务顺序

9.2 PGTransaction

1
2
3
4
5
6
7
8
9
10
// PG 层事务抽象
class PGTransaction {
// 添加操作
void add_write(const hobject_t &oid, uint64_t off, uint64_t len, ...);
void add_setattr(const hobject_t &oid, const string &name, ...);
void add_omap_setkeys(const hobject_t &oid, ...);

// 转换为存储事务
void apply(OSDriver *osdriver, ObjectStore::Transaction *t);
};

9.3 ObjectStore::Transaction

1
2
3
4
5
6
7
8
9
10
// 存储层事务
class Transaction {
// 操作列表
vector<Op> ops;

// 添加操作
void write(const ghobject_t &oid, uint64_t off, uint64_t len, ...);
void setattr(const ghobject_t &oid, const char *name, ...);
void omap_setkeys(const ghobject_t &oid, ...);
};

10. 性能优化

10.1 对象上下文缓存

1
2
3
4
5
6
7
// ObjectContext 缓存对象信息
class ObjectContext {
ObjectState obs; // 对象状态
SnapSetContext *ssc; // 快照集上下文

// 缓存对象信息,避免重复读取
};

10.2 批量操作

1
2
3
4
5
// 批量提交事务
vector<Transaction> tls;
tls.push_back(std::move(t1));
tls.push_back(std::move(t2));
store->queue_transactions(ch, tls, op, handle);

10.3 异步操作

1
2
3
4
5
// 异步读取(EC 池)
pgbackend->objects_read_async(
read_requests,
callback,
fast_read);

11. 错误处理

11.1 读取错误

1
2
3
4
5
6
7
8
9
// 读取失败处理
int r = store->read(ch, oid, offset, len, bl);
if (r < 0) {
if (r == -ENOENT) {
// 对象不存在
} else {
// 其他错误
}
}

11.2 写入错误

1
2
3
4
5
6
// 事务提交失败处理
int r = store->queue_transaction(ch, std::move(t), op);
if (r != 0) {
// 处理错误
derr << "queue_transaction failed: " << cpp_strerror(r) << dendl;
}

12. 总结

12.1 核心机制

  1. 集合管理:每个 PG 对应一个存储集合,通过 CollectionHandle 访问
  2. 事务机制:通过 PGTransactionObjectStore::Transaction 组织操作
  3. 读写分离:读取可以直接调用存储接口,写入通过事务提交
  4. 后端抽象:通过 PGBackend 抽象复制/EC 差异

12.2 关键接口

  • 集合操作open_collection(), flush(), flush_commit()
  • 读取操作read(), getattr(), stat(), omap_get_values()
  • 写入操作:通过 Transactionwrite(), setattr(), omap_setkeys()
  • 事务提交queue_transaction()

12.3 调用流程

  1. 初始化:PG 创建时初始化 coll,加载时打开 ch
  2. 读取:通过 chstore 接口直接读取
  3. 写入:构建 PGTransaction,转换为 ObjectStore::Transaction,通过 PGBackend 提交
  4. 提交store->queue_transaction(ch, t) 提交到底层存储

12.4 相关文件

  • PG.h/cc:PG 基础类,管理集合和存储访问
  • PrimaryLogPG.h/cc:主日志 PG 实现,处理客户端操作
  • PGBackend.h:后端抽象接口
  • ReplicatedBackend.h/cc:复制后端实现
  • ECBackend.h/cc:EC 后端实现
  • ObjectStore.h:存储后端接口定义
  • os/Transaction.h:事务定义