Crimson OSD 代码架构分析:三副本和纠删码实现
1. 总体架构
Crimson OSD采用分层架构设计,核心组件包括:
1 | OSD (osd.h/osd.cc) |
1.1 核心类层次结构
1 | PGBackend (抽象基类) |
2. 三副本(Replication)实现
2.1 架构设计
ReplicatedBackend 是复制存储的核心实现类,负责:
- 事务复制:将写操作复制到所有副本OSD
- 确认机制:等待所有副本确认操作完成
- 一致性保证:通过PG committed-to (PCT) 机制保证一致性
2.2 关键数据结构
2.2.1 pending_on_t - 待处理事务信息
1 | class pending_on_t { |
作用:
- 跟踪一个待处理事务的状态
- 记录哪些副本已确认
- 提供所有副本确认完成的future
2.2.2 pending_transactions_t - 待处理事务映射
1 | using pending_transactions_t = std::map<ceph_tid_t, pending_on_t>; |
作用:
- 维护所有待处理事务的映射
- 通过事务ID快速查找事务状态
2.3 核心流程
2.3.1 提交事务流程(submit_transaction)
1 | rep_op_fut_t submit_transaction( |
执行步骤:
编码事务数据
1
2
3bufferlist encoded_txn_p_bl, encoded_txn_d_bl;
// 编码事务payload和数据部分
txn.encode(encoded_txn_p_bl, encoded_txn_d_bl, pg.min_peer_features());处理日志条目
1
2
3
4
5
6
7
8
9bool is_delete = false;
for (auto &le : log_entries) {
le.mark_unrollbackable(); // 标记为不可回滚
if (le.is_delete()) {
is_delete = true;
}
}
// 更新快照映射
co_await pg.update_snap_map(log_entries, txn);创建待处理事务记录
1
2
3
4
5
6
7const ceph_tid_t tid = shard_services.get_tid();
auto pending_txn = pending_trans.try_emplace(
tid,
pg_shards.size(), // 副本数量(如3)
osd_op_p.at_version, // 版本号
pg.get_last_complete() // 最后完成版本
).first;向所有副本发送复制操作
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
35std::vector<pg_shard_t> to_push_clone; // 需要推送的clone
std::vector<pg_shard_t> to_push_delete; // 需要推送的删除操作
for (const auto& pg_shard : pg_shards) {
if (pg_shard == whoami) {
continue; // 跳过自己
}
MURef<MOSDRepOp> m;
if (pg.should_send_op(pg_shard, hoid)) {
// 副本拥有对象,发送完整的操作
m = new_repop_msg(pg_shard, hoid, encoded_txn_p_bl,
encoded_txn_d_bl, osd_op_p,
min_epoch, map_epoch, log_entries, true, tid);
} else {
// 副本不拥有对象,只发送日志条目(不发送实际操作)
m = new_repop_msg(pg_shard, hoid, encoded_txn_p_bl,
encoded_txn_d_bl, osd_op_p,
min_epoch, map_epoch, log_entries, false, tid);
// 如果对象在副本上缺失,需要后续推送
if (pg.is_missing_on_peer(pg_shard, hoid)) {
if (_new_clone) {
to_push_clone.push_back(pg_shard);
}
if (is_delete) {
to_push_delete.push_back(pg_shard);
}
}
}
// 记录待确认的副本
pending_txn->second.acked_peers.push_back({pg_shard, eversion_t{}});
// 发送消息到副本OSD
sends->emplace_back(
shard_services.send_to_osd(pg_shard.osd, std::move(m), map_epoch));
}记录操作到PG日志
1
2
3
4
5
6
7
8pg.log_operation(
std::move(log_entries),
osd_op_p.pg_trim_to,
osd_op_p.at_version,
osd_op_p.pg_committed_to,
true,
txn,
false);在本地执行事务并等待确认
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
31auto all_completed = interruptor::make_interruptible(
shard_services.get_store().do_transaction(coll, std::move(txn))
).then_interruptible([pending_txn, this] {
// 如果副本数量为1(只有自己),直接完成
if (--pending_txn->second.pending == 0) {
pg.complete_write(pending_txn->second.at_version,
pending_txn->second.last_complete);
pending_txn->second.all_committed.set_value();
return seastar::now();
}
// 等待所有副本ACK(在got_rep_op_reply中处理)
return pending_txn->second.all_committed.get_shared_future();
}).then_interruptible([pending_txn, this, _new_clone, &hoid,
to_push_delete, to_push_clone] {
// 所有副本确认后,处理需要推送的对象
auto acked_peers = std::move(pending_txn->second.acked_peers);
pending_trans.erase(pending_txn);
// 如果有新clone需要推送,加入backfill队列
if (_new_clone && !to_push_clone.empty()) {
pg.enqueue_push_for_backfill(_new_clone->obs.oi.soid,
_new_clone->obs.oi.version,
to_push_clone);
}
// 如果有删除操作需要推送,加入backfill队列
if (!to_push_delete.empty()) {
pg.enqueue_delete_for_backfill(hoid, {}, to_push_delete);
}
// 可能触发PCT更新
maybe_kick_pct_update();
return seastar::now();
});返回future
1
2
3
4
5// 返回发送future和完成future
co_return std::make_tuple(
std::move(sends_complete), // 所有发送操作完成的future
std::move(all_completed) // 所有副本确认完成的future
);
2.3.2 副本确认流程(got_rep_op_reply)
1 | void ReplicatedBackend::got_rep_op_reply(const MOSDRepOpReply& reply) { |
流程说明:
- 收到副本回复时,根据事务ID查找对应的待处理事务
- 更新对应副本的确认状态和最后完成版本
- 减少待确认副本数量
- 当所有副本都确认后:
- 更新PG的最后完成版本
- 设置promise,使等待的future ready
- 清理已完成的待处理事务记录
- 可能触发PCT更新
2.3.3 PG Committed-To (PCT) 更新机制
目的:在IO暂停期间,更新PG的committed-to版本,保证一致性。
1 | interruptible_future<> send_pct_update() { |
触发时机:
- 当repop队列为空时,启动PCT定时器
- 定时器触发后发送PCT更新消息
2.4 智能发送机制
2.4.1 条件发送
根据副本是否拥有对象,采用不同的发送策略:
副本拥有对象 (
pg.should_send_op(pg_shard, hoid) == true):- 发送完整的操作(
send_op = true) - 包含事务数据和日志条目
- 副本直接执行操作
- 发送完整的操作(
副本不拥有对象 (
pg.should_send_op(pg_shard, hoid) == false):- 只发送日志条目(
send_op = false) - 不发送实际操作数据
- 如果对象缺失,标记为需要后续推送(backfill)
- 只发送日志条目(
优势:
- 减少网络传输:不拥有对象的副本不需要接收完整数据
- 支持延迟推送:缺失的对象可以通过backfill机制后续推送
2.4.2 Backfill推送
对于缺失的对象,会加入backfill队列:
1 | // Clone对象推送 |
2.5 消息类型
2.5.1 MOSDRepOp - 复制操作消息
内容:
- 事务数据(编码后的Transaction)
txn_payload:事务payload(Tentacle格式)data:事务数据(或pre-tentacle格式的完整事务)
- PG日志条目(
logbl) - 对象句柄(
hoid) - 版本信息
at_version:事务版本pg_trim_to:PG trim版本pg_committed_to:PG committed-to版本
- PG统计信息(
pg_stats) - Epoch信息
map_epoch:OSDMap epochmin_epoch:最小epoch
标志:
CEPH_OSD_FLAG_ACK:需要ACK确认CEPH_OSD_FLAG_ONDISK:需要磁盘确认
消息格式:
- Tentacle格式(新格式):
txn_payload和data分离 - Pre-tentacle格式(旧格式):所有内容都在
data中
2.5.2 MOSDRepOpReply - 复制操作回复
内容:
- 事务ID(
tid) - 发送者信息(
from) - 确认信息
last_complete:最后完成版本
- 错误信息(如果有)
2.5.3 MOSDPGPCT - PG Committed-To更新
内容:
- PG ID(
pgid) - Committed-to版本(
pg_committed_to)
用途:
- 在IO暂停期间更新PG的committed-to版本
- 保证所有副本知道已提交的版本范围
3. 纠删码(Erasure Code)实现
3.1 架构设计
ECBackend 是纠删码存储的核心实现类,负责:
- 数据编码:将数据分成多个数据分片和校验分片
- 数据解码:从部分分片恢复原始数据
- 分片管理:管理数据分片和校验分片的分布
3.2 关键概念
3.2.1 纠删码参数
- K(数据分片数):原始数据分成的分片数量
- M(校验分片数):生成的校验分片数量
- 条带宽度(stripe_width):每个条带的大小
- EC Profile:纠删码配置,包含K、M、算法等
3.2.2 数据分布
1 | 原始数据 (N bytes) |
容错能力:可以容忍最多M个分片丢失,仍能恢复原始数据。
3.3 核心流程
3.3.1 读取流程(_read)
1 | ll_read_ierrorator::future<ceph::bufferlist> |
实现要点:
- 需要读取至少K个分片才能恢复数据
- 可以并行从多个OSD读取
- 如果某些分片丢失,使用纠删码算法恢复
3.3.2 写入流程(submit_transaction)
1 | rep_op_fut_t ECBackend::submit_transaction( |
实现要点:
- 数据需要先编码成K+M个分片
- 每个分片写入到不同的OSD
- 需要等待所有分片写入完成
3.4 当前实现状态
根据代码分析,ECBackend目前是占位实现:
1 | ECBackend::_read(...) { |
说明:
- ECBackend的框架已搭建,但具体实现尚未完成
- 纠删码的编码/解码逻辑需要进一步实现
- 分片管理和恢复机制需要完善
4. 后端选择机制
4.1 工厂方法(create)
1 | static std::unique_ptr<PGBackend> PGBackend::create( |
选择逻辑:
pool.is_replicated()→ReplicatedBackendpool.is_erasure()→ECBackend
4.2 存储池配置
存储池类型在创建时确定:
- 复制池:
pool_type = replicated,指定副本数(如size=3) - 纠删码池:
pool_type = erasure,指定EC profile(K、M、算法等)
5. 关键设计模式
5.1 策略模式(Strategy Pattern)
- PGBackend 作为抽象策略接口
- ReplicatedBackend 和 ECBackend 作为具体策略实现
- 运行时根据存储池类型选择策略
5.2 工厂模式(Factory Pattern)
- PGBackend::create() 作为工厂方法
- 根据配置创建相应的后端实例
5.3 观察者模式(Observer Pattern)
- pending_on_t 使用
seastar::shared_promise实现异步通知 - 副本确认时通知等待的future
6. 数据一致性保证
6.1 三副本一致性
写操作一致性:
- 主OSD接收写请求
- 向所有副本发送复制操作(MOSDRepOp)
- 在本地执行事务
- 等待所有副本ACK确认(MOSDRepOpReply)
- 更新PG的最后完成版本
- 返回成功给客户端
一致性保证:
- 所有副本必须确认操作完成
- 使用事务ID(tid)跟踪每个操作
- 通过版本号(at_version)保证顺序
读操作一致性:
- 从主OSD读取(或从最近的副本读取)
- 保证读取到已提交的数据
- 使用PG的committed-to版本判断数据可见性
PCT(PG Committed-To)机制:
- 目的:在IO暂停期间,更新PG的committed-to版本
- 触发:当repop队列为空时,启动PCT定时器
- 执行:向acting set中的所有其他OSD发送PCT更新
- 作用:保证所有副本知道已提交的版本范围
1
2
3
4
5
6void maybe_kick_pct_update() {
if (pending_trans.empty() && !pct_timer.armed()) {
// 队列为空且定时器未启动,启动PCT更新
pct_timer.arm(std::chrono::milliseconds(pct_interval));
}
}Acting Set变化处理:
- 当PG的acting set发生变化时,取消所有待处理事务
- 对所有待处理事务设置异常(
actingset_changed) - 清空待处理事务映射
- 取消PCT更新
6.2 纠删码一致性
写操作:
- 编码数据成K+M个分片
- 写入所有分片
- 等待所有分片确认
读操作:
- 读取至少K个分片
- 解码恢复原始数据
- 如果分片丢失,使用其他分片恢复
7. 性能优化
7.1 三副本优化
- 并行发送:向所有副本并行发送复制操作
- 异步确认:使用future/promise异步等待确认
- 批量处理:可以批量处理多个操作
7.2 纠删码优化(待实现)
- 并行读取:从多个OSD并行读取分片
- 增量编码:只编码修改的部分
- 缓存分片:缓存常用的分片
8. 总结
8.1 三副本实现
- ✅ 已完成:ReplicatedBackend实现完整
- ✅ 核心功能:事务复制、确认机制、PCT更新
- ✅ 消息机制:MOSDRepOp、MOSDRepOpReply、MOSDPGPCT
8.2 纠删码实现
- ⚠️ 框架已搭建:ECBackend类结构完整
- ❌ 实现待完成:编码/解码逻辑需要实现
- ❌ 分片管理:分片分布和恢复机制需要完善
8.3 架构优势
- 清晰的抽象:PGBackend提供统一接口
- 灵活扩展:易于添加新的存储后端类型
- 异步设计:充分利用Seastar的异步能力
- 类型安全:使用C++类型系统保证正确性
9. 相关文件
9.1 核心文件
pg_backend.h/cc- PG后端抽象基类replicated_backend.h/cc- 复制后端实现ec_backend.h/cc- 纠删码后端实现pg.h/cc- Placement Group主类osd.h/cc- OSD主类
9.2 消息文件
messages/MOSDRepOp.h- 复制操作消息messages/MOSDRepOpReply.h- 复制操作回复messages/MOSDPGPCT.h- PG Committed-To更新
9.3 辅助文件
acked_peers.h- 已确认副本管理shard_services.h- 分片服务object_context.h- 对象上下文
正在加载留言…