Crimson OSD 代码架构分析:三副本和纠删码实现

Crimson OSD 代码架构分析:三副本和纠删码实现

1. 总体架构

Crimson OSD采用分层架构设计,核心组件包括:

1
2
3
4
5
OSD (osd.h/osd.cc)
└── PG (pg.h/pg.cc) - Placement Group,数据分布和复制的基本单位
└── PGBackend (pg_backend.h) - PG后端抽象基类
├── ReplicatedBackend (replicated_backend.h/cc) - 复制后端(三副本)
└── ECBackend (ec_backend.h/cc) - 纠删码后端

1.1 核心类层次结构

1
2
3
4
5
6
7
8
9
PGBackend (抽象基类)
├── ReplicatedBackend (复制后端)
│ ├── 实现三副本复制
│ ├── 管理副本确认机制
│ └── 处理PG committed-to (PCT) 更新
└── ECBackend (纠删码后端)
├── 实现纠删码编码/解码
├── 管理数据分片和校验分片
└── 处理分片恢复

2. 三副本(Replication)实现

2.1 架构设计

ReplicatedBackend 是复制存储的核心实现类,负责:

  1. 事务复制:将写操作复制到所有副本OSD
  2. 确认机制:等待所有副本确认操作完成
  3. 一致性保证:通过PG committed-to (PCT) 机制保证一致性

2.2 关键数据结构

2.2.1 pending_on_t - 待处理事务信息

1
2
3
4
5
6
7
class pending_on_t {
unsigned pending; // 待确认的副本数量
const eversion_t at_version; // 事务版本
const eversion_t last_complete; // 最后完成版本
crimson::osd::acked_peers_t acked_peers; // 已确认的副本列表
seastar::shared_promise<> all_committed; // 所有副本确认完成的promise
};

作用

  • 跟踪一个待处理事务的状态
  • 记录哪些副本已确认
  • 提供所有副本确认完成的future

2.2.2 pending_transactions_t - 待处理事务映射

1
2
using pending_transactions_t = std::map<ceph_tid_t, pending_on_t>;
pending_transactions_t pending_trans; // key为事务ID

作用

  • 维护所有待处理事务的映射
  • 通过事务ID快速查找事务状态

2.3 核心流程

2.3.1 提交事务流程(submit_transaction)

1
2
3
4
5
6
7
8
rep_op_fut_t submit_transaction(
const std::set<pg_shard_t> &pg_shards, // 所有副本分片
const hobject_t& hoid, // 对象句柄
crimson::osd::ObjectContextRef&& new_clone,
ceph::os::Transaction&& txn, // 事务
osd_op_params_t&& osd_op_p,
epoch_t min_epoch, epoch_t max_epoch,
std::vector<pg_log_entry_t>&& log_entries)

执行步骤

  1. 编码事务数据

    1
    2
    3
    bufferlist encoded_txn_p_bl, encoded_txn_d_bl;
    // 编码事务payload和数据部分
    txn.encode(encoded_txn_p_bl, encoded_txn_d_bl, pg.min_peer_features());
  2. 处理日志条目

    1
    2
    3
    4
    5
    6
    7
    8
    9
    bool 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);
  3. 创建待处理事务记录

    1
    2
    3
    4
    5
    6
    7
    const 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;
  4. 向所有副本发送复制操作

    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
    std::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));
    }
  5. 记录操作到PG日志

    1
    2
    3
    4
    5
    6
    7
    8
    pg.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);
  6. 在本地执行事务并等待确认

    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
    auto 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();
    });
  7. 返回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
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
void ReplicatedBackend::got_rep_op_reply(const MOSDRepOpReply& reply) {
auto it = pending_trans.find(reply.tid);
if (it != pending_trans.end()) {
auto& pending = it->second;
// 查找并更新对应副本的确认状态
for (auto& peer : pending.acked_peers) {
if (peer.shard == reply.from) {
peer.last_complete = reply.last_complete;
break;
}
}
// 减少待确认数量
pending.pending--;

// 如果所有副本都已确认
if (pending.pending == 0) {
// 更新PG的最后完成版本
pg.complete_write(pending.at_version, pending.last_complete);
// 设置promise,使等待的future ready
pending.all_committed.set_value();
// 清理待处理记录
pending_trans.erase(it);
// 可能触发PCT更新
maybe_kick_pct_update();
}
}
}

流程说明

  • 收到副本回复时,根据事务ID查找对应的待处理事务
  • 更新对应副本的确认状态和最后完成版本
  • 减少待确认副本数量
  • 当所有副本都确认后:
    • 更新PG的最后完成版本
    • 设置promise,使等待的future ready
    • 清理已完成的待处理事务记录
    • 可能触发PCT更新

2.3.3 PG Committed-To (PCT) 更新机制

目的:在IO暂停期间,更新PG的committed-to版本,保证一致性。

1
2
3
4
5
6
7
8
9
10
interruptible_future<> send_pct_update() {
// 向acting set中的所有其他OSD发送PCT更新
for (const auto& peer : acting_set) {
if (peer != whoami) {
auto msg = crimson::make_message<MOSDPGPCT>(
pgid, pg.get_pg_committed_to());
shard_services.send_to_osd(peer.osd, msg);
}
}
}

触发时机

  • 当repop队列为空时,启动PCT定时器
  • 定时器触发后发送PCT更新消息

2.4 智能发送机制

2.4.1 条件发送

根据副本是否拥有对象,采用不同的发送策略:

  1. 副本拥有对象 (pg.should_send_op(pg_shard, hoid) == true):

    • 发送完整的操作(send_op = true
    • 包含事务数据和日志条目
    • 副本直接执行操作
  2. 副本不拥有对象 (pg.should_send_op(pg_shard, hoid) == false):

    • 只发送日志条目(send_op = false
    • 不发送实际操作数据
    • 如果对象缺失,标记为需要后续推送(backfill)

优势

  • 减少网络传输:不拥有对象的副本不需要接收完整数据
  • 支持延迟推送:缺失的对象可以通过backfill机制后续推送

2.4.2 Backfill推送

对于缺失的对象,会加入backfill队列:

1
2
3
4
5
6
7
8
9
10
11
12
// Clone对象推送
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);
}

// 删除操作推送
if (!to_push_delete.empty()) {
pg.enqueue_delete_for_backfill(hoid, {}, to_push_delete);
}

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 epoch
    • min_epoch:最小epoch

标志

  • CEPH_OSD_FLAG_ACK:需要ACK确认
  • CEPH_OSD_FLAG_ONDISK:需要磁盘确认

消息格式

  • Tentacle格式(新格式):txn_payloaddata 分离
  • 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 是纠删码存储的核心实现类,负责:

  1. 数据编码:将数据分成多个数据分片和校验分片
  2. 数据解码:从部分分片恢复原始数据
  3. 分片管理:管理数据分片和校验分片的分布

3.2 关键概念

3.2.1 纠删码参数

  • K(数据分片数):原始数据分成的分片数量
  • M(校验分片数):生成的校验分片数量
  • 条带宽度(stripe_width):每个条带的大小
  • EC Profile:纠删码配置,包含K、M、算法等

3.2.2 数据分布

1
2
3
4
5
6
7
原始数据 (N bytes)

分成K个数据分片 (每个 N/K bytes)

生成M个校验分片 (每个 N/K bytes)

总共 K+M 个分片,分布在不同的OSD上

容错能力:可以容忍最多M个分片丢失,仍能恢复原始数据。

3.3 核心流程

3.3.1 读取流程(_read)

1
2
3
4
5
6
7
8
9
10
11
ll_read_ierrorator::future<ceph::bufferlist>
ECBackend::_read(const hobject_t& hoid,
const uint64_t off,
const uint64_t len,
const uint32_t flags)
{
// 1. 确定需要读取的分片
// 2. 从多个OSD并行读取分片
// 3. 如果某些分片不可用,使用其他分片恢复
// 4. 解码并返回原始数据
}

实现要点

  • 需要读取至少K个分片才能恢复数据
  • 可以并行从多个OSD读取
  • 如果某些分片丢失,使用纠删码算法恢复

3.3.2 写入流程(submit_transaction)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
rep_op_fut_t ECBackend::submit_transaction(
const std::set<pg_shard_t> &pg_shards, // 所有分片OSD
const hobject_t& hoid,
crimson::osd::ObjectContextRef&& new_clone,
ceph::os::Transaction&& txn,
osd_op_params_t&& osd_op_p,
epoch_t min_epoch, epoch_t max_epoch,
std::vector<pg_log_entry_t>&& log_entries)
{
// 1. 编码数据:将数据分成K个数据分片
// 2. 生成校验:计算M个校验分片
// 3. 分发分片:将K+M个分片发送到不同的OSD
// 4. 等待确认:等待所有分片写入完成
}

实现要点

  • 数据需要先编码成K+M个分片
  • 每个分片写入到不同的OSD
  • 需要等待所有分片写入完成

3.4 当前实现状态

根据代码分析,ECBackend目前是占位实现

1
2
3
4
5
6
7
8
9
ECBackend::_read(...) {
// todo
return seastar::make_ready_future<bufferlist>();
}

ECBackend::submit_transaction(...) {
// todo
return make_ready_future<rep_op_ret_t>(seastar::now(), seastar::now());
}

说明

  • ECBackend的框架已搭建,但具体实现尚未完成
  • 纠删码的编码/解码逻辑需要进一步实现
  • 分片管理和恢复机制需要完善

4. 后端选择机制

4.1 工厂方法(create)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
static std::unique_ptr<PGBackend> PGBackend::create(
pg_t pgid,
const pg_shard_t pg_shard,
const pg_pool_t& pool, // 存储池信息
crimson::osd::PG &pg,
crimson::os::CollectionRef coll,
crimson::osd::ShardServices& shard_services,
const ec_profile_t& ec_profile, // 纠删码配置
DoutPrefixProvider &dpp)
{
// 根据存储池类型选择后端
if (pool.is_replicated()) {
// 复制存储池 -> ReplicatedBackend
return std::make_unique<ReplicatedBackend>(
pgid, pg_shard, pg, coll, shard_services, dpp);
} else if (pool.is_erasure()) {
// 纠删码存储池 -> ECBackend
return std::make_unique<ECBackend>(
pg_shard.shard, coll, shard_services,
ec_profile, pool.get_stripe_width(), dpp);
}
}

选择逻辑

  • pool.is_replicated()ReplicatedBackend
  • pool.is_erasure()ECBackend

4.2 存储池配置

存储池类型在创建时确定:

  • 复制池pool_type = replicated,指定副本数(如size=3)
  • 纠删码池pool_type = erasure,指定EC profile(K、M、算法等)

5. 关键设计模式

5.1 策略模式(Strategy Pattern)

  • PGBackend 作为抽象策略接口
  • ReplicatedBackendECBackend 作为具体策略实现
  • 运行时根据存储池类型选择策略

5.2 工厂模式(Factory Pattern)

  • PGBackend::create() 作为工厂方法
  • 根据配置创建相应的后端实例

5.3 观察者模式(Observer Pattern)

  • pending_on_t 使用 seastar::shared_promise 实现异步通知
  • 副本确认时通知等待的future

6. 数据一致性保证

6.1 三副本一致性

  1. 写操作一致性

    • 主OSD接收写请求
    • 向所有副本发送复制操作(MOSDRepOp)
    • 在本地执行事务
    • 等待所有副本ACK确认(MOSDRepOpReply)
    • 更新PG的最后完成版本
    • 返回成功给客户端

    一致性保证

    • 所有副本必须确认操作完成
    • 使用事务ID(tid)跟踪每个操作
    • 通过版本号(at_version)保证顺序
  2. 读操作一致性

    • 从主OSD读取(或从最近的副本读取)
    • 保证读取到已提交的数据
    • 使用PG的committed-to版本判断数据可见性
  3. PCT(PG Committed-To)机制

    • 目的:在IO暂停期间,更新PG的committed-to版本
    • 触发:当repop队列为空时,启动PCT定时器
    • 执行:向acting set中的所有其他OSD发送PCT更新
    • 作用:保证所有副本知道已提交的版本范围
    1
    2
    3
    4
    5
    6
    void maybe_kick_pct_update() {
    if (pending_trans.empty() && !pct_timer.armed()) {
    // 队列为空且定时器未启动,启动PCT更新
    pct_timer.arm(std::chrono::milliseconds(pct_interval));
    }
    }
  4. Acting Set变化处理

    • 当PG的acting set发生变化时,取消所有待处理事务
    • 对所有待处理事务设置异常(actingset_changed
    • 清空待处理事务映射
    • 取消PCT更新

6.2 纠删码一致性

  1. 写操作

    • 编码数据成K+M个分片
    • 写入所有分片
    • 等待所有分片确认
  2. 读操作

    • 读取至少K个分片
    • 解码恢复原始数据
    • 如果分片丢失,使用其他分片恢复

7. 性能优化

7.1 三副本优化

  1. 并行发送:向所有副本并行发送复制操作
  2. 异步确认:使用future/promise异步等待确认
  3. 批量处理:可以批量处理多个操作

7.2 纠删码优化(待实现)

  1. 并行读取:从多个OSD并行读取分片
  2. 增量编码:只编码修改的部分
  3. 缓存分片:缓存常用的分片

8. 总结

8.1 三副本实现

  • 已完成:ReplicatedBackend实现完整
  • 核心功能:事务复制、确认机制、PCT更新
  • 消息机制:MOSDRepOp、MOSDRepOpReply、MOSDPGPCT

8.2 纠删码实现

  • ⚠️ 框架已搭建:ECBackend类结构完整
  • 实现待完成:编码/解码逻辑需要实现
  • 分片管理:分片分布和恢复机制需要完善

8.3 架构优势

  1. 清晰的抽象:PGBackend提供统一接口
  2. 灵活扩展:易于添加新的存储后端类型
  3. 异步设计:充分利用Seastar的异步能力
  4. 类型安全:使用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 - 对象上下文

文章互动

阅读 --

留言

0 条留言

正在加载留言…