首页/目录/全部文章

全部文章

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

笔记列表

SeaStore 硬盘存储调用机制

SeaStore 硬盘存储调用机制

1. 概述

SeaStore 通过多层抽象来访问底层存储设备,最终使用 Seastar 框架的 DMA(Direct Memory Access)接口进行实际的 I/O 操作。本文档详细说明从应用层到硬件层的完整调用路径。

2. 调用层次结构

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
应用层 (SeaStore API)


事务管理层 (TransactionManager)


缓存层 (Cache)


Extent放置管理层 (ExtentPlacementManager)

├─► 段式设备路径
│ │
│ ▼
│ SegmentManager (Block/ZBD)
│ │
│ ▼
│ Segment::write() / Segment::read()
│ │
│ ▼
│ seastar::file::dma_write() / dma_read()

└─► 随机块设备路径


RandomBlockManager


seastar::file::dma_write() / dma_read()

3. 设备抽象层

3.1 Device 接口

位置: device.h

Device 是存储设备的抽象接口,定义了统一的读写接口:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
class Device {
public:
// 读取接口
virtual read_ertr::future<> read(
paddr_t addr, // 物理地址
size_t len, // 读取长度
ceph::bufferptr &out // 输出缓冲区
) = 0;

// 设备类型
virtual device_type_t get_device_type() const = 0;

// 后端类型(SEGMENTED 或 RANDOM_BLOCK)
virtual backend_type_t get_backend_type() const = 0;

// 块大小
virtual extent_len_t get_block_size() const = 0;
};

3.2 设备实现类型

SeaStore 支持两种主要的设备类型:

  1. 段式设备 (SEGMENTED):

    • BlockSegmentManager: 普通块设备
    • ZBDSegmentManager: Zoned Block Device (ZBD)
    • EphemeralSegmentManager: 临时内存设备(测试用)
  2. 随机块设备 (RANDOM_BLOCK):

    • RandomBlockManager: 随机访问块设备

4. 段式设备调用路径

4.1 写入路径

4.1.1 调用链

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
TransactionManager::submit_transaction()


Journal::submit_record()


SegmentedJournal::record_submitter::submit()


SegmentAllocator::allocate_segment()


SegmentProvider::allocate_segment()


SegmentManager::open(segment_id)


Segment::write(offset, bufferlist)


BlockSegmentManager::segment_write(paddr, bufferlist)


do_writev(device_id, device, offset, bufferlist, block_size)


seastar::file::dma_write(offset, iovec)


Linux 系统调用: pwritev() / io_submit()

4.1.2 关键代码

Segment::write() (segment_manager/block.cc:372):

1
2
3
4
5
6
7
8
9
10
11
12
13
14
Segment::write_ertr::future<> BlockSegment::write(
segment_off_t offset, ceph::bufferlist bl)
{
auto paddr = paddr_t::make_seg_paddr(id, offset);
// 验证参数
if (offset < write_pointer ||
offset % manager.superblock.block_size != 0 ||
bl.length() % manager.superblock.block_size != 0) {
return crimson::ct_error::invarg::make();
}

write_pointer = offset + bl.length();
return manager.segment_write(paddr, bl);
}

BlockSegmentManager::segment_write() (segment_manager/block.cc:425):

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
Segment::write_ertr::future<> BlockSegmentManager::segment_write(
paddr_t addr,
ceph::bufferlist bl,
bool ignore_check)
{
assert(addr.get_device_id() == get_device_id());
assert((bl.length() % superblock.block_size) == 0);
stats.data_write.increment(bl.length());

// 调用底层写入函数
return do_writev(
get_device_id(),
device, // seastar::file 对象
get_offset(addr), // 物理偏移
std::move(bl),
superblock.block_size);
}

do_writev() (segment_manager/block.cc:87):

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
static write_ertr::future<> do_writev(
device_id_t device_id,
seastar::file &device, // Seastar 文件对象
uint64_t offset,
bufferlist&& bl,
size_t block_size)
{
// 确保缓冲区按块大小对齐
bl.rebuild_aligned(block_size);

return seastar::do_with(
bl.prepare_iovs(), // 准备 I/O 向量
std::move(bl),
[&device, device_id, offset](auto& iovs, auto& bl)
{
// 并行写入多个 I/O 向量
return write_ertr::parallel_for_each(
iovs,
[&device, device_id, offset](auto& p) mutable
{
auto off = offset + p.offset;
auto len = p.length;
auto& iov = p.iov;

// 调用 Seastar 的 DMA 写入接口
return device.dma_write(off, std::move(iov))
.handle_exception([device_id, off, len](auto e) {
ERROR("dma_write got error -- {}", e);
return crimson::ct_error::input_output_error::make();
})
.then([device_id, off, len](size_t written) {
if (written != len) {
ERROR("dma_write len inconsistent");
return crimson::ct_error::input_output_error::make();
}
return write_ertr::now();
});
});
});
}

4.2 读取路径

4.2.1 调用链

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
Cache::read_extent()


ExtentPlacementManager::read(paddr, len, bufferptr)


Device::read(paddr, len, bufferptr)


SegmentManager::read(paddr, len, bufferptr)


do_read(device_id, device, offset, len, bufferptr)


seastar::file::dma_read(offset, buffer, len)


Linux 系统调用: preadv() / io_submit()

4.2.2 关键代码

do_read() (segment_manager/block.cc:137):

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
static read_ertr::future<> do_read(
device_id_t device_id,
seastar::file &device,
uint64_t offset,
size_t len,
bufferptr &bptr)
{
assert(len <= bptr.length());

// 调用 Seastar 的 DMA 读取接口
return device.dma_read(
offset,
bptr.c_str(),
len
).handle_exception([device_id, offset, len](auto e) {
ERROR("dma_read got error -- {}", e);
return crimson::ct_error::input_output_error::make();
}).then([device_id, offset, len](auto result) {
if (result != len) {
ERROR("read len inconsistent");
return crimson::ct_error::input_output_error::make();
}
return read_ertr::now();
});
}

5. Seastar DMA 接口

5.1 文件对象创建

位置: segment_manager/block.cc:258

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
static open_device_ret open_device(const std::string &path)
{
// 获取文件状态
return seastar::file_stat(path, seastar::follow_symlink::yes)
.then([&path](auto stat) mutable {
// 以 DMA 模式打开文件
return seastar::open_file_dma(
path,
seastar::open_flags::rw | seastar::open_flags::dsync
).then([stat](auto file) mutable {
return file.size().then([stat, file](auto size) mutable {
stat.size = size;
// 获取 DMA 写入对齐要求
stat.block_size = file.disk_write_dma_alignment();
return std::make_pair(file, stat);
});
});
});
}

5.2 DMA 操作特点

  1. 零拷贝: DMA 操作直接从用户缓冲区传输到内核,减少数据拷贝
  2. 异步 I/O: 所有操作都是异步的,返回 seastar::future
  3. 对齐要求: 缓冲区和偏移必须满足 DMA 对齐要求
  4. 批量操作: 支持 writev/readv 进行批量 I/O

5.3 Seastar 底层实现

Seastar 的 dma_read/dma_write 最终会调用:

  1. Linux AIO: 使用 io_submit()io_getevents() 进行异步 I/O
  2. io_uring: 如果系统支持,使用更高效的 io_uring 接口
  3. 直接 I/O: 绕过页缓存,直接访问存储设备

6. 随机块设备调用路径

6.1 RandomBlockManager

位置: random_block_manager/block_rb_manager.cc

随机块设备也使用类似的 DMA 接口:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
RandomBlockManager::write_ertr::future<>
RandomBlockManager::write(
paddr_t addr,
ceph::bufferptr &bptr)
{
// 计算物理偏移
auto offset = get_offset(addr);

// 调用 DMA 写入
return device.dma_write(offset, bptr.c_str(), bptr.length())
.handle_exception([](auto e) {
return crimson::ct_error::input_output_error::make();
});
}

7. 设备初始化流程

7.1 设备创建

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
Device::make_device(device_path, device_type)

├─► 如果是 SEGMENTED 类型
│ │
│ ▼
│ SegmentManager::get_segment_manager()
│ │
│ ├─► BlockSegmentManager::create()
│ ├─► ZBDSegmentManager::create()
│ └─► EphemeralSegmentManager::create()

└─► 如果是 RANDOM_BLOCK 类型


RandomBlockManager::create()

7.2 设备挂载

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
Device::mount()


BlockSegmentManager::mount()


open_device(device_path) // 打开文件


read_superblock() // 读取超级块


tracker->read_in() // 读取段状态跟踪器


设备就绪,可以接受 I/O

8. I/O 优化机制

8.1 批量 I/O

SeaStore 使用 writev/readv 进行批量 I/O 操作:

1
2
3
4
5
6
7
// 准备多个 I/O 向量
auto iovs = bl.prepare_iovs();

// 并行执行多个 I/O
return parallel_for_each(iovs, [&device](auto& iov) {
return device.dma_write(offset, iov);
});

8.2 对齐优化

所有 I/O 操作都按照块大小对齐:

1
2
3
4
5
6
// 确保缓冲区对齐
bl.rebuild_aligned(block_size);

// 验证偏移对齐
assert(offset % block_size == 0);
assert(len % block_size == 0);

8.3 异步 I/O

所有 I/O 操作都是异步的,使用 Seastar 的 future/promise 机制:

1
2
3
4
5
return device.dma_write(offset, buffer, len)
.then([](size_t written) {
// I/O 完成后的处理
return seastar::now();
});

9. 错误处理

9.1 I/O 错误类型

  • input_output_error: 底层 I/O 错误
  • invarg: 无效参数(如未对齐)
  • enospc: 空间不足
  • erange: 地址超出范围

9.2 错误处理机制

1
2
3
4
5
6
7
8
9
10
11
12
13
return device.dma_write(offset, buffer, len)
.handle_exception([](auto e) {
// 捕获异常并转换为错误类型
ERROR("dma_write got error -- {}", e);
return crimson::ct_error::input_output_error::make();
})
.then([](size_t written) {
// 验证写入长度
if (written != len) {
return crimson::ct_error::input_output_error::make();
}
return write_ertr::now();
});

10. 性能特性

10.1 DMA 优势

  1. 零拷贝: 数据直接从用户缓冲区传输到设备,无需内核拷贝
  2. 异步操作: 不阻塞线程,充分利用 CPU
  3. 批量操作: 支持向量 I/O,减少系统调用次数

10.2 Seastar 优化

  1. 事件驱动: 基于 epoll/io_uring 的事件驱动模型
  2. 协程支持: 使用协程简化异步代码
  3. Shard 隔离: 每个 CPU 核心独立处理 I/O

11. 总结

SeaStore 的存储调用机制具有以下特点:

  1. 多层抽象: 从 Device 接口到具体实现,再到 Seastar DMA 接口
  2. 异步 I/O: 所有操作都是异步的,提高并发性能
  3. 零拷贝: 使用 DMA 减少数据拷贝开销
  4. 类型支持: 支持段式设备和随机块设备
  5. 错误处理: 完善的错误处理和恢复机制

通过这种设计,SeaStore 能够充分利用现代存储设备的性能,同时保持代码的清晰和可维护性。

Crimson 代码架构文档

Crimson 代码架构文档

目录

  1. 概述
  2. 核心设计理念
  3. 目录结构
  4. 主要模块
  5. 关键组件
  6. 数据流
  7. 依赖关系

概述

Crimson 是 Ceph 的下一代 OSD 实现,基于 Seastar 异步框架构建。它采用事件驱动、无锁编程模型,旨在提供更高的性能和更好的可扩展性。

主要特点

  • 基于 Seastar: 使用 Seastar 异步框架实现高性能 I/O
  • 无锁设计: 采用单线程、事件驱动的架构
  • 模块化: 清晰的模块划分,便于维护和扩展
  • 类型安全: 使用现代 C++ 特性,提供类型安全的异步操作

核心设计理念

1. 异步编程模型

  • 使用 seastar::futureseastar::promise 进行异步操作
  • 使用 crimson::interruptible_future 支持可中断的异步操作
  • 使用 crimson::errorator 处理错误类型

2. 单线程架构

  • 每个 CPU 核心运行一个 Seastar 线程
  • 通过消息传递进行跨核心通信
  • 避免锁竞争,提高性能

3. 存储抽象

  • FuturizedStore: 存储接口的异步版本
  • Seastore: 基于 Seastar 的新存储后端
  • AlienStore: 兼容传统 BlueStore 的适配器

目录结构

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
67
68
69
70
71
72
73
74
75
76
77
crimson/
├── admin/ # 管理接口
│ ├── admin_socket.h/cc
│ ├── osd_admin.h/cc
│ └── pg_commands.h/cc

├── auth/ # 认证模块
│ ├── AuthClient.h
│ ├── AuthServer.h
│ ├── DummyAuth.h
│ └── KeyRing.h/cc

├── common/ # 通用工具和基础设施
│ ├── errorator.h # 错误处理
│ ├── interruptible_future.h # 可中断的 future
│ ├── coroutine.h # 协程支持
│ ├── operation.h/cc # 操作抽象
│ ├── throttle.h/cc # 限流器
│ ├── log.h/cc # 日志系统
│ └── ...

├── crush/ # CRUSH 算法
│ └── CrushLocation.h/cc

├── mgr/ # Manager 客户端
│ └── client.h/cc

├── mon/ # Monitor 客户端
│ └── MonClient.h/cc

├── net/ # 网络层
│ ├── Messenger.h/cc # 消息传递
│ ├── Socket.h/cc # Socket 抽象
│ ├── ProtocolV2.h/cc # 协议实现
│ └── ...

├── os/ # 对象存储层
│ ├── futurized_store.h/cc # 存储接口
│ ├── seastore/ # Seastore 存储后端
│ │ ├── seastore.h/cc
│ │ ├── transaction_manager.h/cc
│ │ ├── cache.h/cc
│ │ ├── journal.h/cc
│ │ ├── lba_manager.h/cc
│ │ ├── onode_manager.h
│ │ ├── omap_manager.h/cc
│ │ ├── collection_manager.h/cc
│ │ └── ...
│ ├── cyanstore/ # CyanStore (测试用)
│ └── alienstore/ # AlienStore (BlueStore 适配器)

├── osd/ # OSD 核心模块
│ ├── osd.h/cc # OSD 主类
│ ├── pg.h/cc # Placement Group
│ ├── pg_map.h/cc # PG 映射管理
│ ├── shard_services.h/cc # Shard 服务
│ ├── pg_backend.h/cc # PG 后端
│ ├── recovery_backend.h/cc # 恢复后端
│ ├── ops_executer.h/cc # 操作执行器
│ ├── osd_operations/ # OSD 操作
│ │ ├── client_request.h/cc
│ │ ├── peering_event.h/cc
│ │ ├── replicated_request.h/cc
│ │ └── ...
│ ├── scheduler/ # 调度器
│ │ ├── scheduler.h/cc
│ │ └── mclock_scheduler.h/cc
│ ├── scrub/ # Scrub 功能
│ │ ├── pg_scrubber.h/cc
│ │ ├── scrub_machine.h/cc
│ │ └── scrub_validator.h/cc
│ └── ...

└── tools/ # 工具程序
├── store_nbd/ # NBD 工具
├── store_bench/ # 性能测试工具
└── objectstore/ # 对象存储工具

主要模块

1. OSD 模块 (crimson/osd/)

1.1 OSD 主类 (osd.h/cc)

  • 职责: OSD 进程的主入口和生命周期管理
  • 关键组件:
    • cluster_msgr: 集群内部通信
    • public_msgr: 客户端/监控通信
    • hb_front_msgr / hb_back_msgr: 心跳消息
    • monc: Monitor 客户端
    • mgrc: Manager 客户端
    • store: 存储后端

1.2 Placement Group (pg.h/cc)

  • 职责: 管理单个 PG 的状态和操作
  • 关键功能:
    • PG 状态机管理
    • 客户端请求处理
    • Peering 和恢复
    • Scrub 操作
    • 快照管理

1.3 Shard Services (shard_services.h/cc)

  • 职责: 为每个 shard 提供服务接口
  • 功能:
    • OSDMap 管理
    • PG 创建和删除
    • 消息分发

1.4 OSD Operations (osd_operations/)

各种 OSD 操作的实现:

  • client_request: 客户端请求处理
  • peering_event: Peering 事件处理
  • replicated_request: 副本请求处理
  • background_recovery: 后台恢复
  • scrub_events: Scrub 事件

2. 存储层 (crimson/os/)

2.1 FuturizedStore (futurized_store.h/cc)

  • 职责: 存储接口的异步抽象
  • 接口:
    • create(): 创建集合
    • open(): 打开集合
    • read() / write(): 异步读写
    • omap_*(): OMap 操作

2.2 Seastore (os/seastore/)

基于 Seastar 的新存储后端:

核心组件:

  • TransactionManager (transaction_manager.h/cc): 事务管理器

    • 提供事务接口
    • 管理缓存和日志
    • 处理事务提交和回滚
  • Cache (cache.h/cc): 缓存管理

    • 管理内存中的 extent
    • LRU 淘汰策略
    • 脏页管理
  • Journal (journal.h/cc): 日志系统

    • 写前日志 (WAL)
    • 日志回放
    • 日志清理
  • LBAManager (lba_manager.h/cc): 逻辑块地址管理

    • 逻辑地址到物理地址映射
    • B-tree 索引
  • OnodeManager (onode_manager.h): 对象节点管理

    • 对象元数据管理
    • 对象树结构
  • OMapManager (omap_manager.h/cc): OMap 管理

    • 对象扩展属性
    • 键值对存储
  • CollectionManager (collection_manager.h/cc): 集合管理

    • PG 集合管理
    • 集合元数据

2.3 AlienStore (os/alienstore/)

  • 职责: 兼容传统 BlueStore 的适配器
  • 功能: 在 Seastar 环境中运行 BlueStore

2.4 CyanStore (os/cyanstore/)

  • 职责: 内存存储后端(主要用于测试)

3. 网络层 (crimson/net/)

3.1 Messenger (Messenger.h/cc)

  • 职责: 消息传递抽象
  • 功能:
    • 连接管理
    • 消息发送和接收
    • 协议处理

3.2 Socket (Socket.h/cc)

  • 职责: Socket 抽象
  • 实现: 基于 Seastar 的异步 Socket

3.3 ProtocolV2 (ProtocolV2.h/cc)

  • 职责: Ceph 协议 V2 实现
  • 功能: 消息序列化、压缩、加密

4. 通用模块 (crimson/common/)

4.1 Errorator (errorator.h)

  • 职责: 类型安全的错误处理
  • 特点: 编译时错误类型检查

4.2 InterruptibleFuture (interruptible_future.h)

  • 职责: 可中断的异步操作
  • 用途: 支持操作取消和中断

4.3 Operation (operation.h/cc)

  • 职责: 操作抽象基类
  • 功能: 操作跟踪、阻塞管理

4.4 Throttle (throttle.h/cc)

  • 职责: 限流器
  • 用途: 控制并发操作数量

4.5 Log (log.h/cc)

  • 职责: 日志系统
  • 特点: 集成 Seastar 日志后端

5. 调度器 (crimson/osd/scheduler/)

5.1 Scheduler (scheduler.h/cc)

  • 职责: 操作调度抽象
  • 功能: 优先级调度、限流

5.2 MClockScheduler (mclock_scheduler.h/cc)

  • 职责: MClock 调度算法实现
  • 特点: 支持 QoS、公平性保证

关键组件

1. OSD 启动流程

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
main() (main.cc)

初始化 Seastar 应用

创建 OSD 实例 (OSD::create())

初始化 Messenger (cluster_msgr, public_msgr)

连接 Monitor (monc->start())

加载 OSDMap

初始化存储后端 (store->mount())

启动 PG 管理 (pg_map)

开始服务

2. 客户端请求处理流程

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
客户端请求到达 (public_msgr)

OSD::ms_dispatch()

根据消息类型分发

ClientRequest 操作

PG::handle_op()

OpsExecuter::execute_op()

PGBackend::submit_transaction()

Store::queue_transaction()

返回响应

3. Peering 流程

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

PeeringEvent 操作

PG::handle_peering_event()

PeeringState 状态机

根据状态执行相应操作

更新 PG 状态

4. 恢复流程

1
2
3
4
5
6
7
8
9
10
11
12
13
触发恢复

BackgroundRecovery 操作

PGRecovery::start_recovery_ops()

RecoveryBackend::start_recovery_ops()

ReplicatedRecoveryBackend::recover_object()

从其他 OSD 拉取对象

写入本地存储

5. Seastore 事务流程

1
2
3
4
5
6
7
8
9
10
11
12
13
创建事务 (TransactionManager::create_transaction())

读取/修改 extents

提交事务 (TransactionManager::submit_transaction())

写入 Journal

更新 LBA 映射

更新 Cache

异步刷新到设备

数据流

写操作数据流

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

Public Messenger

OSD::ms_dispatch()

ClientRequest

PG::handle_op()

OpsExecuter::execute_op()

PGBackend::submit_transaction()

Seastore::queue_transaction()

TransactionManager::submit_transaction()

Journal::submit_record()

设备写入

读操作数据流

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

Public Messenger

OSD::ms_dispatch()

ClientRequest

PG::handle_op()

OpsExecuter::execute_op()

Seastore::read()

TransactionManager::read_extent()

Cache::get_extent()

设备读取 (如需要)

返回数据

依赖关系

核心依赖

  • Seastar: 异步框架
  • Boost: C++ 库支持
  • Ceph Common: Ceph 通用库

模块依赖关系

1
2
3
4
5
6
7
8
9
10
11
12
13
crimson-osd
├── crimson-admin
├── crimson-os
│ ├── crimson-seastore
│ ├── crimson-cyanstore
│ └── crimson-alienstore (可选)
├── crimson
│ ├── crimson-common
│ ├── crimson-auth
│ ├── crimson-mon
│ ├── crimson-mgr
│ └── crimson-net
└── crimson-common

编译配置

  • WITH_CRIMSON=1: 启用 Crimson 支持
  • WITH_BLUESTORE: 启用 BlueStore 支持(AlienStore)
  • WITH_TESTS: 启用测试支持

关键设计模式

1. Future/Promise 模式

  • 所有异步操作返回 seastar::future
  • 使用 co_await.then() 处理异步结果

2. 操作模式 (Operation Pattern)

  • 所有操作继承自 Operation
  • 支持操作跟踪、阻塞、取消

3. 状态机模式

  • PG 使用状态机管理状态转换
  • PeeringState 处理 PG 状态

4. 策略模式

  • PGBackend 抽象,支持 ReplicatedBackend 和 ECBackend
  • Scheduler 抽象,支持不同调度算法

5. 适配器模式

  • AlienStore 适配 BlueStore
  • FuturizedStore 适配传统存储接口

性能优化

1. 零拷贝

  • 使用 Seastar 的零拷贝网络栈
  • 减少内存拷贝开销

2. 批处理

  • 批量提交事务
  • 合并小 I/O

3. 预取

  • 智能预取策略
  • 减少 I/O 等待

4. 缓存

  • 多级缓存
  • LRU 淘汰策略

测试

测试目录结构

1
2
3
4
5
6
7
test/crimson/
├── test_*.cc # 单元测试
├── seastore/ # Seastore 测试
│ ├── test_seastore.cc
│ ├── test_transaction_manager.cc
│ └── ...
└── gtest_seastar.h/cc # GoogleTest + Seastar 集成

测试框架

  • GoogleTest: 单元测试框架
  • Seastar Test: Seastar 测试支持
  • CBT: 性能测试工具

扩展点

1. 存储后端

  • 实现 FuturizedStore 接口
  • 支持新的存储引擎

2. 调度器

  • 实现 Scheduler 接口
  • 支持新的调度算法

3. 网络协议

  • 扩展 Messenger 实现
  • 支持新协议版本

参考资料


版本信息

  • 文档生成时间: 2024
  • 基于代码库: devCephMain
  • Crimson 版本: 基于当前代码库

维护者

本文档应随代码更新而更新。主要更新点:

  • 新增模块时更新目录结构
  • 架构变更时更新设计说明
  • API 变更时更新接口文档

Crimson OSD 三副本实现分析

Crimson OSD 三副本实现分析

概述

虽然 main.cc 是 Crimson OSD 的入口点,但三副本(replication)的实现并不在 main.cc 中。main.cc 主要负责 OSD 的启动和初始化,而三副本的核心逻辑在 ReplicatedBackend 类中实现。

架构概览

1
2
3
4
5
6
7
main.cc (启动 OSD)

OSD::start()

PG 创建和管理

ReplicatedBackend (三副本实现)

三副本实现位置

1. 核心类:ReplicatedBackend

文件位置

  • cephMain/src/crimson/osd/replicated_backend.h
  • cephMain/src/crimson/osd/replicated_backend.cc

继承关系

1
ReplicatedBackend : public PGBackend

2. 关键方法:submit_transaction

这是三副本写入的核心方法,负责将事务提交到所有副本。

三副本工作流程

1. 副本集合确定

副本集合由 PG 的 acting set 决定,通常包含 3 个 OSD:

  • Primary OSD(主副本,通常是当前 OSD)
  • Secondary OSD(第二副本)
  • Tertiary OSD(第三副本)
1
2
// 在 ReplicatedBackend::submit_transaction 中
const std::set<pg_shard_t> &pg_shards // 包含所有副本的 OSD 列表

2. 写操作流程

步骤 1: 创建待处理事务

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
// replicated_backend.cc:94-119
ReplicatedBackend::rep_op_fut_t
ReplicatedBackend::submit_transaction(
const std::set<pg_shard_t> &pg_shards, // 副本集合(通常 3 个)
const hobject_t& hoid,
crimson::osd::ObjectContextRef &&new_clone,
ceph::os::Transaction&& t,
osd_op_params_t&& opp,
epoch_t min_epoch, epoch_t map_epoch,
std::vector<pg_log_entry_t>&& logv)
{
// 1. 生成事务 ID
const ceph_tid_t tid = shard_services.get_tid();

// 2. 创建待处理事务记录
auto pending_txn = pending_trans.try_emplace(
tid,
pg_shards.size(), // 待确认的副本数量
osd_op_p.at_version,
pg.get_last_complete()
).first;

// 3. 编码事务
bufferlist encoded_txn_p_bl, encoded_txn_d_bl;
txn.encode(encoded_txn_p_bl, encoded_txn_d_bl, pg.min_peer_features());
}

步骤 2: 发送到所有副本

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
// replicated_backend.cc:136-167
// 遍历所有副本 OSD
for (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(..., false, tid);
}

// 记录待确认的副本
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));
}

步骤 3: 本地提交

1
2
3
4
5
6
7
8
9
10
// replicated_backend.cc:169-176
// 主副本本地提交事务
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);

步骤 4: 等待所有副本确认

1
2
3
4
5
6
7
8
9
10
11
12
13
// replicated_backend.cc:178-197
auto all_completed = interruptor::make_interruptible(
shard_services.get_store().do_transaction(coll, std::move(txn))
).then_interruptible([this, peers=pending_txn->second.weak_from_this()] {
if (--peers->pending == 0) {
// 所有副本已确认(包括自己)
pg.complete_write(peers->at_version, peers->last_complete);
peers->all_committed.set_value();
return seastar::now();
}
// 等待其他副本的确认
return peers->all_committed.get_shared_future();
});

3. 副本确认处理

当副本 OSD 完成写入后,会发送 MOSDRepOpReply 消息:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
// replicated_backend.cc:231-249
void ReplicatedBackend::got_rep_op_reply(const MOSDRepOpReply& reply)
{
// 查找对应的事务
auto found = pending_trans.find(reply.get_tid());
auto& peers = found->second;

// 更新副本的确认状态
for (auto& peer : peers.acked_peers) {
if (peer.shard == reply.from) {
peer.last_complete_ondisk = reply.get_last_complete_ondisk();
pg.update_peer_last_complete_ondisk(
peer.shard, peer.last_complete_ondisk);

// 检查是否所有副本都已确认
if (--peers.pending == 0) {
pg.complete_write(peers.at_version, peers.last_complete);
peers.all_committed.set_value(); // 唤醒等待的 future
}
}
}
}

关键数据结构

1. pending_on_t - 待处理事务状态

1
2
3
4
5
6
7
8
// replicated_backend.h:51-71
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
};

2. acked_peers_t - 已确认副本列表

1
2
3
4
5
6
// acked_peers.h:9-13
struct peer_shard_t {
pg_shard_t shard; // 副本 OSD
eversion_t last_complete_ondisk; // 磁盘上的最后完成版本
};
using acked_peers_t = std::vector<peer_shard_t>;

3. pending_transactions_t - 待处理事务映射

1
2
3
// replicated_backend.h:72
using pending_transactions_t = std::map<ceph_tid_t, pending_on_t>;
pending_transactions_t pending_trans; // 事务 ID -> 状态

三副本保证机制

1. 同步写入

  • 主副本必须等待所有副本(包括自己)确认后才完成写操作
  • 使用 shared_promiseshared_future 实现多等待者同步

2. 版本控制

  • 每个写操作都有唯一的版本号(at_version
  • 副本必须按顺序应用操作,保证一致性

3. 故障处理

1
2
3
4
5
6
7
8
9
// replicated_backend.cc:221-229
void ReplicatedBackend::on_actingset_changed(bool same_primary)
{
// 当副本集合改变时,取消所有待处理的事务
for (auto& [tid, pending_txn] : pending_trans) {
pending_txn.all_committed.set_exception(e_actingset_changed);
}
pending_trans.clear();
}

消息类型

1. MOSDRepOp - 副本操作请求

主副本发送给其他副本的操作请求,包含:

  • 事务数据(txn_payloaddata
  • 日志条目(log_entries
  • 版本信息
  • PG 统计信息

2. MOSDRepOpReply - 副本操作响应

副本 OSD 完成写入后发送的确认消息,包含:

  • 事务 ID(tid
  • 最后完成版本(last_complete_ondisk
  • 操作结果

与 main.cc 的关系

虽然 main.cc 不直接实现三副本逻辑,但它负责:

  1. 初始化 OSD
1
2
3
4
5
// main.cc:210-213
crimson::osd::OSD osd(
whoami, nonce, std::ref(should_stop.abort_source()),
std::ref(*store), cluster_msgr, client_msgr,
hb_front_msgr, hb_back_msgr);
  1. 启动 OSD
1
2
// main.cc:239
osd.start().get();
  1. 创建 Messenger
1
2
3
4
// main.cc:189-204
// 创建集群通信和客户端通信的 Messenger
crimson::net::MessengerRef cluster_msgr, client_msgr;
crimson::net::MessengerRef hb_front_msgr, hb_back_msgr;

这些组件为三副本通信提供了基础设施。

完整写操作流程示例

假设有 3 个副本:OSD 0(主)、OSD 1、OSD 2

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
1. 客户端发送写请求到 OSD 0(主副本)

2. OSD 0 的 PG 处理请求,创建事务

3. ReplicatedBackend::submit_transaction()
├─ 创建 pending_txn (pending = 2,等待 OSD 1 和 OSD 2)
├─ 编码事务数据
├─ 发送 MOSDRepOp 到 OSD 1
├─ 发送 MOSDRepOp 到 OSD 2
├─ 本地提交事务
└─ 等待所有副本确认

4. OSD 1 接收 MOSDRepOp,执行写入,发送 MOSDRepOpReply

5. OSD 2 接收 MOSDRepOp,执行写入,发送 MOSDRepOpReply

6. OSD 0 收到所有回复
├─ got_rep_op_reply() 处理每个回复
├─ pending 减为 0
└─ all_committed.set_value() 唤醒等待

7. 写操作完成,返回客户端

性能优化

1. 异步发送

所有副本操作消息异步发送,不阻塞主流程:

1
2
sends->emplace_back(
shard_services.send_to_osd(pg_shard.osd, std::move(m), map_epoch));

2. 批量处理

使用 when_all_succeed 等待所有发送完成:

1
2
3
auto sends_complete = seastar::when_all_succeed(
sends->begin(), sends->end()
);

3. 条件发送

根据副本状态决定发送完整操作还是简化操作:

1
2
3
4
5
if (pg.should_send_op(pg_shard, hoid)) {
// 发送完整操作
} else {
// 发送简化操作(副本已有数据)
}

总结

三副本实现的核心特点:

  1. 位置:主要在 ReplicatedBackend 类中,不在 main.cc
  2. 机制:主副本发送操作到所有副本,等待所有确认
  3. 同步:使用 shared_promise/shared_future 实现同步
  4. 容错:副本集合变化时取消待处理事务
  5. 性能:异步发送,批量等待,条件优化

main.cc 的作用是启动整个系统,为三副本提供运行环境(OSD、Messenger、Store 等),但具体的副本逻辑由 ReplicatedBackend 实现。

Ceph Crimson 网络调用和消息派发到OSD Shard分析

Ceph Crimson 网络调用和消息派发到OSD Shard分析

概述

本文档分析Ceph Crimson中网络层如何被调用,以及消息如何从网络层派发到OSD的shard的完整流程。

1. 网络层架构

1.1 核心组件

  • SocketMessenger: 消息传递器,管理连接和消息路由
  • SocketConnection: 表示一个网络连接
  • IOHandler: 处理连接的I/O操作(读/写消息)
  • ProtocolV2: 实现Ceph消息协议V2版本
  • FrameAssemblerV2: 负责消息帧的组装和解析
  • ChainedDispatchers: 链式分发器,将消息分发给多个Dispatcher

1.2 Shard管理

Crimson使用Seastar的shard模型,每个shard是一个独立的执行上下文(通常对应一个CPU核心)。网络连接可以在不同的shard之间迁移。

2. 网络监听和连接建立

2.1 启动监听

1
2
3
SocketMessenger::start(
const dispatchers_t& _dispatchers) {
assert(seastar::this_shard_id() == sid);

流程:

  1. SocketMessenger::start() 被调用,传入dispatchers列表
  2. 如果已绑定地址,调用 ShardedServerSocket::accept() 开始接受连接
  3. ShardedServerSocket 在所有shard上创建监听器(如果dispatch_only_on_this_shard=false

2.2 接受连接

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
67
68
69
70
71
72
seastar::future<>
ShardedServerSocket::accept(accept_func_t &&_fn_accept)
{
ceph_assert_always(seastar::this_shard_id() == primary_sid);
logger().debug("ShardedServerSocket({})::accept()...", listen_addr);
return this->container().invoke_on_all([_fn_accept](auto &ss) {
assert(ss.listener);
ss.fn_accept = _fn_accept;
// gate accepting
// ShardedServerSocket::shutdown() will drain the continuations in the gate
// so ignore the returned future
std::ignore = seastar::with_gate(ss.shutdown_gate, [&ss] {
return seastar::keep_doing([&ss] {
return ss.listener->accept(
).then([&ss](seastar::accept_result accept_result) {
#ifndef NDEBUG
if (ss.dispatch_only_on_primary_sid) {
// see seastar::listen_options::set_fixed_cpu()
ceph_assert_always(seastar::this_shard_id() == ss.primary_sid);
}
#endif
auto [socket, paddr] = std::move(accept_result);
entity_addr_t peer_addr;
peer_addr.set_sockaddr(&paddr.as_posix_sockaddr());
peer_addr.set_type(ss.listen_addr.get_type());
SocketRef _socket = std::make_unique<Socket>(
std::move(socket), Socket::side_t::acceptor,
peer_addr.get_port(), Socket::construct_tag{});
logger().debug("ShardedServerSocket({})::accept(): accepted peer {}, "
"socket {}, dispatch_only_on_primary_sid = {}",
ss.listen_addr, peer_addr, fmt::ptr(_socket.get()),
ss.dispatch_only_on_primary_sid);
std::ignore = seastar::with_gate(
ss.shutdown_gate,
[socket=std::move(_socket), peer_addr, &ss]() mutable {
return ss.fn_accept(std::move(socket), peer_addr
).handle_exception([&ss, peer_addr](auto eptr) {
const char *e_what;
try {
std::rethrow_exception(eptr);
} catch (std::exception &e) {
e_what = e.what();
}
logger().error("ShardedServerSocket({})::accept(): "
"fn_accept(s, {}) got unexpected exception {}",
ss.listen_addr, peer_addr, e_what);
ceph_abort();
});
});
});
}).handle_exception_type([&ss](const std::system_error& e) {
if (e.code() == std::errc::connection_aborted ||
e.code() == std::errc::invalid_argument) {
logger().debug("ShardedServerSocket({})::accept(): stopped ({})",
ss.listen_addr, e.what());
} else {
throw;
}
}).handle_exception([&ss](auto eptr) {
const char *e_what;
try {
std::rethrow_exception(eptr);
} catch (std::exception &e) {
e_what = e.what();
}
logger().error("ShardedServerSocket({})::accept(): "
"got unexpected exception {}", ss.listen_addr, e_what);
ceph_abort();
});
});
});
}

当新连接到达时:

  1. ShardedServerSocket::accept() 在循环中接受连接
  2. 创建 Socket 对象
  3. 调用 fn_accept 回调(由 SocketMessenger::accept() 提供)
  4. SocketMessenger::accept() 创建 SocketConnection 并开始握手

3. 消息读取流程

3.1 启动消息读取

当连接建立并完成握手后,IOHandler 开始读取消息:

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
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
void IOHandler::do_in_dispatch()
{
shard_states->enter_in_dispatching();
shard_states->dispatch_in_background(
"do_in_dispatch", conn, [this, &ctx=*shard_states] {
return seastar::keep_doing([this, &ctx] {
return frame_assembler->read_main_preamble<false>(
).then([this, &ctx](auto ret) {
switch (ret.tag) {
case Tag::MESSAGE: {
size_t msg_size = get_msg_size(*ret.rx_frame_asm);
return seastar::futurize_invoke([this] {
// throttle_message() logic
if (!conn.policy.throttler_messages) {
return seastar::now();
}
// TODO: message throttler
ceph_abort_msg("TODO");
return seastar::now();
}).then([this, msg_size] {
// throttle_bytes() logic
if (!conn.policy.throttler_bytes) {
return seastar::now();
}
if (!msg_size) {
return seastar::now();
}
logger().trace("{} wants {} bytes from policy throttler {}/{}",
conn, msg_size,
conn.policy.throttler_bytes->get_current(),
conn.policy.throttler_bytes->get_max());
return conn.policy.throttler_bytes->get(msg_size);
}).then([this, msg_size, &ctx] {
// TODO: throttle_dispatch_queue() logic
utime_t throttle_stamp{seastar::lowres_system_clock::now()};
return read_message(ctx, throttle_stamp, msg_size);
});
}
case Tag::ACK:
return frame_assembler->read_frame_payload<false>(
).then([this](auto payload) {
// handle_message_ack() logic
auto ack = AckFrame::Decode(payload->back());
logger().debug("{} GOT AckFrame: seq={}", conn, ack.seq());
ack_out_sent(ack.seq());
});
case Tag::KEEPALIVE2:
return frame_assembler->read_frame_payload<false>(
).then([this](auto payload) {
// handle_keepalive2() logic
auto keepalive_frame = KeepAliveFrame::Decode(payload->back());
logger().debug("{} GOT KeepAliveFrame: timestamp={}",
conn, keepalive_frame.timestamp());
// notify keepalive ack
next_keepalive_ack = keepalive_frame.timestamp();
if (seastar::this_shard_id() == get_shard_id()) {
notify_out_dispatch();
}

last_keepalive = seastar::lowres_system_clock::now();
});
case Tag::KEEPALIVE2_ACK:
return frame_assembler->read_frame_payload<false>(
).then([this](auto payload) {
// handle_keepalive2_ack() logic
auto keepalive_ack_frame = KeepAliveFrameAck::Decode(payload->back());
auto _last_keepalive_ack =
seastar::lowres_system_clock::time_point{keepalive_ack_frame.timestamp()};
set_last_keepalive_ack(_last_keepalive_ack);
logger().debug("{} GOT KeepAliveFrameAck: timestamp={}",
conn, _last_keepalive_ack);
});
default: {
logger().warn("{} do_in_dispatch() received unexpected tag: {}",
conn, static_cast<uint32_t>(ret.tag));
abort_in_fault();
}
}
});
}).handle_exception([this, &ctx](std::exception_ptr eptr) {
const char *e_what;
try {
std::rethrow_exception(eptr);
} catch (std::exception &e) {
e_what = e.what();
}

auto io_state = ctx.get_io_state();
if (io_state == io_state_t::open) {
auto cc_seq = proto_crosscore.prepare_submit();
logger().info("{} do_in_dispatch(): fault at {}, {}, going to delay -- {}, "
"send {} notify_out_fault()",
conn, io_state, io_stat_printer{*this}, e_what, cc_seq);
do_set_io_state(io_state_t::delay);
shard_states->dispatch_in_background(
"notify_out_fault(in)", conn, [this, cc_seq, eptr] {
auto states = get_states();
return seastar::smp::submit_to(
conn.get_messenger_shard_id(), [this, cc_seq, eptr, states] {
return handshake_listener->notify_out_fault(
cc_seq, "do_in_dispatch", eptr, states);
});
});
} else {
if (io_state != io_state_t::switched) {
logger().info("{} do_in_dispatch(): fault at {}, {} -- {}",
conn, io_state, io_stat_printer{*this}, e_what);
} else {
logger().info("{} do_in_dispatch(): fault at {} -- {}",
conn, io_state, e_what);
}
}
}).finally([&ctx] {
ctx.exit_in_dispatching();
});
});
}

流程:

  1. do_in_dispatch() 进入消息读取循环
  2. 调用 frame_assembler->read_main_preamble() 读取消息头
  3. 根据消息类型(MESSAGE、ACK、KEEPALIVE等)进行不同处理
  4. 对于MESSAGE类型,调用 read_message() 读取完整消息

3.2 读取和解析消息

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
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
seastar::future<>
IOHandler::read_message(
shard_states_t &ctx,
utime_t throttle_stamp,
std::size_t msg_size)
{
return frame_assembler->read_frame_payload<false>(
).then([this, throttle_stamp, msg_size, &ctx](auto payload) {
if (unlikely(ctx.get_io_state() != io_state_t::open)) {
logger().debug("{} triggered {} during read_message()",
conn, ctx.get_io_state());
abort_protocol();
}

utime_t recv_stamp{seastar::lowres_system_clock::now()};

// we need to get the size before std::moving segments data
auto msg_frame = MessageFrame::Decode(*payload);
// XXX: paranoid copy just to avoid oops
ceph_msg_header2 current_header = msg_frame.header();

logger().trace("{} got {} + {} + {} byte message,"
" envelope type={} src={} off={} seq={}",
conn,
msg_frame.front_len(),
msg_frame.middle_len(),
msg_frame.data_len(),
(uint16_t)current_header.type,
conn.get_peer_name(),
(uint16_t)current_header.data_off,
(uint32_t)current_header.seq);

ceph_msg_header header{current_header.seq,
current_header.tid,
current_header.type,
current_header.priority,
current_header.version,
ceph_le32(msg_frame.front_len()),
ceph_le32(msg_frame.middle_len()),
ceph_le32(msg_frame.data_len()),
current_header.data_off,
conn.get_peer_name(),
current_header.compat_version,
current_header.reserved,
ceph_le32(0)};
ceph_msg_footer footer{ceph_le32(0), ceph_le32(0),
ceph_le32(0), ceph_le64(0), current_header.flags};

Message *message = decode_message(nullptr, 0, header, footer,
msg_frame.front(), msg_frame.middle(), msg_frame.data(), nullptr);
if (!message) {
logger().warn("{} decode message failed", conn);
abort_in_fault();
}

// store reservation size in message, so we don't get confused
// by messages entering the dispatch queue through other paths.
message->set_dispatch_throttle_size(msg_size);

message->set_throttle_stamp(throttle_stamp);
message->set_recv_stamp(recv_stamp);
message->set_recv_complete_stamp(utime_t{seastar::lowres_system_clock::now()});

// check received seq#. if it is old, drop the message.
// note that incoming messages may skip ahead. this is convenient for the
// client side queueing because messages can't be renumbered, but the (kernel)
// client will occasionally pull a message out of the sent queue to send
// elsewhere. in that case it doesn't matter if we "got" it or not.
uint64_t cur_seq = in_seq;
if (message->get_seq() <= cur_seq) {
logger().error("{} got old message {} <= {} {}, discarding",
conn, message->get_seq(), cur_seq, *message);
if (HAVE_FEATURE(conn.features, RECONNECT_SEQ) &&
local_conf()->ms_die_on_old_message) {
ceph_assert(0 == "old msgs despite reconnect_seq feature");
}
return seastar::now();
} else if (message->get_seq() > cur_seq + 1) {
logger().error("{} missed message? skipped from seq {} to {}",
conn, cur_seq, message->get_seq());
if (local_conf()->ms_die_on_skipped_message) {
ceph_assert(0 == "skipped incoming seq");
}
}

// note last received message.
in_seq = message->get_seq();
if (conn.policy.lossy) {
logger().debug("{} <== #{} === {} ({})",
conn,
message->get_seq(),
*message,
message->get_type());
} else {
logger().debug("{} <== #{},{} === {} ({})",
conn,
message->get_seq(),
(uint32_t)current_header.ack_seq,
*message,
*message,
message->get_type());
}

// notify ack
if (!conn.policy.lossy) {
++ack_left;
notify_out_dispatch();
}

ack_out_sent(current_header.ack_seq);

// TODO: change MessageRef with seastar::shared_ptr
auto msg_ref = MessageRef{message, false};
assert(ctx.get_io_state() == io_state_t::open);
assert(get_io_state() == io_state_t::open);
ceph_assert_always(conn_ref);

// throttle the reading process by the returned future
return dispatchers.ms_dispatch(conn_ref, std::move(msg_ref));
// user can make changes
});
}

关键步骤:

  1. 读取消息帧的payload
  2. 解码消息头(MessageFrame::Decode
  3. 调用 decode_message() 创建 Message 对象
  4. 验证消息序列号
  5. 调用 dispatchers.ms_dispatch() 派发消息

4. 消息派发到Dispatcher

4.1 ChainedDispatchers派发

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
seastar::future<>
ChainedDispatchers::ms_dispatch(ConnectionRef conn,
MessageRef m) {
try {
for (auto& dispatcher : dispatchers) {
auto dispatched = dispatcher->ms_dispatch(conn, m);
if (dispatched.has_value()) {
return std::move(*dispatched
).handle_exception([conn] (std::exception_ptr eptr) {
logger().error("{} got unexpected exception in ms_dispatch() throttling {}",
*conn, eptr);
ceph_abort();
});
}
}
} catch (...) {
logger().error("{} got unexpected exception in ms_dispatch() {}",
*conn, std::current_exception());
ceph_abort();
}
if (!dispatchers.empty()) {
logger().error("ms_dispatch unhandled message {}", *m);
}
return seastar::now();
}

ChainedDispatchers 按顺序遍历所有dispatchers:

  • 如果某个dispatcher返回非空的future,表示它处理了该消息,停止遍历
  • 如果所有dispatchers都不处理,记录错误

4.2 OSD的Dispatcher实现

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
std::optional<seastar::future<>>
OSD::ms_dispatch(crimson::net::ConnectionRef conn, MessageRef m)
{
if (pg_shard_manager.is_stopping()) {
return seastar::now();
}
auto maybe_ret = do_ms_dispatch(conn, std::move(m));
if (!maybe_ret.has_value()) {
return std::nullopt;
}

gate.dispatch_in_background(
__func__, *this, [ret=std::move(maybe_ret.value())]() mutable {
return std::move(ret);
});
return seastar::now();
}

OSD的 ms_dispatch() 实现:

  1. 检查OSD是否正在停止
  2. 调用 do_ms_dispatch() 处理消息
  3. 如果返回future,在后台执行

4.3 跨Shard消息处理

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
std::optional<seastar::future<>>
OSD::do_ms_dispatch(
crimson::net::ConnectionRef conn,
MessageRef m)
{
if (seastar::this_shard_id() != PRIMARY_CORE) {
switch (m->get_type()) {
case CEPH_MSG_OSD_MAP:
case MSG_COMMAND:
case MSG_OSD_MARK_ME_DOWN:
// FIXME: order is not guaranteed in this path
return conn.get_foreign(
).then([this, m=std::move(m)](auto f_conn) {
return seastar::smp::submit_to(PRIMARY_CORE,
[f_conn=std::move(f_conn), m=std::move(m), this]() mutable {
auto conn = make_local_shared_foreign(std::move(f_conn));
auto ret = do_ms_dispatch(conn, std::move(m));
assert(ret.has_value());
return std::move(ret.value());
});
});
}
}

switch (m->get_type()) {

重要机制:

  • 如果消息不在PRIMARY_CORE上,某些消息类型(如OSD_MAP)会被转发到PRIMARY_CORE
  • 使用 seastar::smp::submit_to() 跨shard提交任务

5. Shard迁移机制

5.1 连接迁移到新Shard

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
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
seastar::future<>
IOHandler::to_new_sid(
cc_seq_t cc_seq,
seastar::shard_id new_sid,
ConnectionFRef conn_fref,
std::optional<bool> is_replace)
{
ceph_assert_always(seastar::this_shard_id() == get_shard_id());
if (!proto_crosscore.proceed_or_wait(cc_seq)) {
logger().debug("{} got {} to_new_sid(), wait at {}",
conn, cc_seq, proto_crosscore.get_in_seq());
return proto_crosscore.wait(cc_seq
).then([this, cc_seq, new_sid, is_replace,
conn_fref=std::move(conn_fref)]() mutable {
return to_new_sid(cc_seq, new_sid, std::move(conn_fref), is_replace);
});
}

bool is_accept_or_connect = is_replace.has_value();
logger().debug("{} got {} to_new_sid_1(new_sid={}, {}) at {}",
conn, cc_seq, new_sid,
fmt::format("{}",
is_accept_or_connect ?
(*is_replace ? "accept(replace)" : "accept(!replace)") :
"connect"),
io_stat_printer{*this});
auto next_cc_seq = ++cc_seq;

if (get_io_state() != io_state_t::drop) {
ceph_assert_always(conn_ref);
if (new_sid != seastar::this_shard_id()) {
dispatchers.ms_handle_shard_change(conn_ref, new_sid, is_accept_or_connect);
// user can make changes
}
} else {
// it is possible that both io_handler and protocolv2 are
// trying to close each other from different cores simultaneously.
assert(!protocol_is_connected);
}

if (get_io_state() != io_state_t::drop) {
if (is_accept_or_connect) {
// protocol_is_connected can be from true to true here if the replacing is
// happening to a connected connection.
} else {
ceph_assert_always(protocol_is_connected == false);
}
protocol_is_connected = true;
} else {
assert(!protocol_is_connected);
}

bool is_dropped = false;
if (get_io_state() == io_state_t::drop) {
is_dropped = true;
}
ceph_assert_always(get_io_state() != io_state_t::open);

// apply the switching atomically
ceph_assert_always(conn_ref);
conn_ref.reset();
auto prv_sid = get_shard_id();
ceph_assert_always(maybe_prv_shard_states == nullptr);
maybe_prv_shard_states = std::move(shard_states);
shard_states = shard_states_t::create_from_previous(
*maybe_prv_shard_states, new_sid);
assert(new_sid == get_shard_id());
// broadcast shard change to all the io waiters, atomically.
io_crosscore.reset_wait();

return seastar::smp::submit_to(new_sid,
[this, next_cc_seq, is_dropped, prv_sid, is_replace, conn_fref=std::move(conn_fref)]() mutable {
logger().debug("{} got {} to_new_sid_2(prv_sid={}, is_dropped={}, {}) at {}",
conn, next_cc_seq, prv_sid, is_dropped,
fmt::format("{}",
is_replace.has_value() ?
(*is_replace ? "accept(replace)" : "accept(!replace)") :
"connect"),
io_stat_printer{*this});

ceph_assert_always(seastar::this_shard_id() == get_shard_id());
ceph_assert_always(get_io_state() != io_state_t::open);
ceph_assert_always(!maybe_dropped_sid.has_value());
ceph_assert_always(proto_crosscore.proceed_or_wait(next_cc_seq));

if (is_dropped) {
ceph_assert_always(get_io_state() == io_state_t::drop);
ceph_assert_always(shard_states->assert_closed_and_exit());
maybe_dropped_sid = prv_sid;
// cleanup_prv_shard() will be done in a follow-up close_io()
} else {
// possible at io_state_t::drop

// previous shard is not cleaned,
// but close_io() is responsible to clean up the current shard,
// so cleanup the previous shard here.
shard_states->dispatch_in_background(
"cleanup_prv_sid", conn, [this, prv_sid] {
return cleanup_prv_shard(prv_sid);
});
maybe_notify_out_dispatch();
}

ceph_assert_always(!conn_ref);
// assign even if already dropping
conn_ref = make_local_shared_foreign(std::move(conn_fref));

if (get_io_state() != io_state_t::drop) {
if (is_replace.has_value()) {
dispatchers.ms_handle_accept(conn_ref, prv_sid, *is_replace);
} else {
dispatchers.ms_handle_connect(conn_ref, prv_sid);
}
// user can make changes
}
});
}

Shard迁移流程:

  1. 在旧shard上:保存当前状态,通知dispatchers shard变化
  2. 创建新shard的状态对象
  3. 提交任务到新shard
  4. 在新shard上:恢复连接引用,通知dispatchers连接已迁移

6. 消息发送流程

6.1 发送消息

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
seastar::future<> IOHandler::send(MessageURef _msg)
{
// may be invoked from any core
MessageFRef msg = seastar::make_foreign(std::move(_msg));
auto cc_seq = io_crosscore.prepare_submit();
auto source_core = seastar::this_shard_id();
logger().debug("{} send {} send() to core {} -- {}",
conn, cc_seq, get_shard_id(), *msg);
// sid may be changed on-the-fly during the submission
if (source_core == get_shard_id()) {
return do_send(cc_seq, source_core, std::move(msg));
} else {
return seastar::smp::submit_to(
get_shard_id(),
[this, cc_seq, source_core, msg=std::move(msg)]() mutable {
return send_recheck_shard(cc_seq, source_core, std::move(msg));
});
}
}

发送流程:

  1. 如果调用者不在连接的shard上,使用 smp::submit_to() 提交到正确的shard
  2. 调用 do_send() 将消息加入发送队列
  3. 通知输出调度器开始发送

7. 关键数据结构

7.1 Shard States

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
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
class shard_states_t {
public:
shard_states_t(seastar::shard_id _sid, io_state_t state)
: sid{_sid}, io_state{state}, gate{_sid} {}

seastar::shard_id get_shard_id() const {
return sid;
}

io_state_t get_io_state() const {
assert(seastar::this_shard_id() == sid);
return io_state;
}

void set_io_state(io_state_t new_state) {
assert(seastar::this_shard_id() == sid);
assert(io_state != new_state);
pr_io_state_changed.set_value();
pr_io_state_changed = seastar::promise<>();
if (io_state == io_state_t::open) {
// from open
if (out_dispatching) {
ceph_assert_always(!out_exit_dispatching.has_value());
out_exit_dispatching = seastar::promise<>();
}
}
io_state = new_state;
}

seastar::future<> wait_state_change() {
assert(seastar::this_shard_id() == sid);
return pr_io_state_changed.get_future();
}

template <typename Func>
void dispatch_in_background(
const char *what, SocketConnection &who, Func &&func) {
assert(seastar::this_shard_id() == sid);
ceph_assert_always(!gate.is_closed());
gate.dispatch_in_background(what, who, std::move(func));
}

void enter_in_dispatching() {
assert(seastar::this_shard_id() == sid);
assert(io_state == io_state_t::open);
ceph_assert_always(!in_exit_dispatching.has_value());
in_exit_dispatching = seastar::promise<>();
}

void exit_in_dispatching() {
assert(seastar::this_shard_id() == sid);
assert(io_state != io_state_t::open);
ceph_assert_always(in_exit_dispatching.has_value());
in_exit_dispatching->set_value();
in_exit_dispatching = std::nullopt;
}

bool try_enter_out_dispatching(SocketConnection &conn) {
assert(seastar::this_shard_id() == sid);
if (out_dispatching) {
// already dispatching out
return false;
}
switch (io_state) {
case io_state_t::open:
[[fallthrough]];
case io_state_t::delay:
out_dispatching = true;
return true;
case io_state_t::drop:
[[fallthrough]];
case io_state_t::switched:
// do not dispatch out
return false;
default:
crimson::get_logger(ceph_subsys_ms).error(
"{} try_enter_out_dispatching() got wrong io_state {}",
conn, io_state);
ceph_abort_msg("impossible");
}
}

void notify_out_dispatching_stopped(
const char *what, SocketConnection &conn);

void exit_out_dispatching(
const char *what, SocketConnection &conn) {
assert(seastar::this_shard_id() == sid);
ceph_assert_always(out_dispatching);
out_dispatching = false;
notify_out_dispatching_stopped(what, conn);
}

seastar::future<> wait_io_exit_dispatching();

seastar::future<> close() {
assert(seastar::this_shard_id() == sid);
assert(!gate.is_closed());
return gate.close();
}

bool assert_closed_and_exit() const {
assert(seastar::this_shard_id() == sid);
if (gate.is_closed()) {
ceph_assert_always(io_state == io_state_t::drop ||
io_state == io_state_t::switched);
ceph_assert_always(!out_dispatching);
ceph_assert_always(!out_exit_dispatching);
ceph_assert_always(!in_exit_dispatching);
return true;
} else {
return false;
}
}

static shard_states_ref_t create(
seastar::shard_id sid, io_state_t state) {
return std::make_unique<shard_states_t>(sid, state);
}

static shard_states_ref_t create_from_previous(
shard_states_t &prv_states, seastar::shard_id new_sid);

private:
const seastar::shard_id sid;
io_state_t io_state;

crimson::common::Gated gate;
seastar::promise<> pr_io_state_changed;
bool out_dispatching = false;
std::optional<seastar::promise<>> out_exit_dispatching;
std::optional<seastar::promise<>> in_exit_dispatching;
};

shard_states_t 管理每个shard的状态:

  • io_state: delay/open/drop/switched
  • gate: 用于同步后台任务
  • out_dispatching/in_dispatching: 跟踪I/O调度状态

8. 总结

8.1 消息接收流程

  1. 网络层接受连接ShardedServerSocket::accept()
  2. 创建连接对象SocketMessenger::accept()SocketConnection
  3. 握手完成ProtocolV2 完成握手
  4. 启动消息读取IOHandler::do_in_dispatch()
  5. 读取消息帧FrameAssemblerV2::read_main_preamble()
  6. 解析消息IOHandler::read_message()
  7. 派发给DispatcherChainedDispatchers::ms_dispatch()
  8. OSD处理消息OSD::ms_dispatch()OSD::do_ms_dispatch()

8.2 Shard分配机制

  • 连接可以在不同shard之间迁移
  • 使用 seastar::smp::submit_to() 跨shard提交任务
  • 每个shard维护独立的状态(shard_states_t
  • OSD的某些消息必须在PRIMARY_CORE上处理

8.3 关键设计点

  1. 异步I/O: 所有网络操作都是异步的,使用Seastar的future/promise模型
  2. Shard隔离: 每个shard独立处理,避免锁竞争
  3. 状态管理: 使用状态机管理连接和I/O状态
  4. 错误处理: 完善的异常处理和连接恢复机制

CRUSH算法实现原理详细分析

CRUSH算法实现原理详细分析

目录

  1. CRUSH算法概述
  2. 核心数据结构
  3. 哈希函数
  4. Bucket算法
  5. 映射规则(Rule)
  6. 映射执行流程
  7. 代码实现分析
  8. 性能优化

CRUSH算法概述

1.1 什么是CRUSH

CRUSH (Controlled Replication Under Scalable Hashing) 是Ceph中用于数据分布的核心算法。它是一个伪随机数据分布算法,能够高效地将输入值(通常是数据对象)分布到异构的、结构化的存储集群中。

1.2 核心特点

  • 确定性:给定相同的输入,总是产生相同的输出
  • 伪随机性:分布看起来是随机的,但实际上是确定性的
  • 可扩展性:支持大规模集群
  • 容错性:能够处理节点故障和恢复
  • 权重支持:根据设备权重进行负载均衡

1.3 算法原理

CRUSH算法通过以下步骤将对象映射到OSD:

1
2
3
4
5
6
7
8
9
对象名 (object_name) 

PG ID (pgid = hash(object_name) % num_pgs)

CRUSH输入 (x = hash(pgid, pool_id))

CRUSH规则 (rule)

OSD列表 (osd_list)

1.4 关键概念

  • Item:CRUSH层次结构中的节点,可以是设备(OSD)或桶(Bucket)
  • Bucket:包含其他Item的容器,形成层次结构
  • Rule:定义如何从层次结构中选择Item的规则序列
  • Weight:Item的权重,用于负载均衡
  • Type:Item的类型,用于故障域隔离

核心数据结构

2.1 crush_map

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
struct crush_map {
struct crush_bucket **buckets; // 桶数组
struct crush_rule **rules; // 规则数组
__s32 max_buckets; // 最大桶数量
__u32 max_rules; // 最大规则数量
__s32 max_devices; // 最大设备数量

// 可调参数
__u32 choose_total_tries; // 总重试次数
__u32 choose_local_tries; // 本地重试次数
__u32 choose_local_fallback_tries; // 本地回退重试次数
__u32 chooseleaf_descend_once; // chooseleaf是否只下降一次
__u8 chooseleaf_vary_r; // chooseleaf是否变化r值
__u8 chooseleaf_stable; // chooseleaf稳定模式

size_t working_size; // 工作空间大小
};

说明

  • buckets:存储所有桶的指针数组,桶ID为负数(-1, -2, …)
  • rules:存储所有规则的指针数组
  • max_devices:最大设备ID + 1
  • 可调参数用于控制映射行为和稳定性

2.2 crush_bucket

1
2
3
4
5
6
7
8
9
struct crush_bucket {
__s32 id; // 桶ID,< 0且唯一
__u16 type; // 桶类型,> 0,由调用者定义
__u8 alg; // 项目选择算法
__u8 hash; // 哈希函数类型
__u32 weight; // 16.16定点累积子项权重
__u32 size; // items数组大小
__s32 *items; // 子项数组:< 0是桶,>= 0是设备
};

说明

  • id:桶的唯一标识符,必须为负数
  • type:桶的类型,用于规则匹配(如rack、host等)
  • alg:选择算法(uniform、list、straw2等)
  • weight:所有子项的累积权重
  • items:子项数组,负数表示桶,非负数表示设备

2.3 桶类型结构

2.3.1 Uniform Bucket

1
2
3
4
struct crush_bucket_uniform {
struct crush_bucket h;
__u32 item_weight; // 每个项目的权重(所有项目相同)
};

特点

  • 所有项目权重相同
  • 选择速度最快 O(1)
  • 添加/删除项目时数据移动较大

2.3.2 List Bucket

1
2
3
4
5
struct crush_bucket_list {
struct crush_bucket h;
__u32 *item_weights; // 每个项目的权重
__u32 *sum_weights; // 累积权重
};

特点

  • 支持不同权重
  • 选择速度 O(n)
  • 添加项目时数据移动最优
  • 删除项目时数据移动较大

2.3.3 Straw2 Bucket

1
2
3
4
struct crush_bucket_straw2 {
struct crush_bucket h;
__u32 *item_weights; // 每个项目的权重
};

特点

  • 支持不同权重
  • 选择速度 O(n)
  • 添加/删除/重权重时数据移动最优
  • 推荐使用的算法

2.4 crush_rule

1
2
3
4
5
6
7
8
9
10
11
struct crush_rule {
__u32 len; // 步骤数量
__u8 type; // 规则类型
struct crush_rule_step steps[0]; // 步骤数组
};

struct crush_rule_step {
__u32 op; // 操作码
__s32 arg1; // 参数1
__s32 arg2; // 参数2
};

操作码类型

  • CRUSH_RULE_TAKE:选择起始桶
  • CRUSH_RULE_CHOOSE_FIRSTN:选择N个项目(深度优先)
  • CRUSH_RULE_CHOOSE_INDEP:选择N个项目(广度优先)
  • CRUSH_RULE_CHOOSELEAF_FIRSTN:选择N个叶子(深度优先)
  • CRUSH_RULE_CHOOSELEAF_INDEP:选择N个叶子(广度优先)
  • CRUSH_RULE_EMIT:输出结果

哈希函数

3.1 哈希函数实现

CRUSH使用Robert Jenkins的哈希函数(rjenkins1),位于hash.c

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
// 哈希混合函数
#define crush_hashmix(a, b, c) do { \
a = a-b; a = a-c; a = a^(c>>13); \
b = b-c; b = b-a; b = b^(a<<8); \
c = c-a; c = c-b; c = c^(b>>13); \
a = a-b; a = a-c; a = a^(c>>12); \
b = b-c; b = b-a; b = b^(a<<16); \
c = c-a; c = c-b; c = c^(b>>5); \
a = a-b; a = a-c; a = a^(c>>3); \
b = b-c; b = b-a; b = b^(a<<10); \
c = c-a; c = c-b; c = c^(b>>15); \
} while (0)

// 单参数哈希
__u32 crush_hash32_rjenkins1(__u32 a) {
__u32 hash = crush_hash_seed ^ a;
__u32 b = a;
__u32 x = 231232;
__u32 y = 1232;
crush_hashmix(b, x, hash);
crush_hashmix(y, a, hash);
return hash;
}

// 多参数哈希
__u32 crush_hash32_3(int type, __u32 a, __u32 b, __u32 c);
__u32 crush_hash32_4(int type, __u32 a, __u32 b, __u32 c, __u32 d);

特点

  • 使用混合函数确保良好的分布
  • 支持多个输入参数
  • 确定性:相同输入产生相同输出

3.2 哈希函数使用

在CRUSH映射中,哈希函数用于:

  1. 桶内项目选择hash(x, bucket_id, r) 选择桶内项目
  2. 排列生成hash(x, bucket_id, position) 生成排列
  3. 权重检查hash(x, item) 检查项目是否”out”

Bucket算法

4.1 Uniform Bucket算法

实现位置mapper.c:115-119

1
2
3
4
5
6
7
static int bucket_uniform_choose(
const struct crush_bucket_uniform *bucket,
struct crush_work_bucket *work,
int x, int r)
{
return bucket_perm_choose(&bucket->h, work, x, r);
}

算法原理

  1. 所有项目权重相同
  2. 使用排列选择算法
  3. 时间复杂度:O(1)(优化后)

排列选择算法bucket_perm_choose):

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
static int bucket_perm_choose(const struct crush_bucket *bucket,
struct crush_work_bucket *work,
int x, int r)
{
unsigned int pr = r % bucket->size;

// 优化:r=0的情况
if (pr == 0) {
s = crush_hash32_3(bucket->hash, x, bucket->id, 0) % bucket->size;
work->perm[0] = s;
return bucket->items[s];
}

// 生成排列
for (i = 0; i < bucket->size; i++)
work->perm[i] = i;

// 计算排列到位置pr
for (p = 0; p <= pr; p++) {
i = crush_hash32_3(bucket->hash, x, bucket->id, p) % (bucket->size - p);
if (i) {
// 交换
swap(work->perm[p], work->perm[p + i]);
}
}

return bucket->items[work->perm[pr]];
}

特点

  • 快速:O(1)平均情况
  • 均匀分布
  • 添加/删除项目时数据移动大

4.2 List Bucket算法

实现位置mapper.c:122-145

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
static int bucket_list_choose(const struct crush_bucket_list *bucket,
int x, int r)
{
for (i = bucket->h.size-1; i >= 0; i--) {
// 计算哈希值
w = crush_hash32_4(bucket->h.hash, x, bucket->h.items[i], r, bucket->h.id);
w &= 0xffff;

// 缩放权重
w *= bucket->sum_weights[i];
w = w >> 16;

// 检查是否选择
if (w < bucket->item_weights[i]) {
return bucket->h.items[i];
}
}
return bucket->h.items[0];
}

算法原理

  1. 从列表尾部(最新添加的项目)开始
  2. 对每个项目计算哈希值
  3. 根据累积权重决定是否选择该项目
  4. 如果选择,返回该项目;否则继续

特点

  • 时间复杂度:O(n)
  • 添加项目时数据移动最优
  • 删除项目时数据移动较大

4.3 Straw2 Bucket算法

实现位置mapper.c:342-365

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
static int bucket_straw2_choose(
const struct crush_bucket_straw2 *bucket,
int x, int r,
const struct crush_choose_arg *arg,
int position)
{
unsigned int i, high = 0;
__s64 draw, high_draw = 0;
__u32 *weights = get_choose_arg_weights(bucket, arg, position);
__s32 *ids = get_choose_arg_ids(bucket, arg);

for (i = 0; i < bucket->h.size; i++) {
if (weights[i]) {
// 生成指数分布随机变量
draw = generate_exponential_distribution(
bucket->h.hash, x, ids[i], r, weights[i]);
} else {
draw = S64_MIN; // 权重为0,不选择
}

// 选择最大的draw值
if (i == 0 || draw > high_draw) {
high = i;
high_draw = draw;
}
}

return bucket->h.items[high];
}

指数分布生成generate_exponential_distribution):

1
2
3
4
5
6
7
8
9
10
11
12
13
static inline __s64 generate_exponential_distribution(
int type, int x, int y, int z, int weight)
{
// 生成随机数
unsigned int u = crush_hash32_3(type, x, y, z);
u &= 0xffff;

// 计算自然对数(使用查找表)
__s64 ln = crush_ln(u) - 0x1000000000000ll;

// 除以权重(16.16定点)
return div64_s64(ln, weight);
}

算法原理

  1. 为每个项目生成一个指数分布的随机变量(”straw”长度)
  2. 权重越大,straw长度期望越大
  3. 选择straw长度最大的项目

数学原理

  • 使用指数分布的逆变换采样
  • 如果每个OSD的请求间隔服从指数分布,则PG分布与权重成正比
  • 参考:指数分布最小值的分布

特点

  • 时间复杂度:O(n)
  • 添加/删除/重权重时数据移动最优
  • 推荐使用

映射规则(Rule)

5.1 规则结构

规则由一系列步骤组成,每个步骤执行一个操作:

1
2
3
4
5
6
7
8
// 示例规则:三副本规则
rule replicated_rule {
id 0
type replicated
step take root // 从root桶开始
step chooseleaf firstn 0 type host // 从每个host选择1个OSD
step emit // 输出结果
}

5.2 规则步骤类型

5.2.1 TAKE

1
2
3
4
case CRUSH_RULE_TAKE:
w[0] = curstep->arg1; // 选择起始桶
wsize = 1;
break;

功能:选择规则执行的起始桶

5.2.2 CHOOSE_FIRSTN / CHOOSELEAF_FIRSTN

深度优先选择

  • CHOOSE_FIRSTN:选择N个桶
  • CHOOSELEAF_FIRSTN:选择N个叶子(OSD)

特点

  • 深度优先:先选择一个桶,再递归选择
  • 副本间有依赖关系
  • 适合副本存储

5.2.3 CHOOSE_INDEP / CHOOSELEAF_INDEP

广度优先选择

  • CHOOSE_INDEP:选择N个桶
  • CHOOSELEAF_INDEP:选择N个叶子(OSD)

特点

  • 广度优先:同时选择所有副本
  • 副本间独立
  • 适合纠删码

5.3 规则执行流程

规则执行在crush_do_rule_no_retry中实现:

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
static int crush_do_rule_no_retry(...)
{
// 初始化工作空间
w[0] = ...; // 起始桶
wsize = 1;

// 遍历规则步骤
for (step = 0; step < rule->len; step++) {
switch (curstep->op) {
case CRUSH_RULE_TAKE:
// 选择起始桶
break;

case CRUSH_RULE_CHOOSELEAF_FIRSTN:
case CRUSH_RULE_CHOOSE_FIRSTN:
// 深度优先选择
osize += crush_choose_firstn(...);
break;

case CRUSH_RULE_CHOOSELEAF_INDEP:
case CRUSH_RULE_CHOOSE_INDEP:
// 广度优先选择
crush_choose_indep(...);
break;

case CRUSH_RULE_EMIT:
// 输出结果
break;
}
}

return result_len;
}

映射执行流程

6.1 整体流程

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
crush_do_rule()

crush_do_rule_no_retry()

遍历规则步骤

TAKE: 选择起始桶

CHOOSE/CHOOSELEAF: 选择项目

├─► crush_choose_firstn() (深度优先)
│ ↓
│ 遍历每个副本位置
│ ↓
│ 选择桶内项目
│ ↓
│ crush_bucket_choose()
│ ↓
│ 根据桶类型选择算法
│ ├─► bucket_uniform_choose()
│ ├─► bucket_list_choose()
│ └─► bucket_straw2_choose()

└─► crush_choose_indep() (广度优先)

同时选择所有副本

crush_bucket_choose()

根据桶类型选择算法

6.2 crush_choose_firstn(深度优先)

实现位置mapper.c:441-629

核心逻辑

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
static int crush_choose_firstn(...)
{
// 对每个副本位置
for (rep = stable ? 0 : outpos; rep < numrep && count > 0; rep++) {
ftotal = 0;
do {
retry_descent = 0;
in = bucket; // 从起始桶开始

do {
retry_bucket = 0;
r = rep + parent_r + ftotal; // 计算r值

// 选择桶内项目
item = crush_bucket_choose(in, work, x, r, ...);

// 检查类型
if (itemtype != type) {
// 继续下降
in = map->buckets[-1-item];
retry_bucket = 1;
continue;
}

// 检查冲突
if (collide) {
// 重试
retry_bucket = 1;
}

// 如果是chooseleaf,递归选择叶子
if (recurse_to_leaf && item < 0) {
if (crush_choose_firstn(...) <= outpos) {
reject = 1;
}
}

// 检查是否out
if (is_out(...)) {
reject = 1;
}

} while (retry_bucket);
} while (retry_descent);

// 保存结果
out[outpos] = item;
outpos++;
}

return outpos;
}

关键点

  1. r值计算r = rep + parent_r + ftotal
    • rep:副本位置
    • parent_r:父级r值
    • ftotal:总失败次数
  2. 冲突检测:检查是否与已选择的项目冲突
  3. 重试机制
    • retry_bucket:桶内重试
    • retry_descent:下降重试
  4. out检测:检查项目是否可用(基于权重)

6.3 crush_choose_indep(广度优先)

实现位置mapper.c:636-824

核心逻辑

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
static void crush_choose_indep(...)
{
// 初始化所有位置为未定义
for (rep = outpos; rep < endpos; rep++) {
out[rep] = CRUSH_ITEM_UNDEF;
}

// 尝试选择
for (ftotal = 0; left > 0 && ftotal < tries; ftotal++) {
// 对每个未定义的位置
for (rep = outpos; rep < endpos; rep++) {
if (out[rep] != CRUSH_ITEM_UNDEF)
continue;

in = bucket;
for (;;) {
// 计算r值(独立于其他副本)
r = rep + parent_r;
if (in->alg == CRUSH_BUCKET_UNIFORM &&
in->size % numrep == 0) {
r += (numrep+1) * ftotal;
} else {
r += numrep * ftotal;
}

// 选择项目
item = crush_bucket_choose(in, work, x, r, ...);

// 检查类型、冲突、out等
// ...

// 保存结果
out[rep] = item;
left--;
break;
}
}
}
}

关键点

  1. 独立选择:每个副本位置独立选择
  2. r值计算:考虑副本数量和失败次数
  3. 广度优先:同时处理所有副本位置

6.4 is_out函数

实现位置mapper.c:405-419

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
static int is_out(const struct crush_map *map,
const __u32 *weight, int weight_max,
int item, int x)
{
if (item >= weight_max)
return 1; // 超出范围,out
if (weight[item] >= 0x10000)
return 0; // 权重>=1.0,in
if (weight[item] == 0)
return 1; // 权重=0,out

// 基于哈希的概率检查
if ((crush_hash32_2(CRUSH_HASH_RJENKINS1, x, item) & 0xffff)
< weight[item])
return 0; // 在概率范围内,in
return 1; // 超出概率范围,out
}

功能:检查项目是否”out”(不可用)

原理

  • 权重=0:总是out
  • 权重>=1.0:总是in
  • 0<权重<1.0:基于哈希的概率检查

代码实现分析

7.1 文件结构

1
2
3
4
5
6
7
8
9
10
11
crush/
├── crush.h # 核心数据结构定义
├── crush.c # 数据结构管理
├── mapper.h # 映射函数声明
├── mapper.c # 映射算法实现(核心)
├── builder.h # 构建函数声明
├── builder.c # CRUSH map构建
├── hash.h # 哈希函数声明
├── hash.c # 哈希函数实现
├── CrushWrapper.h # C++包装类
└── CrushWrapper.cc # C++包装实现

7.2 关键函数调用链

7.2.1 映射调用链

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
crush_do_rule()

crush_do_rule_no_retry()

遍历规则步骤

crush_choose_firstn() / crush_choose_indep()

crush_bucket_choose()

bucket_uniform_choose() / bucket_list_choose() / bucket_straw2_choose()

bucket_perm_choose() / generate_exponential_distribution()

crush_hash32_*()

7.2.2 桶选择函数

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
// mapper.c:368-399
static int crush_bucket_choose(
const struct crush_bucket *in,
struct crush_work_bucket *work,
int x, int r,
const struct crush_choose_arg *arg,
int position)
{
switch (in->alg) {
case CRUSH_BUCKET_UNIFORM:
return bucket_uniform_choose(...);
case CRUSH_BUCKET_LIST:
return bucket_list_choose(...);
case CRUSH_BUCKET_STRAW2:
return bucket_straw2_choose(...);
default:
return in->items[0];
}
}

7.3 工作空间管理

工作空间结构

1
2
3
4
5
6
7
8
9
struct crush_work {
struct crush_work_bucket **work; // 每个桶的工作空间
};

struct crush_work_bucket {
__u32 perm_x; // 排列对应的x值
__u32 perm_n; // 已计算的排列元素数
__u32 *perm; // 排列数组
};

初始化

1
2
3
4
5
6
7
8
9
10
11
size_t crush_work_size(const struct crush_map *map, int result_max)
{
// 计算所需工作空间大小
// 包括:work结构 + 桶指针数组 + 每个桶的工作空间
}

void crush_init_workspace(const struct crush_map *m, void *v)
{
// 初始化工作空间
// 分配每个桶的工作空间
}

7.4 重试机制

重试层次

  1. 本地重试local_retries):桶内重试,避免冲突
  2. 本地回退重试local_fallback_retries):桶内穷举搜索
  3. 下降重试retry_descent):重新开始下降过程
  4. 总重试choose_total_tries):总重试次数限制

重试逻辑crush_choose_firstn):

1
2
3
4
5
6
7
8
9
if (collide && flocal <= local_retries)
retry_bucket = 1; // 本地重试
else if (local_fallback_retries > 0 &&
flocal <= in->size + local_fallback_retries)
retry_bucket = 1; // 本地回退重试
else if (ftotal < tries)
retry_descent = 1; // 下降重试
else
skip_rep = 1; // 放弃

性能优化

8.1 Uniform Bucket优化

优化点

  1. r=0优化:直接计算第一个元素,避免完整排列
  2. 延迟排列:只在需要时计算排列元素
1
2
3
4
5
6
7
// 优化:r=0的情况
if (pr == 0) {
s = crush_hash32_3(bucket->hash, x, bucket->id, 0) % bucket->size;
work->perm[0] = s;
work->perm_n = 0xffff; // 标记
return bucket->items[s];
}

8.2 工作空间重用

优化

  • 工作空间可以在多次调用间重用
  • 只要CRUSH map不变,工作空间就有效
  • 减少内存分配开销

8.3 哈希函数优化

优化

  • 使用查找表加速自然对数计算(crush_ln
  • 内联函数减少函数调用开销
  • 位运算优化

8.4 选择算法选择

性能对比

算法 选择速度 添加项目 删除项目 重权重
Uniform O(1)
List O(n)
Straw2 O(n)

建议

  • 小规模、权重相同的桶:使用Uniform
  • 大规模、频繁变化的桶:使用Straw2
  • 只添加、不删除的场景:使用List

8.5 规则优化

优化建议

  1. 减少规则步骤:步骤越少,执行越快
  2. 合理使用类型:利用类型快速过滤
  3. 避免过深层次:层次越深,递归开销越大

总结

核心要点

  1. CRUSH是确定性伪随机算法:相同输入总是产生相同输出
  2. 支持多种桶算法:Uniform、List、Straw2,各有优缺点
  3. 规则驱动:通过规则定义映射行为
  4. 权重支持:根据权重进行负载均衡
  5. 容错机制:通过重试处理冲突和故障

算法优势

  • 可扩展性:支持大规模集群
  • 灵活性:通过规则和权重灵活控制
  • 稳定性:Straw2算法在变化时数据移动最小
  • 性能:Uniform算法选择速度最快

应用场景

  • 副本存储:使用FIRSTN模式
  • 纠删码:使用INDEP模式
  • 故障域隔离:通过类型和规则实现
  • 负载均衡:通过权重实现

参考资料

  1. CRUSH论文:http://www.ssrc.ucsc.edu/Papers/weil-sc06.pdf
  2. Ceph源码:cephMain/src/crush/
  3. CRUSH算法文档:Ceph官方文档

libaio、io_uring、SPDK 对比分析

libaio、io_uring、SPDK 对比分析

1. 概述

本文档对比分析三种高性能 I/O 接口:libaioio_uringSPDK,重点阐述它们的技术特点、性能差异和适用场景。

2. 基本概念

2.1 libaio(Linux Asynchronous I/O)

定义:Linux 传统的异步 I/O 库,基于 POSIX AIO 标准。

特点

  • 内核级异步 I/O 接口
  • 使用 io_setupio_submitio_getevents 等系统调用
  • 支持文件和块设备的异步 I/O
  • 较老的接口,但兼容性好

版本历史

  • 在 Linux 2.5+ 引入
  • 已有近 20 年历史

2.2 io_uring(Linux 5.1+)

定义:Linux 5.1 引入的新的异步 I/O 接口,旨在提供更高效的系统调用。

特点

  • 共享内存队列(shared memory ring buffers)
  • 批量系统调用(batch system calls)
  • 支持多种操作(I/O、网络、文件系统等)
  • 更少的系统调用开销

版本历史

  • Linux 5.1 引入基础功能
  • Linux 5.5+ 增加更多功能
  • Linux 5.19+ 性能进一步优化

2.3 SPDK(Storage Performance Development Kit)

定义:Intel 开发的高性能存储开发工具包,提供用户态、轮询模式的存储应用框架。

特点

  • 用户态驱动:绕过内核,直接访问硬件
  • 轮询模式:不使用中断,主动轮询完成队列
  • 零拷贝:最小化数据拷贝
  • DPDK 集成:基于 DPDK 环境抽象层

版本历史

  • 2014 年由 Intel 开源
  • 持续演进,支持多种存储协议和硬件

3. 架构对比

3.1 libaio 架构

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
┌─────────────────────────────────────────┐
│ 用户应用程序 │
│ - libaio API (io_setup, io_submit) │
└─────────────────────────────────────────┘
↓ 系统调用
┌─────────────────────────────────────────┐
│ Linux 内核 │
│ - io_submit() 提交 I/O 请求 │
│ - io_getevents() 获取完成事件 │
│ - 中断驱动完成通知 │
└─────────────────────────────────────────┘

┌─────────────────────────────────────────┐
│ 存储设备 │
│ - 块设备 / 文件系统 │
└─────────────────────────────────────────┘

关键特点

  • 系统调用:每次提交和获取需要系统调用
  • 中断模式:使用中断通知完成
  • 内核态:I/O 处理在内核中进行

3.2 io_uring 架构

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
┌─────────────────────────────────────────┐
│ 用户应用程序 │
│ - io_uring API (io_uring_setup) │
│ - 共享内存队列 (SQ/CQ) │
└─────────────────────────────────────────┘
↓ mmap 共享内存
┌─────────────────────────────────────────┐
│ Linux 内核 │
│ - io_uring_enter() 批量提交/获取 │
│ - 轮询模式 (IORING_SETUP_IOPOLL) │
│ - 中断模式 (默认) │
└─────────────────────────────────────────┘

┌─────────────────────────────────────────┐
│ 存储设备 │
│ - 块设备 / 文件系统 │
└─────────────────────────────────────────┘

关键特点

  • 共享内存:使用 mmap 共享队列
  • 批量操作:一次系统调用处理多个请求
  • 轮询模式:可选轮询模式(IORING_SETUP_IOPOLL)
  • 内核态:I/O 处理仍在内核中进行

3.3 SPDK 架构

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
┌─────────────────────────────────────────┐
│ 用户应用程序 │
│ - SPDK API (spdk_nvme_ns_cmd_read) │
│ - 用户态驱动 │
└─────────────────────────────────────────┘
↓ 直接访问(VFIO/IOMMU)
┌─────────────────────────────────────────┐
│ Linux 内核 │
│ - VFIO (Virtual Function I/O) │
│ - IOMMU 管理 │
│ - 最小内核参与 │
└─────────────────────────────────────────┘
↓ 直接 DMA
┌─────────────────────────────────────────┐
│ 存储设备 (NVMe SSD) │
│ - PCIe 设备 │
│ - 直接寄存器访问 │
└─────────────────────────────────────────┘

关键特点

  • 用户态驱动:完全在用户空间运行
  • 直接硬件访问:通过 VFIO/IOMMU 直接访问设备
  • 轮询模式:始终使用轮询,无中断
  • 绕过内核:最小化内核参与

4. 技术细节对比

4.1 I/O 模型

特性 libaio io_uring SPDK
I/O 模式 异步(回调) 异步(队列) 异步(回调/轮询)
完成通知 中断 + 轮询 中断 + 轮询(可选) 轮询(始终)
系统调用 频繁(每次提交/获取) 批量(一次处理多个) 极少(初始化时)
内核参与 完全在内核 完全在内核 最小化(仅 VFIO)

4.2 性能特征

特性 libaio io_uring SPDK
延迟 较高(~10-50μs) 低(~5-20μs) 极低(~1-5μs)
吞吐量 中等 极高
CPU 开销 中等 极低(轮询模式)
上下文切换 频繁 较少 极少

4.3 适用场景

特性 libaio io_uring SPDK
块设备 ✅ 支持 ✅ 支持 ✅ 支持
文件系统 ✅ 支持 ✅ 支持 ❌ 不支持
网络 I/O ❌ 不支持 ✅ 支持 ✅ 支持(通过 DPDK)
NVMe 设备 ✅ 通过内核驱动 ✅ 通过内核驱动 ✅ 用户态驱动
通用性 ✅ 通用 ✅ 通用 ⚠️ 特定场景

4.4 API 复杂度

特性 libaio io_uring SPDK
API 复杂度 中等 较高
学习曲线 平缓 中等 陡峭
代码示例 丰富 较少 较少
文档 完善 较少 完善

5. 代码示例对比

5.1 libaio 示例

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
#include <libaio.h>
#include <fcntl.h>
#include <unistd.h>

// 初始化
io_context_t ctx = 0;
io_setup(128, &ctx);

// 提交 I/O
struct iocb iocb, *iocbp = &iocb;
io_prep_pread(&iocb, fd, buf, size, offset);
io_submit(ctx, 1, &iocbp);

// 获取完成
struct io_event events[1];
io_getevents(ctx, 1, 1, events, NULL);

// 清理
io_destroy(ctx);

特点

  • 每次提交需要系统调用 io_submit()
  • 每次获取需要系统调用 io_getevents()
  • 系统调用开销较大

5.2 io_uring 示例

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
#include <liburing.h>

// 初始化
struct io_uring ring;
io_uring_queue_init(32, &ring, 0);

// 提交 I/O(无系统调用,直接写入共享内存)
struct io_uring_sqe *sqe = io_uring_get_sqe(&ring);
io_uring_prep_read(sqe, fd, buf, size, offset);
io_uring_sqe_set_data(sqe, buf);
io_uring_submit(&ring); // 批量提交,可能无系统调用

// 获取完成(无系统调用,直接读取共享内存)
struct io_uring_cqe *cqe;
io_uring_wait_cqe(&ring, &cqe); // 可能无系统调用
io_uring_cqe_seen(&ring, cqe);

// 清理
io_uring_queue_exit(&ring);

特点

  • 使用共享内存队列,减少系统调用
  • 批量提交/获取,提高效率
  • 支持轮询模式(IORING_SETUP_IOPOLL)

5.3 SPDK 示例

1
2
3
4
5
6
7
8
9
10
11
12
13
14
#include "spdk/nvme.h"

// 初始化(一次系统调用,VFIO 设置)
spdk_nvme_probe(NULL, probe_cb, attach_cb, NULL);

// 提交 I/O(无系统调用,直接写入 NVMe 寄存器)
spdk_nvme_ns_cmd_read(ns, qpair, buf, lba, lba_count,
read_complete, NULL, 0);

// 获取完成(轮询,无系统调用)
spdk_nvme_qpair_process_completions(qpair, 32);

// 清理
spdk_nvme_detach(ctrlr);

特点

  • 初始化时一次性系统调用(VFIO 设置)
  • I/O 操作无系统调用,直接访问硬件
  • 轮询模式,零延迟响应

6. 性能对比

6.1 延迟对比(假设场景:4KB 随机读)

接口 延迟 说明
libaio ~10-50μs 系统调用 + 中断延迟
io_uring(中断模式) ~5-20μs 批量系统调用 + 中断延迟
io_uring(轮询模式) ~2-10μs 批量系统调用 + 轮询
SPDK ~1-5μs 直接硬件访问 + 轮询

延迟来源

  • libaio:系统调用开销 + 中断延迟 + 上下文切换
  • io_uring:批量系统调用(减少) + 中断/轮询 + 上下文切换(减少)
  • SPDK:直接硬件访问 + 轮询(无中断,无上下文切换)

6.2 吞吐量对比(假设场景:顺序读)

接口 吞吐量 说明
libaio 中等(~2-4 GB/s) 受系统调用限制
io_uring 高(~5-8 GB/s) 批量操作提高效率
SPDK 极高(~8-15 GB/s) 零拷贝 + 轮询

6.3 CPU 利用率对比

接口 CPU 利用率 说明
libaio 中等(~30-50%) 系统调用 + 中断处理
io_uring 低(~20-40%) 批量操作 + 较少中断
SPDK 可调(~10-100%) 轮询模式,CPU 占用可控

注意:SPDK 轮询模式虽然延迟低,但 CPU 占用较高,需要平衡性能和 CPU 使用率。

7. 适用场景分析

7.1 libaio 适用场景

优点

  • 兼容性好,支持较老的 Linux 内核
  • API 简单,易于使用
  • 支持文件系统和块设备

缺点

  • 性能较低,系统调用开销大
  • 不支持网络 I/O
  • 功能有限

适用场景

  • 对性能要求不高的应用
  • 需要兼容老版本内核
  • 简单的异步 I/O 需求

7.2 io_uring 适用场景

优点

  • 性能好,批量操作减少系统调用
  • 功能丰富,支持 I/O、网络、文件系统
  • 支持轮询模式,可降低延迟
  • 通用性强

缺点

  • 需要 Linux 5.1+ 内核
  • API 相对复杂
  • 轮询模式仍在内核中

适用场景

  • 高性能存储应用
  • 需要低延迟的 I/O
  • 通用异步 I/O 需求
  • 现代 Linux 系统的最佳选择

7.3 SPDK 适用场景

优点

  • 性能极佳,微秒级延迟
  • 用户态驱动,绕过内核
  • 零拷贝,高吞吐量
  • 支持多种存储协议

缺点

  • 仅支持块设备,不支持文件系统
  • 需要 DPDK,依赖复杂
  • 轮询模式 CPU 占用高
  • 学习曲线陡峭

适用场景

  • 极致性能要求(延迟 < 10μs)
  • NVMe 设备直接访问
  • 存储服务器/存储网关
  • 云存储后端
  • 不适合文件系统 I/O

8. 在 Ceph/BlueStore 中的应用

8.1 BlueStore 中的选择

当前状态

  • BlueStore 使用 KernelDevice(通过 libaio)访问块设备
  • 支持通过 libaio 进行异步 I/O

可能的选择

  1. io_uring:可以替换 libaio,提高性能
  2. SPDK:通过 SPDK BDEV 模块,获得极致性能

8.2 SPDK 在 Ceph 中的使用

SPDK BDEV 模块

  • SPDK 提供了 AIO BDEV 模块(module/bdev/aio/
  • 可以封装 libaio,提供统一的 BDEV 接口
  • 也可以使用 SPDK NVMe BDEV,直接访问 NVMe 设备

架构

1
2
3
4
5
6
BlueStore
└── BlockDevice
├── KernelDevice (libaio) ← 当前方式
├── io_uring Device (可选) ← 未来可能
└── SPDK BDEV (可选) ← 高性能选项
└── SPDK NVMe Driver

9. 选择建议

9.1 如何选择

  1. 如果只需要文件系统 I/O

    • ❌ SPDK(不支持)
    • ✅ libaio 或 io_uring
  2. 如果只需要块设备 I/O

    • ✅ 都可以,按性能需求选择:
      • 性能要求低:libaio
      • 性能要求中等:io_uring
      • 性能要求极高:SPDK
  3. 如果需要网络 I/O

    • ❌ libaio(不支持)
    • ✅ io_uring 或 SPDK(通过 DPDK)
  4. 如果需要极致性能(延迟 < 10μs)

    • ❌ libaio、io_uring(无法达到)
    • ✅ SPDK(唯一选择)

9.2 性能 vs 复杂度权衡

1
2
3
4
性能:     libaio < io_uring < SPDK
复杂度: libaio < io_uring < SPDK
兼容性: libaio > io_uring > SPDK
通用性: libaio < io_uring > SPDK

9.3 推荐方案

  1. 通用场景:使用 io_uring(Linux 5.1+)

    • 性能好,兼容性强
    • 现代 Linux 系统的最佳选择
  2. 高性能存储:使用 SPDK

    • 极致性能,微秒级延迟
    • 适合存储服务器、存储网关
  3. 兼容性要求:使用 libaio

    • 支持老版本内核
    • 简单易用

10. 总结

10.1 核心差异

特性 libaio io_uring SPDK
定位 传统异步 I/O 现代异步 I/O 极致性能存储
内核参与 完全内核 完全内核 最小化内核
系统调用 频繁 批量 极少
延迟 10-50μs 2-20μs 1-5μs
适用场景 通用 通用 存储专用

10.2 选择建议

  1. 现代应用:优先使用 io_uring(Linux 5.1+)
  2. 极致性能:使用 SPDK(存储服务器)
  3. 兼容性:使用 libaio(老版本内核)

10.3 未来趋势

  • libaio:逐步被 io_uring 替代
  • io_uring:成为 Linux 异步 I/O 的主流选择
  • SPDK:在存储领域持续发展,性能不断优化

11. 参考资料

  • libaio: Linux man page io_setup(2)
  • io_uring: Linux man page io_uring_setup(2)
  • SPDK: https://spdk.io/
  • BlueStore: Ceph BlueStore 源码
  • SPDK BDEV AIO: module/bdev/aio/

RBD对象映射优化机制分析

RBD对象映射优化机制分析

问题

问题: 当RBD客户端的object-map为空时,是否不会向RADOS发起请求?

答案: 不会跳过。当object-map为空或无效时,RBD采用保守策略,仍然会发起RADOS请求。只有当object-map存在且明确标记对象为不存在时,才会跳过RADOS请求。


核心机制

1. object_may_exist() 方法

位置: librbd/ObjectMap.cc:154-177

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
template <typename I>
bool ObjectMap<I>::object_may_exist(uint64_t object_no) const
{
ceph_assert(ceph_mutex_is_locked(m_image_ctx.image_lock));

// 如果对象映射被禁用或无效,回退到默认逻辑(保守策略)
if (!m_image_ctx.test_features(RBD_FEATURE_OBJECT_MAP,
m_image_ctx.image_lock)) {
return true; // 保守策略:假设对象可能存在
}

bool flags_set;
int r = m_image_ctx.test_flags(m_image_ctx.snap_id,
RBD_FLAG_OBJECT_MAP_INVALID,
m_image_ctx.image_lock, &flags_set);
if (r < 0 || flags_set) {
return true; // 如果object-map无效,返回true(保守策略)
}

uint8_t state = (*this)[object_no];
bool exists = (state != OBJECT_NONEXISTENT);
ldout(m_image_ctx.cct, 20) << "object_no=" << object_no << " r=" << exists
<< dendl;
return exists;
}

关键逻辑:

  1. 如果object-map未启用: 返回true(保守策略)
  2. 如果object-map无效: 返回true(保守策略)
  3. 如果object-map有效: 根据对象状态返回
    • OBJECT_NONEXISTENT → 返回false
    • 其他状态 → 返回true

2. 读取请求优化

位置: librbd/io/ObjectRequest.cc:223-256

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
template <typename I>
void ObjectReadRequest<I>::read_object() {
I *image_ctx = this->m_ictx;

std::shared_lock image_locker{image_ctx->image_lock};
auto read_snap_id = this->m_io_context->get_read_snap();

// 关键优化点:如果object-map标记对象不存在,直接读取父镜像
if (read_snap_id == image_ctx->snap_id &&
image_ctx->object_map != nullptr &&
!image_ctx->object_map->object_may_exist(this->m_object_no)) {
// 跳过RADOS读取,直接尝试从父镜像读取
image_ctx->asio_engine->post([this]() { read_parent(); });
return;
}
image_locker.unlock();

// 如果object-map为空或标记对象存在,正常发起RADOS读取请求
neorados::ReadOp read_op;
// ... 构建读取操作
image_ctx->rados_api.execute(
{data_object_name(this->m_ictx, this->m_object_no)},
*this->m_io_context, std::move(read_op), nullptr,
// ...
);
}

优化逻辑:

  • 条件1: read_snap_id == image_ctx->snap_id (读取当前快照)
  • 条件2: image_ctx->object_map != nullptr (object-map存在)
  • 条件3: !object_may_exist() (对象不存在)

当所有条件满足时: 跳过RADOS读取请求,直接尝试从父镜像读取

当object-map为空时: object_may_exist()返回true不会跳过RADOS请求


3. 写入请求优化

位置: librbd/io/ObjectRequest.cc:412-437

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
template <typename I>
void AbstractObjectWriteRequest<I>::send() {
I *image_ctx = this->m_ictx;

{
std::shared_lock image_lock{image_ctx->image_lock};
if (image_ctx->object_map == nullptr) {
m_object_may_exist = true; // object-map为空,假设对象存在
} else {
m_object_may_exist = image_ctx->object_map->object_may_exist(
this->m_object_no);
}
}

// 优化:如果对象不存在且操作是no-op,直接完成
if (!m_object_may_exist && is_no_op_for_nonexistent_object()) {
ldout(image_ctx->cct, 20) << "skipping no-op on nonexistent object"
<< dendl;
this->async_finish(0);
return; // 跳过RADOS写入请求
}

pre_write_object_map_update();
}

优化逻辑:

  • 如果object-map为空: m_object_may_exist = true不会跳过写入请求
  • 如果object-map存在且对象不存在:
    • 对于no-op操作(如全零写入、discard等),可能跳过写入
    • 对于实际写入操作,仍然会发起RADOS请求(需要创建对象)

对象映射状态

对象状态定义

1
2
3
4
5
6
enum {
OBJECT_NONEXISTENT = 0, // 对象不存在
OBJECT_EXISTS = 1, // 对象存在
OBJECT_PENDING = 2, // 对象待定(正在创建)
OBJECT_EXISTS_CLEAN = 3 // 对象存在且干净(未修改)
};

object-map为空的情况

当object-map为空时,意味着:

  1. 未启用object-map功能: RBD_FEATURE_OBJECT_MAP未设置
  2. object-map对象不存在: RADOS中不存在rbd_object_map.<image_id>对象
  3. object-map无效: RBD_FLAG_OBJECT_MAP_INVALID标志被设置

在这些情况下,object_may_exist()都会返回true,采用保守策略。


优化效果

1. 读取优化

场景: 读取稀疏镜像(大部分对象不存在)

优化前:

1
读取对象 → RADOS读取请求 → 返回ENOENT → 读取父镜像

优化后 (object-map有效且标记对象不存在):

1
读取对象 → 检查object-map → 跳过RADOS请求 → 直接读取父镜像

性能提升: 减少不必要的RADOS请求,降低延迟

2. 写入优化

场景: 对不存在对象执行no-op操作(如全零写入、discard)

优化前:

1
写入对象 → RADOS写入请求 → 创建对象 → 执行操作

优化后 (object-map有效且标记对象不存在):

1
写入对象 → 检查object-map → 判断为no-op → 直接完成(跳过RADOS请求)

性能提升: 避免创建不必要的对象


代码流程总结

读取请求流程

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
ObjectReadRequest::read_object()

├─→ 检查object-map
│ │
│ ├─→ object-map为空/无效
│ │ └─→ object_may_exist() = true
│ │ └─→ 发起RADOS读取请求 ✓
│ │
│ └─→ object-map有效
│ │
│ ├─→ object_may_exist() = true (对象存在)
│ │ └─→ 发起RADOS读取请求 ✓
│ │
│ └─→ object_may_exist() = false (对象不存在)
│ └─→ 跳过RADOS请求,直接读取父镜像 ✗

写入请求流程

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
AbstractObjectWriteRequest::send()

├─→ 检查object-map
│ │
│ ├─→ object-map为空
│ │ └─→ m_object_may_exist = true
│ │ └─→ 发起RADOS写入请求 ✓
│ │
│ └─→ object-map有效
│ │
│ ├─→ object_may_exist() = true
│ │ └─→ 发起RADOS写入请求 ✓
│ │
│ └─→ object_may_exist() = false
│ │
│ ├─→ is_no_op_for_nonexistent_object() = true
│ │ └─→ 跳过RADOS请求,直接完成 ✗
│ │
│ └─→ is_no_op_for_nonexistent_object() = false
│ └─→ 发起RADOS写入请求(需要创建对象) ✓

结论

  1. object-map为空时:

    • object_may_exist()返回true(保守策略)
    • 仍然会发起RADOS请求
    • 不会跳过任何I/O操作
  2. object-map有效时:

    • 如果对象标记为不存在,可能跳过RADOS请求
    • 读取:直接读取父镜像
    • 写入:no-op操作可能直接完成
  3. 优化前提:

    • object-map必须存在且有效
    • 必须启用RBD_FEATURE_OBJECT_MAP功能
    • object-map不能标记为无效
  4. 设计理念:

    • 保守策略: 当不确定时,假设对象存在,避免跳过实际存在的对象
    • 性能优化: 当确定对象不存在时,跳过不必要的RADOS请求

相关代码位置

  • ObjectMap::object_may_exist(): librbd/ObjectMap.cc:154-177
  • 读取优化: librbd/io/ObjectRequest.cc:223-256
  • 写入优化: librbd/io/ObjectRequest.cc:412-437
  • ObjectMap定义: librbd/ObjectMap.h

文档版本: 1.0
最后更新: 2024年

MDS 业务功能详细分析

MDS 业务功能详细分析

1. 概述

MDS(Metadata Server,元数据服务器)是 CephFS 的核心组件,负责管理分布式文件系统的元数据,包括目录树、文件属性、权限、快照等。本文档详细分析 ceph_mds.cc 的业务功能以及 MDS 的完整功能集。

1.1 ceph_mds.cc 的作用

ceph_mds.cc 是 MDS 守护进程的入口点,主要负责:

  • MDS 守护进程的启动和初始化
  • 信号处理
  • 守护进程化
  • MDSDaemon 的创建和管理

1.2 MDS 架构

1
2
3
4
5
6
7
8
9
10
11
12
13
ceph_mds.cc (入口)

MDSDaemon (守护进程管理)

MDSRank (Rank 管理)
├── Server (客户端请求处理)
├── MDCache (元数据缓存)
├── MDLog (元数据日志)
├── Locker (锁管理)
├── SnapServer/SnapClient (快照管理)
├── MDBalancer (负载均衡)
├── Migrator (目录迁移)
└── SessionMap (会话管理)

2. ceph_mds.cc 业务功能分析

2.1 启动流程

2.1.1 初始化阶段

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
int main(int argc, const char **argv)
{
// 1. 设置线程名称
ceph_pthread_setname("ceph-mds");

// 2. 全局初始化
auto cct = global_init(NULL, args,
CEPH_ENTITY_TYPE_MDS,
CODE_ENVIRONMENT_DAEMON, 0);

// 3. NUMA 节点亲和性配置
// 4. 参数解析
// 5. 守护进程化(Preforker)
// 6. 网络地址选择
// 7. MDS ID 验证
// 8. 创建 Messenger
// 9. 创建 MDSDaemon
// 10. 初始化 MDSDaemon
}

关键步骤

  1. NUMA 亲和性配置

    • 支持将 MDS 绑定到特定 NUMA 节点
    • 提高性能,减少跨节点内存访问
  2. MDS ID 验证

    • MDS ID 不能为空
    • MDS ID 不能以数字开头(避免与 Rank 混淆)
  3. 守护进程化

    • 使用 Preforker 实现守护进程化
    • 父进程等待子进程初始化后退出

2.1.2 消息传递器配置

1
2
3
4
5
6
7
8
// 创建消息传递器
Messenger *msgr = Messenger::create(...);

// 配置消息策略
msgr->set_default_policy(Messenger::Policy::lossy_client(required));
msgr->set_policy(entity_name_t::TYPE_MON, ...); // Monitor:有损客户端
msgr->set_policy(entity_name_t::TYPE_MDS, ...); // MDS:无损对等体
msgr->set_policy(entity_name_t::TYPE_CLIENT, ...); // 客户端:有状态服务器

消息策略说明

  • Monitor:有损客户端,连接失败可重连
  • MDS:无损对等体,保证消息可靠传输
  • Client:有状态服务器,维护客户端会话

2.1.3 信号处理

1
2
3
4
// 注册信号处理器
register_async_signal_handler(SIGHUP, sighup_handler); // 重新加载配置
register_async_signal_handler_oneshot(SIGINT, handle_mds_signal); // 中断
register_async_signal_handler_oneshot(SIGTERM, handle_mds_signal); // 终止

信号处理

  • SIGHUP:重新加载配置
  • SIGINT/SIGTERM:优雅关闭 MDS

2.2 MDSDaemon 创建

1
2
3
4
5
// 创建 MDSDaemon
mds = new MDSDaemon(g_conf()->name.get_id().c_str(), msgr, &mc, ctxpool);

// 初始化
r = mds->init();

MDSDaemon 职责

  • 管理 MDSRank
  • 处理 MDSMap 更新
  • 处理管理命令
  • 心跳和健康检查

3. MDS 核心功能模块

3.1 元数据管理(MDCache)

功能:管理文件系统元数据的缓存

核心组件

  • CInode:Inode 缓存对象
  • CDir:目录缓存对象
  • CDentry:目录项缓存对象

主要功能

  1. 元数据缓存

    • 缓存 Inode、目录、目录项
    • LRU 淘汰策略
    • 内存压力管理
  2. 路径遍历

    • 根据路径查找 Inode
    • 支持多种遍历标志(发现、锁定等)
    • 跨 MDS 路径解析
  3. 目录碎片化

    • 大目录自动碎片化
    • 目录合并
    • 碎片管理
  4. Stray 管理

    • 处理删除的文件
    • Stray 目录管理
    • 延迟清理

3.2 客户端请求处理(Server)

功能:处理客户端发送的文件系统操作请求

支持的操作

3.2.1 文件操作

  • LOOKUP:查找文件/目录
  • GETATTR:获取文件属性
  • SETATTR:设置文件属性
  • OPEN:打开文件
  • CREATE:创建文件
  • MKNOD:创建设备节点
  • UNLINK:删除文件
  • LINK:创建硬链接
  • SYMLINK:创建符号链接
  • READLINK:读取符号链接

3.2.2 目录操作

  • MKDIR:创建目录
  • RMDIR:删除目录
  • READDIR:读取目录
  • RENAME:重命名文件/目录

3.2.3 扩展属性操作

  • GETXATTR:获取扩展属性
  • SETXATTR:设置扩展属性
  • LISTXATTR:列出扩展属性
  • REMOVEXATTR:删除扩展属性

3.2.4 快照操作

  • MKSNAP:创建快照
  • RMSNAP:删除快照
  • LSSNAP:列出快照
  • LOOKUPSNAP:查找快照

3.2.5 其他操作

  • FLUSH:刷新数据
  • FSYNC:同步文件系统
  • SETFILELOCK:设置文件锁
  • GETFILELOCK:获取文件锁

请求处理流程

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

Server::dispatch_client_request()

Server::handle_client_request()

路径解析(MDCache::path_traverse)

权限检查(Locker::check_caps)

执行操作(do_* 方法)

日志记录(MDLog::submit_entry)

发送回复(send_reply)

3.3 锁管理(Locker)

功能:管理文件系统的各种锁

锁类型

  1. SimpleLock:简单锁

    • 用于文件大小、时间戳等属性
    • 支持读锁和写锁
  2. ScatterLock:分散锁

    • 用于目录统计信息
    • 支持分散更新
  3. FileLock:文件锁

    • 用于文件范围锁定
    • 支持共享锁和排他锁
  4. LocalLock:本地锁

    • 用于本地操作
    • 不跨 MDS 传播

锁管理功能

  • 锁获取和释放
  • 锁升级和降级
  • 锁超时处理
  • 死锁检测

3.4 元数据日志(MDLog)

功能:记录元数据变更日志,用于恢复和复制

日志事件类型

  1. EUpdate:更新事件

    • 记录元数据变更
    • 包含完整的变更信息
  2. ESession:会话事件

    • 记录客户端会话
    • 用于会话恢复
  3. ESubtreeMap:子树映射事件

    • 记录目录树分布
    • 用于负载均衡
  4. EFragment:碎片事件

    • 记录目录碎片化
    • 用于目录管理
  5. ESegment:段事件

    • 日志段管理
    • 用于日志轮转

日志功能

  • 日志提交
  • 日志回放
  • 日志轮转
  • 日志同步

3.5 快照管理(SnapServer/SnapClient)

功能:管理文件系统快照

核心组件

  • SnapServer:快照服务器(主 MDS)
  • SnapClient:快照客户端(从 MDS)
  • SnapRealm:快照域

主要功能

  1. 快照创建

    • 创建快照
    • 快照元数据管理
    • 快照版本控制
  2. 快照删除

    • 删除快照
    • 清理快照数据
    • 快照引用计数
  3. 快照查询

    • 列出快照
    • 查找快照
    • 快照差异计算
  4. 快照回滚

    • 回滚到快照
    • 数据恢复

3.6 负载均衡(MDBalancer)

功能:在多个 MDS 之间平衡元数据负载

负载均衡策略

  1. 基于目录的负载均衡

    • 将目录树分布到不同 MDS
    • 根据负载情况迁移目录
  2. 动态负载均衡

    • 监控各 MDS 负载
    • 自动迁移热点目录
  3. 负载均衡算法

    • 基于目录大小的均衡
    • 基于访问频率的均衡
    • 基于 MDS 能力的均衡

负载均衡触发条件

  • MDS 负载差异超过阈值
  • 新 MDS 加入
  • MDS 故障恢复

3.7 目录迁移(Migrator)

功能:在 MDS 之间迁移目录子树

迁移类型

  1. 导出(Export)

    • 将目录子树导出到其他 MDS
    • 更新目录映射
    • 通知客户端
  2. 导入(Import)

    • 从其他 MDS 导入目录子树
    • 重建缓存
    • 更新引用

迁移流程

1
2
3
4
5
6
1. 选择迁移目标
2. 冻结目录操作
3. 导出目录数据
4. 更新目录映射
5. 通知客户端
6. 完成迁移

3.8 会话管理(SessionMap)

功能:管理客户端会话

会话状态

  • OPEN:会话打开
  • STALE:会话过期
  • CLOSING:会话关闭中
  • CLOSED:会话已关闭

会话功能

  • 会话创建和销毁
  • 会话超时处理
  • 会话恢复
  • 会话统计

3.9 权限管理(Capability)

功能:管理客户端对文件和目录的访问权限

Capability 类型

  1. CAP_AUTH:认证权限
  2. CAP_READ:读权限
  3. CAP_WRITE:写权限
  4. CAP_EXEC:执行权限
  5. CAP_ANY:任意权限

Capability 管理

  • Capability 授予
  • Capability 回收
  • Capability 刷新
  • Capability 迁移

3.10 文件锁管理(flock)

功能:管理 POSIX 文件锁

锁类型

  • F_RDLCK:读锁(共享锁)
  • F_WRLCK:写锁(排他锁)
  • F_UNLCK:解锁

锁管理

  • 锁获取和释放
  • 锁冲突检测
  • 锁超时处理
  • 死锁检测

3.11 恢复队列(RecoveryQueue)

功能:管理文件大小恢复

恢复场景

  • MDS 重启后恢复文件大小
  • 从日志中恢复文件大小
  • 修复不一致的文件大小

3.12 清理队列(PurgeQueue)

功能:管理已删除文件的清理

清理流程

  1. 文件删除后进入 Stray 目录
  2. 延迟清理(等待引用释放)
  3. 从 OSD 删除数据
  4. 清理元数据

3.13 打开文件表(OpenFileTable)

功能:跟踪打开的文件

功能

  • 记录打开的文件
  • 文件引用计数
  • 文件关闭通知
  • 文件迁移处理

3.14 指标处理(MetricsHandler)

功能:收集和报告 MDS 性能指标

指标类型

  • 请求延迟
  • 请求吞吐量
  • 缓存命中率
  • 内存使用
  • CPU 使用

3.15 健康检查(Beacon)

功能:向 Monitor 发送心跳,报告 MDS 健康状态

健康信息

  • MDS 状态(Active、Standby、Stopped)
  • MDS 负载
  • MDS 版本
  • 故障信息

4. MDS 状态管理

4.1 MDS 状态

  1. Standby:待机状态

    • MDS 已启动但未分配 Rank
    • 等待成为 Active 或 Standby-Replay
  2. Standby-Replay:待机回放状态

    • 跟随 Active MDS 回放日志
    • 准备接管
  3. Active:活跃状态

    • 处理客户端请求
    • 管理分配的 Rank
  4. Stopped:停止状态

    • MDS 已停止
    • 等待清理

4.2 Rank 管理

Rank:MDS 的职责范围,每个 Rank 由一个 Active MDS 管理

Rank 分配

  • Monitor 分配 Rank
  • 一个 Rank 只能有一个 Active MDS
  • 可以有多个 Standby MDS

Rank 功能

  • 目录树分片
  • 负载均衡
  • 故障隔离

5. MDS 关键特性

5.1 高可用性

  1. Standby MDS

    • 多个 Standby MDS 待命
    • Active MDS 故障时自动接管
  2. 日志回放

    • Standby-Replay MDS 回放日志
    • 快速故障恢复
  3. 故障检测

    • Beacon 心跳机制
    • 自动故障转移

5.2 可扩展性

  1. 多 MDS 支持

    • 多个 Active MDS
    • 目录树分片
  2. 动态负载均衡

    • 自动迁移热点目录
    • 平衡 MDS 负载
  3. 目录碎片化

    • 大目录自动碎片化
    • 提高并发性能

5.3 一致性保证

  1. 日志机制

    • 所有元数据变更记录日志
    • 支持日志回放
  2. 锁机制

    • 分布式锁管理
    • 保证操作原子性
  3. 版本控制

    • 元数据版本管理
    • 冲突检测和解决

5.4 性能优化

  1. 元数据缓存

    • 内存缓存热点元数据
    • LRU 淘汰策略
  2. 批量操作

    • 批量提交日志
    • 批量更新元数据
  3. 异步 I/O

    • 异步处理客户端请求
    • 提高并发性能

6. MDS 管理功能

6.1 Admin Socket 命令

通过 Admin Socket 可以执行以下命令:

  1. status:显示 MDS 状态
  2. lockup:锁定测试
  3. exit:退出 MDS
  4. respawn:重启 MDS
  5. dump cache:转储缓存
  6. flush journal:刷新日志
  7. scrub:数据校验
  8. export dir:导出目录
  9. import dir:导入目录

6.2 监控和统计

  1. 性能计数器

    • 请求延迟
    • 请求吞吐量
    • 缓存命中率
  2. 内存统计

    • Inode 数量
    • 目录数量
    • 目录项数量
  3. 操作统计

    • 各种操作的计数
    • 操作延迟分布

7. MDS 与其他组件的交互

7.1 与 Monitor 的交互

  • MDSMap 同步:接收 MDSMap 更新
  • 心跳报告:通过 Beacon 发送心跳
  • 状态报告:报告 MDS 状态和负载

7.2 与 OSD 的交互

  • 元数据存储:将元数据存储在 OSD
  • 文件数据:文件数据存储在 OSD
  • 快照数据:快照数据存储在 OSD

7.3 与客户端的交互

  • 请求处理:处理客户端文件系统操作
  • 会话管理:管理客户端会话
  • 权限管理:管理客户端权限

7.4 与其他 MDS 的交互

  • 目录迁移:在 MDS 之间迁移目录
  • 日志同步:同步日志到 Standby MDS
  • 负载均衡:协调负载均衡

8. 总结

8.1 ceph_mds.cc 的核心职责

  1. 启动管理:MDS 守护进程的启动和初始化
  2. 信号处理:处理系统信号(SIGHUP、SIGINT、SIGTERM)
  3. 守护进程化:实现守护进程化
  4. MDSDaemon 管理:创建和管理 MDSDaemon

8.2 MDS 的完整功能集

  1. 元数据管理:Inode、目录、目录项的管理
  2. 客户端服务:处理客户端文件系统操作
  3. 锁管理:分布式锁管理
  4. 日志管理:元数据变更日志
  5. 快照管理:文件系统快照
  6. 负载均衡:多 MDS 负载均衡
  7. 目录迁移:目录子树迁移
  8. 会话管理:客户端会话管理
  9. 权限管理:Capability 管理
  10. 文件锁:POSIX 文件锁
  11. 恢复机制:文件大小恢复
  12. 清理机制:已删除文件清理
  13. 监控统计:性能指标收集

8.3 MDS 的关键特性

  • 高可用性:Standby MDS、日志回放、故障转移
  • 可扩展性:多 MDS、负载均衡、目录碎片化
  • 一致性:日志机制、锁机制、版本控制
  • 高性能:元数据缓存、批量操作、异步 I/O

8.4 相关文件

  • ceph_mds.cc:MDS 守护进程入口
  • MDSDaemon.h/cc:MDS 守护进程管理
  • MDSRank.h/cc:MDS Rank 管理
  • Server.h/cc:客户端请求处理
  • MDCache.h/cc:元数据缓存
  • MDLog.h/cc:元数据日志
  • Locker.h/cc:锁管理
  • SnapServer.h/cc:快照服务器
  • MDBalancer.h/cc:负载均衡
  • Migrator.h/cc:目录迁移

BlueStore硬盘IO调用链分析

BlueStore硬盘IO调用链分析

目录

  1. 问题定义
  2. 总体结论
  3. BlueStore如何发起读写
  4. BlockDevice如何选择后端
  5. 内核AIO路径
  6. SPDK路径
  7. 回调与完成通知
  8. 读写调用链小结

问题定义

本文回答的问题是:

  • BlueStore如何调用aio接口读写硬盘
  • BlueStore如何调用SPDK接口读写硬盘
  • 这两条路径分别在代码中的哪些函数里发生

这里的重点不是 BlueStore 的完整对象模型,而是 从 BlueStore 发起一次磁盘读写,到最终落到内核 AIO 或 SPDK NVMe 命令 的调用链。


总体结论

BlueStore 自己并不直接调用 libaioio_uringspdk_nvme_ns_cmd_readv/writev

它采用的是一层统一抽象:

1
2
3
4
BlueStore
-> BlockDevice 抽象接口
-> KernelDevice(内核 AIO / io_uring)
-> NVMEDevice(SPDK)

也就是说:

  • BlueStore 只调用 bdev->aio_read() / bdev->aio_write() / bdev->aio_submit()
  • 真正决定走内核后端还是 SPDK 后端的是 BlockDevice::create()
  • 真正下发到底层设备的是 KernelDeviceNVMEDevice

BlueStore如何发起读写

1. 打开块设备时创建 bdev

BlueStore 在 _open_bdev() 中创建主块设备对象:

1
2
3
4
5
6
7
8
bdev = BlockDevice::create(
cct,
p,
aio_cb,
static_cast<void*>(this),
discard_cb,
static_cast<void*>(this),
"bluestore");

这里有两个关键信息:

  • bdev 的静态类型是 BlockDevice*
  • 创建时传入了完成回调 aio_cb

因此 BlueStore 后续读写都只需要调用 bdev 抽象接口,而不用关心底层到底是内核设备还是 SPDK 设备。

2. 读路径如何发起

BlueStore 在 _do_read() 中,会把需要从块设备读取的物理区间收集起来,然后调用:

  • bdev->aio_read(...)
  • bdev->aio_submit(&ioc)
  • ioc.aio_wait()

这表示:

  1. 先把一个或多个读请求挂入 IOContext
  2. 再统一提交
  3. 然后等待完成

3. 写路径如何发起

写路径也一样。

BlueStore 不会每产生一个小写就立刻自己直接打到底层,而是:

  1. 先调用 bdev->aio_write(...) 把 IO 挂入当前事务的 IOContext
  2. 等事务状态机推进到需要下发数据 IO 时
  3. 再由 _txc_aio_submit() 调用 bdev->aio_submit(&txc->ioc)

所以 BlueStore 的设计是:

  • 读路径:边组装边挂请求,最后统一 submit
  • 写路径:事务收集写 IO,状态机推进时统一 submit

BlockDevice如何选择后端

1. 入口函数

设备类型选择发生在:

  • BlockDevice::create()
  • BlockDevice::create_with_type()

2. 选择逻辑

BlockDevice::create() 会先看配置项 bdev_type

  • 如果用户显式配置了 aio,就走内核设备
  • 如果显式配置了 spdk,就走 SPDK 设备
  • 如果没有显式配置,则自动探测

自动探测逻辑是:

  • 如果路径支持 NVMEDevice::support(path),则走 spdk
  • 否则默认走 aio

3. 分派结果

最终在 create_with_type() 中实例化具体对象:

  • KernelDevice
  • NVMEDevice
  • PMEMDevice

对于本文主题,重点是前两者:

  • KernelDevice:传统内核块设备路径
  • NVMEDevice:SPDK NVMe 直通路径

内核AIO路径

1. 具体实现类

内核后端由 KernelDevice 实现,文件在:

  • src/blk/kernel/KernelDevice.cc
  • src/blk/kernel/KernelDevice.h

2. aio_write() 做什么

KernelDevice::aio_write() 的核心动作是:

  1. 检查 IO 对齐
  2. 必要时把 bufferlist 重建成满足 direct I/O 对齐要求的形式
  3. 生成 aio_t
  4. 把请求挂到 ioc->pending_aios

也就是说,aio_write() 本身通常只是“登记请求”,不是立刻提交。

3. aio_read() 做什么

KernelDevice::aio_read() 逻辑类似:

  1. 分配对齐缓冲
  2. 构造 aio_t
  3. 调用 aio.preadv(off, len)
  4. 把它挂到 IOContext::pending_aios

4. aio_submit() 做什么

KernelDevice::aio_submit() 才是真正提交批量 IO 的地方。

它会:

  1. pending_aios 挪到 running_aios
  2. 增加 num_running
  3. 调用 io_queue->submit_batch(...)

这里的 io_queue 有两种实现:

  • ioring_queue_t
  • aio_queue_t

初始化时如果系统支持且配置启用 bdev_ioring,优先用 io_uring
否则回退到传统 libaio

所以“内核 AIO 路径”实际上又分成:

  • io_uring 路径
  • libaio 路径

但对 BlueStore 来说,这两者都被 KernelDevice 屏蔽了。

5. 完成通知

内核设备有自己的 AIO 处理线程 _aio_thread()

底层提交后,完成事件由内核队列返回,KernelDevice 会:

  1. 回收完成的 aio_t
  2. 更新 IOContext
  3. 在适当时机唤醒等待者
  4. 通过创建时注册的回调把完成事件交回上层

因此 BlueStore 并不自己轮询内核完成队列,而是由 KernelDevice 负责。


SPDK路径

1. 具体实现类

SPDK 后端由 NVMEDevice 实现,文件在:

  • src/blk/spdk/NVMEDevice.cc
  • src/blk/spdk/NVMEDevice.h

2. aio_write() 做什么

NVMEDevice::aio_write() 并不直接调用 SPDK 写命令。

它会先调用 write_split(...),把大写请求拆成若干个 Task

  • 每个 Task 记录偏移、长度、命令类型
  • 再把这些 Task 追加到 IOContext

这一步本质上是“把写请求转换成 SPDK 可提交的任务链表”。

3. aio_read() 做什么

NVMEDevice::aio_read() 也类似。

它会:

  1. 准备目标缓冲
  2. 调用 make_read_tasks(...)
  3. 把读请求拆成多个 Task
  4. 追加进 IOContext

4. aio_submit() 做什么

NVMEDevice::aio_submit() 会:

  1. IOContext 里取出挂好的任务链
  2. num_pending 转成 num_running
  3. 创建或复用当前线程的 SharedDriverQueueData
  4. 调用 queue_t._aio_handle(t, ioc)

5. _aio_handle() 如何真正调用 SPDK

SharedDriverQueueData::_aio_handle() 是 SPDK 路径的核心。

对每个任务:

  • 写命令调用 spdk_nvme_ns_cmd_writev(...)
  • 读命令调用 spdk_nvme_ns_cmd_readv(...)
  • flush 调用 spdk_nvme_ns_cmd_flush(...)

这就是 BlueStore 最终触发 SPDK 读写硬盘的真正位置。

6. SPDK完成处理方式

SPDK 路径不是依赖内核 io_getevents,而是用户态轮询:

  • spdk_nvme_qpair_process_completions(...)

_aio_handle() 在循环里持续轮询 completions:

  1. 如果队列里有在飞 IO,就处理 completions
  2. 如果没有完成,就按配置 sleep 很短时间
  3. 完成后由 SPDK completion callback 回调 io_complete(...)

所以 SPDK 路径的特征是:

  • 用户态 qpair
  • 用户态提交
  • 用户态 polling completion

回调与完成通知

BlueStore 在创建 BlockDevice 时传了:

  • aio_cb
  • discard_cb

其中 aio_cb 的作用是把设备层完成通知重新送回 BlueStore:

  1. 设备层完成某个 AIO
  2. 调用注册回调
  3. BlueStore 的 AioContext::aio_finish(store) 被触发
  4. 事务状态机继续向前推进

因此:

  • BlueStore 不直接处理内核完成队列
  • BlueStore 也不直接处理 SPDK completion queue
  • 这些都先由具体设备实现处理,再通过回调交回 BlueStore

读写调用链小结

1. 读路径

1
2
3
4
5
6
7
8
9
10
BlueStore::_do_read()
-> bdev->aio_read(...)
-> bdev->aio_submit(&ioc)
-> KernelDevice::aio_submit()
-> io_uring/libaio submit

-> NVMEDevice::aio_submit()
-> spdk_nvme_ns_cmd_readv()
-> ioc.aio_wait()
-> aio_cb -> BlueStore::AioContext::aio_finish()

2. 写路径

1
2
3
4
5
6
7
8
9
10
11
BlueStore::_write() / _do_write() / Writer.cc
-> bdev->aio_write(...)
-> BlueStore::_txc_aio_submit()
-> bdev->aio_submit(&txc->ioc)
-> KernelDevice::aio_submit()
-> io_uring/libaio submit

-> NVMEDevice::aio_submit()
-> spdk_nvme_ns_cmd_writev()
-> 完成回调 aio_cb
-> 事务状态机继续推进

3. 一句话总结

BlueStore 的关键点只有一句话:

BlueStore 只面向 BlockDevice 抽象编程,而 KernelDeviceNVMEDevice 分别把同一套 aio_read/aio_write/aio_submit 接口落到内核 AIO 与 SPDK NVMe 命令。


附:关键源码位置

  • src/os/bluestore/BlueStore.cc
    • _open_bdev()
    • _do_read()
    • _txc_aio_submit()
  • src/blk/BlockDevice.cc
    • BlockDevice::create()
    • BlockDevice::create_with_type()
  • src/blk/BlockDevice.h
    • BlockDevice::aio_read()
    • BlockDevice::aio_write()
    • BlockDevice::aio_submit()
  • src/blk/kernel/KernelDevice.cc
    • KernelDevice::aio_read()
    • KernelDevice::aio_write()
    • KernelDevice::aio_submit()
    • _aio_thread()
  • src/blk/spdk/NVMEDevice.cc
    • NVMEDevice::aio_read()
    • NVMEDevice::aio_write()
    • NVMEDevice::aio_submit()
    • SharedDriverQueueData::_aio_handle()

BlueStore软件架构流程分析

BlueStore软件架构流程分析

目录

  1. BlueStore概述
  2. 整体架构
  3. 核心组件
  4. 数据模型
  5. 写入流程
  6. 读取流程
  7. 事务处理
  8. 空间管理
  9. 缓存机制
  10. 压缩和校验

BlueStore概述

1.1 什么是BlueStore

BlueStore 是Ceph的默认对象存储后端,直接在块设备上存储对象数据,避免了传统文件系统的开销。它是为高性能SSD设计的存储后端。

1.2 核心特点

  • 直接管理块设备:绕过文件系统,直接在块设备上操作
  • 元数据存储在RocksDB:使用RocksDB存储所有元数据
  • BlueFS文件系统:轻量级文件系统,用于存储RocksDB的WAL和SST文件
  • 写时分配(COW):支持克隆和快照
  • 内联压缩:支持数据压缩
  • 校验和:支持数据完整性校验
  • 多设备支持:支持WAL、DB、慢速设备分离

1.3 与传统文件系统的区别

特性 传统文件系统(FileStore) BlueStore
数据存储 通过文件系统 直接在块设备
元数据 文件系统元数据 RocksDB
性能 受文件系统限制 更高性能
开销 双重写入(数据+元数据) 单次写入
适用场景 通用场景 SSD优化

整体架构

2.1 架构层次

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
┌─────────────────────────────────────────────────────────┐
│ ObjectStore接口 │
│ (OSD层调用的存储接口) │
└───────────────────────┬───────────────────────────────────┘

┌───────────────────────▼───────────────────────────────────┐
│ BlueStore核心 │
│ ┌──────────────┐ ┌──────────────┐ ┌──────────────┐ │
│ │ Collection │ │ Onode │ │ Buffer │ │
│ │ (集合) │ │ (对象节点) │ │ (缓冲区) │ │
│ └──────────────┘ └──────────────┘ └──────────────┘ │
│ ┌──────────────┐ ┌──────────────┐ ┌──────────────┐ │
│ │ Extent │ │ Blob │ │ SharedBlob │ │
│ │ (逻辑范围) │ │ (数据块) │ │ (共享块) │ │
│ └──────────────┘ └──────────────┘ └──────────────┘ │
└───────────────────────┬───────────────────────────────────┘

┌───────────────┼───────────────┐
│ │ │
┌───────▼──────┐ ┌─────▼──────┐ ┌─────▼──────┐
│ Allocator │ │ Freelist │ │ BlueFS │
│ (分配器) │ │ Manager │ │ (文件系统) │
│ │ │ (空闲管理) │ │ │
└───────┬──────┘ └─────┬──────┘ └─────┬──────┘
│ │ │
┌───────▼───────────────▼───────────────▼──────┐
│ RocksDB (元数据) │
│ (通过BlueFS存储WAL和SST) │
└───────────────────────┬───────────────────────┘

┌───────────────────────▼───────────────────────┐
│ 块设备 (BlockDevice) │
│ ┌──────────┐ ┌──────────┐ ┌──────────┐ │
│ │ WAL │ │ DB │ │ Slow │ │
│ │ (快速) │ │ (快速) │ │ (慢速) │ │
│ └──────────┘ └──────────┘ └──────────┘ │
└───────────────────────────────────────────────┘

2.2 设备布局

BlueStore支持三种设备类型:

  1. WAL设备:存储RocksDB的WAL(Write-Ahead Log)文件

    • 要求:低延迟、高IOPS
    • 通常:NVMe SSD
  2. DB设备:存储RocksDB的SST文件

    • 要求:中等延迟、高IOPS
    • 通常:SSD
  3. 慢速设备(Slow):存储对象数据

    • 要求:大容量
    • 通常:HDD或大容量SSD

设备布局示例

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
块设备布局:
┌─────────────────────────────────────────┐
│ Label (设备标签) │
│ - OSD UUID │
│ - 设备大小 │
│ - 创建时间 │
└─────────────────────────────────────────┘
┌─────────────────────────────────────────┐
│ BlueFS (如果使用BlueFS) │
│ - WAL文件 │
│ - SST文件 │
└─────────────────────────────────────────┘
┌─────────────────────────────────────────┐
│ BlueStore数据区 │
│ - 对象数据 │
│ - 空闲空间由Allocator管理 │
└─────────────────────────────────────────┘

核心组件

3.1 BlueStore类

位置BlueStore.h / BlueStore.cc

职责

  • 实现ObjectStore接口
  • 管理所有存储操作(read、write、transaction等)
  • 协调各个子组件

关键成员

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
class BlueStore : public ObjectStore {
// 设备管理
BlockDevice *bdev; // 主数据设备
BlockDevice *bdev_slow; // 慢速设备
BlockDevice *bdev_wal; // WAL设备
BlockDevice *bdev_db; // DB设备

// 元数据存储
KeyValueDB *db; // RocksDB实例

// 文件系统
BlueFS *bluefs; // BlueFS文件系统

// 空间管理
Allocator *alloc; // 空间分配器
FreelistManager *fm; // 空闲列表管理器

// 缓存
OnodeCache *onode_cache; // Onode缓存
BufferCache *buffer_cache; // 缓冲区缓存

// 集合管理
ceph::unordered_map<coll_t, CollectionRef> coll_map;
};

3.2 Collection(集合)

位置BlueStore.h 内部结构

职责

  • 管理一个PG(Placement Group)的对象集合
  • 维护Onode缓存
  • 管理事务上下文

关键成员

1
2
3
4
5
6
7
8
9
10
11
12
13
14
struct Collection {
coll_t cid; // 集合ID
bluestore_cnode_t cnode; // 集合节点元数据

// Onode管理
OnodeCacheShard *onode_cache; // Onode缓存分片
ceph::unordered_map<ghobject_t, OnodeRef> onode_map;

// 共享Blob管理
SharedBlobSet shared_blob_set;

// 统计信息
PerfCounters *logger;
};

3.3 Onode(对象节点)

位置BlueStore.h 内部结构

职责

  • 表示一个对象的元数据
  • 管理对象的逻辑到物理映射(ExtentMap)
  • 管理对象的缓冲区(BufferSpace)

关键成员

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
struct Onode {
ghobject_t oid; // 对象ID
bluestore_onode_t onode; // 持久化元数据

// 逻辑到物理映射
ExtentMap extent_map; // Extent映射表

// 缓冲区管理
BufferSpace buffer_space; // 缓冲区空间

// OMAP管理
OMapSpace omap_space; // OMAP空间

// 引用计数
std::atomic_int nref;
};

Onode元数据结构bluestore_onode_t):

1
2
3
4
5
6
7
8
9
struct bluestore_onode_t {
uint64_t nid; // 节点ID
uint64_t size; // 对象大小
utime_t mtime; // 修改时间
utime_t atime; // 访问时间
uint32_t expected_object_size; // 期望对象大小
uint32_t alloc_hint_flags; // 分配提示标志
// ... 其他元数据
};

3.4 Extent(逻辑范围)

位置BlueStore.h 内部结构

职责

  • 表示对象的一个逻辑范围
  • 映射到物理Blob

关键成员

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

3.5 Blob(数据块)

位置BlueStore.h 内部结构

职责

  • 表示物理存储的数据块
  • 管理物理extent(pextent)
  • 支持压缩和校验

关键成员

1
2
3
4
5
6
7
8
9
10
struct Blob {
uint64_t id; // Blob ID
bluestore_blob_t blob; // 持久化Blob元数据

// 物理extent列表
PExtentVector extents; // 物理extent向量

// 共享Blob引用
SharedBlobRef shared_blob; // 共享Blob引用
};

Blob元数据结构bluestore_blob_t):

1
2
3
4
5
6
7
8
9
10
11
struct bluestore_blob_t {
uint64_t id; // Blob ID
uint32_t logical_length; // 逻辑长度
uint32_t compressed_length; // 压缩后长度
uint8_t compression; // 压缩算法
uint8_t csum_type; // 校验和类型
uint8_t csum_chunk_order; // 校验和块大小
uint32_t flags; // 标志位
PExtentVector extents; // 物理extent列表
bluestore_extent_ref_map_t ref_map; // 引用映射
};

3.6 SharedBlob(共享Blob)

位置BlueStore.h 内部结构

职责

  • 管理共享Blob的引用计数
  • 支持写时复制(COW)
  • 管理共享Blob的缓冲区缓存

关键成员

1
2
3
4
5
6
7
struct SharedBlob {
std::atomic_int nref; // 引用计数
bluestore_shared_blob_t *persistent; // 持久化元数据

// 引用映射
bluestore_extent_ref_map_t ref_map; // 引用映射表
};

3.7 Buffer(缓冲区)

位置BlueStore.h 内部结构

职责

  • 缓存对象数据
  • 管理写入状态
  • 支持读写缓存

状态

  • STATE_CLEAN:干净状态(已写入磁盘)
  • STATE_WRITING:正在写入
  • STATE_DIRTY:脏数据(待写入)

3.8 BlueFS(文件系统)

位置BlueFS.h / BlueFS.cc

职责

  • 轻量级文件系统
  • 管理RocksDB的WAL和SST文件
  • 直接在块设备上操作

关键特性

  • 日志结构文件系统
  • 支持三种设备(WAL、DB、Slow)
  • 文件分配和回收

3.9 Allocator(分配器)

位置Allocator.h / 各种实现

职责

  • 管理空闲空间
  • 分配物理extent
  • 支持多种分配算法

分配器类型

  • StupidAllocator:简单分配器
  • BitmapAllocator:位图分配器
  • AvlAllocator:AVL树分配器
  • BtreeAllocator:B树分配器
  • Btree2Allocator:B树2分配器
  • HybridAllocator:混合分配器

3.10 FreelistManager(空闲列表管理器)

位置FreelistManager.h / 实现

职责

  • 持久化空闲空间信息到RocksDB
  • 管理空闲空间的分配和释放
  • 支持多种实现(Bitmap、Btree等)

数据模型

4.1 逻辑到物理映射

1
2
3
4
5
6
7
8
9
10
11
12
13
对象 (Object)

├─► Onode (对象元数据)
│ │
│ └─► ExtentMap (逻辑范围映射)
│ │
│ ├─► Extent 1: [0, 1MB) → Blob A
│ ├─► Extent 2: [1MB, 2MB) → Blob B
│ └─► Extent 3: [2MB, 3MB) → Blob A (共享)

└─► BufferSpace (缓冲区空间)

└─► Buffer缓存 (内存中的数据)

4.2 物理存储结构

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
Blob (数据块)

├─► 物理Extent列表
│ ├─► PExtent 1: [offset=100MB, length=1MB]
│ ├─► PExtent 2: [offset=200MB, length=1MB]
│ └─► PExtent 3: [offset=500MB, length=1MB]

├─► 压缩信息
│ ├─► 压缩算法: snappy
│ ├─► 原始长度: 2MB
│ └─► 压缩长度: 1.5MB

└─► 校验和信息
├─► 校验和类型: crc32c
└─► 校验和块: 每64KB一个校验和

4.3 RocksDB键空间

BlueStore使用以下键前缀:

1
2
3
4
5
6
7
8
PREFIX_SUPER = "S"        // 超级块信息
PREFIX_STAT = "T" // 统计信息
PREFIX_COLL = "C" // Collection元数据
PREFIX_OBJ = "O" // Onode元数据
PREFIX_OMAP = "M" // OMAP键值对
PREFIX_DEFERRED = "L" // 延迟事务
PREFIX_ALLOC = "B" // 空闲空间(FreelistManager)
PREFIX_SHARED_BLOB = "X" // 共享Blob元数据

4.4 对象键编码

对象键的编码格式:

1
2
3
[shard_id + 0x80] + [pool_id + 2^63] + [hash(bit_reversed)] + 
[namespace(escaped)] + [key(escaped)] + ['<'|'='|'>'] +
[object_name(escaped)] + [snap] + [generation] + ['o']

写入流程

5.1 写入流程概览

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
客户端写入请求


BlueStore::write()

├─► 获取或创建Collection
├─► 获取或创建Onode
└─► 创建事务上下文 (TransContext)


Writer::do_write()

├─► 处理缓冲区 (BufferSpace)
├─► 分配空间 (Allocator)
├─► 创建Blob
├─► 压缩数据 (可选)
├─► 计算校验和
└─► 调度I/O


BlueStore::_txc_add_transaction()

├─► 准备事务 (prepare)
├─► 提交I/O (aio_submit)
├─► 等待I/O完成 (aio_wait)
└─► 提交元数据 (kv_commit)


事务完成

5.2 详细写入步骤

步骤1:接收写入请求

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
int BlueStore::write(
coll_t cid, // 集合ID
const ghobject_t& oid, // 对象ID
uint64_t offset, // 偏移
size_t len, // 长度
const bufferlist& bl, // 数据
uint32_t fadvise_flags) // 建议标志
{
// 1. 获取Collection
CollectionRef c = _get_collection(cid);

// 2. 获取或创建Onode
OnodeRef o = c->get_onode(oid, true);

// 3. 创建事务上下文
TransContext *txc = _get_trans_context(c);

// 4. 执行写入
Writer w(this, txc, wctx, o);
w.do_write(offset, bl);

// 5. 提交事务
_txc_add_transaction(txc);
return 0;
}

步骤2:Writer处理写入

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
void Writer::do_write(uint32_t location, bufferlist& data)
{
// 1. 处理数据分割和压缩
blob_vec blobs;
_split_data(location, data, blobs);

// 2. 对齐到磁盘块
_align_to_disk_block(location, ref_end, blobs);

// 3. 处理现有Blob的修改
_do_put_blobs(location, data_end, ref_end, blobs, after_punch_it);

// 4. 分配新空间(如果需要)
_defer_or_allocate(need_size);

// 5. 创建新Blob
_do_put_new_blobs(location, ref_end, bd_it, bd_end);
}

步骤3:空间分配

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
void Writer::_defer_or_allocate(uint32_t need_size)
{
// 1. 检查是否需要延迟写入
if (should_defer) {
// 延迟写入:先写入缓冲区,稍后分配空间
do_deferred = true;
return;
}

// 2. 从Allocator分配空间
PExtentVector extents;
int64_t allocated = alloc->allocate(
need_size, // 需要的大小
block_size, // 块大小
max_alloc_size, // 最大分配大小
hint, // 分配提示
&extents); // 输出的extent列表

// 3. 记录分配的extent
allocated.insert(allocated.end(), extents.begin(), extents.end());
}

步骤4:Blob创建和数据放置

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
BlobRef Writer::_blob_create_with_data(
uint32_t in_blob_offset,
bufferlist& disk_data)
{
// 1. 创建Blob
BlobRef b = new Blob();
b->id = generate_blob_id();

// 2. 分配物理extent
PExtentVector extents;
_get_disk_space(disk_data.length(), extents);
b->extents = extents;

// 3. 设置Blob元数据
b->blob.logical_length = disk_data.length();
b->blob.compressed_length = compressed ? compressed_size : 0;
b->blob.compression = compression_algorithm;

// 4. 计算校验和
if (csum_type) {
calculate_checksums(b, disk_data);
}

return b;
}

步骤5:调度I/O

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
void Writer::_schedule_io(
const PExtentVector& disk_extents,
bufferlist data)
{
// 1. 将数据分割到extent
for (auto& e : disk_extents) {
uint64_t disk_offset = e.offset;
uint32_t length = e.length;

// 2. 创建AIO请求
IOContext *ioc = new IOContext();
ioc->aio_write(disk_offset, data);

// 3. 添加到事务上下文
txc->aio_write_queue.push_back(ioc);
}
}

步骤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
void BlueStore::_txc_add_transaction(TransContext *txc)
{
// 1. 准备阶段
_txc_prepare(txc);
├─► 更新Onode元数据
├─► 更新ExtentMap
├─► 更新Blob元数据
└─► 准备RocksDB事务

// 2. 提交I/O
_txc_submit_aio(txc);
└─► 提交所有AIO请求到块设备

// 3. 等待I/O完成
_txc_aio_wait(txc);
└─► 等待所有AIO完成

// 4. 提交元数据
_txc_commit_kv(txc);
└─► 提交RocksDB事务

// 5. 完成阶段
_txc_finish(txc);
├─► 更新缓存
├─► 释放资源
└─► 更新统计信息
}

5.3 小写优化

对于小写(小于bluestore_small_io_size),BlueStore使用特殊优化:

  1. 内联写入:小数据直接写入Onode的OMAP
  2. 延迟分配:先写入缓冲区,稍后批量分配空间
  3. 合并写入:多个小写合并为一个大写

5.4 大写处理

对于大写(大于bluestore_big_io_size),BlueStore:

  1. 直接分配:立即分配空间
  2. 直接写入:直接写入块设备,不经过缓冲区
  3. 支持压缩:如果启用压缩,会压缩数据

读取流程

6.1 读取流程概览

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
客户端读取请求


BlueStore::read()

├─► 获取Collection
├─► 获取Onode
└─► 从BufferSpace读取

├─► 检查缓冲区缓存
│ ├─► 命中:直接返回
│ └─► 未命中:继续

├─► 查找ExtentMap
│ └─► 确定需要读取的Blob

├─► 从Blob读取数据
│ ├─► 查找物理extent
│ ├─► 调度AIO读取
│ ├─► 等待I/O完成
│ ├─► 验证校验和
│ └─► 解压缩(如果需要)

└─► 更新缓存
└─► 将数据加入BufferSpace

6.2 详细读取步骤

步骤1:接收读取请求

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 BlueStore::read(
coll_t cid, // 集合ID
const ghobject_t& oid, // 对象ID
uint64_t offset, // 偏移
size_t len, // 长度
bufferlist& bl, // 输出缓冲区
uint32_t op_flags) // 操作标志
{
// 1. 获取Collection
CollectionRef c = _get_collection(cid);

// 2. 获取Onode
OnodeRef o = c->get_onode(oid, false);
if (!o) {
return -ENOENT;
}

// 3. 从BufferSpace读取
ready_regions_t ready;
interval_set<uint32_t> missing;
o->buffer_space.read(
cache, // 缓存
offset, // 偏移
len, // 长度
ready, // 已就绪区域
missing); // 缺失区域

// 4. 从磁盘读取缺失区域
if (!missing.empty()) {
_read_from_disk(o, missing, bl);
}

return 0;
}

步骤2:从BufferSpace读取

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
void BufferSpace::read(
BufferCacheShard* cache,
uint32_t offset,
uint32_t length,
ready_regions_t& res,
interval_set<uint32_t>& res_intervals)
{
// 1. 查找覆盖范围的Buffer
auto it = _data_lower_bound(offset);
uint32_t pos = offset;
uint32_t end = offset + length;

while (pos < end && it != buffer_map.end()) {
Buffer& b = *it;

// 2. 检查Buffer是否覆盖当前位置
if (b.offset <= pos && pos < b.offset + b.length) {
// Buffer覆盖,添加到结果
uint32_t b_off = pos - b.offset;
uint32_t b_len = min(b.length - b_off, end - pos);
res[pos] = b.data.substr(b_off, b_len);
res_intervals.insert(pos, b_len);
pos += b_len;
} else if (pos < b.offset) {
// 有间隙,需要从磁盘读取
res_intervals.insert(pos, b.offset - pos);
pos = b.offset;
} else {
++it;
}
}

// 3. 处理末尾间隙
if (pos < end) {
res_intervals.insert(pos, end - pos);
}
}

步骤3:从磁盘读取

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
void BlueStore::_read_from_disk(
OnodeRef o,
interval_set<uint32_t>& missing,
bufferlist& bl)
{
// 1. 遍历缺失区域
for (auto& p : missing) {
uint32_t offset = p.first;
uint32_t length = p.second;

// 2. 查找ExtentMap
auto extents = o->extent_map.get_containing_extents(offset, length);

// 3. 对每个Extent读取
for (auto& ext : extents) {
BlobRef blob = ext.blob;
uint32_t blob_offset = ext.blob_offset + (offset - ext.logical_offset);

// 4. 从Blob读取
_read_blob(blob, blob_offset, length, bl);
}
}
}

步骤4:从Blob读取

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
void BlueStore::_read_blob(
BlobRef blob,
uint32_t blob_offset,
uint32_t length,
bufferlist& bl)
{
// 1. 查找物理extent
PExtentVector extents = blob->get_extents_containing(blob_offset, length);

// 2. 调度AIO读取
vector<IOContext*> ios;
for (auto& e : extents) {
uint64_t disk_offset = e.offset + (blob_offset - e.logical_offset);
IOContext *ioc = new IOContext();
ioc->aio_read(disk_offset, length);
ios.push_back(ioc);
}

// 3. 提交AIO
bdev->aio_submit(ios);

// 4. 等待完成
for (auto ioc : ios) {
ioc->wait();

// 5. 验证校验和
if (blob->blob.csum_type) {
verify_checksum(blob, ioc->bl);
}

// 6. 解压缩(如果需要)
if (blob->blob.compression) {
decompress(blob, ioc->bl);
}

bl.append(ioc->bl);
}

// 7. 更新缓存
o->buffer_space.did_read(cache, blob_offset, std::move(bl));
}

6.3 读取优化

  1. 预读(Read-ahead):预测性读取相邻数据
  2. 缓存命中:优先从BufferSpace读取
  3. 并行读取:多个extent并行读取
  4. 校验和验证:读取时验证数据完整性

事务处理

7.1 事务模型

BlueStore使用两阶段提交:

  1. 准备阶段(Prepare)

    • 更新内存中的元数据
    • 准备RocksDB事务
    • 调度AIO操作
  2. 提交阶段(Commit)

    • 等待AIO完成
    • 提交RocksDB事务
    • 更新缓存状态

7.2 TransContext(事务上下文)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
struct TransContext {
CollectionRef collection; // 所属集合
ObjectStore::Transaction t; // 事务对象

// I/O队列
list<IOContext*> aio_write_queue; // 写入队列
list<IOContext*> aio_read_queue; // 读取队列

// 统计信息
uint64_t bytes_written;
uint64_t bytes_read;

// 状态
enum {
STATE_PREPARE, // 准备中
STATE_AIO_WAIT, // 等待AIO
STATE_IO_DONE, // I/O完成
STATE_KV_QUEUED, // KV已排队
STATE_KV_COMMITTING, // KV提交中
STATE_KV_DONE, // KV完成
STATE_FINISHING, // 完成中
STATE_DONE // 完成
} state;
};

7.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
26
27
_txc_add_transaction(txc)


_txc_prepare(txc)
├─► 更新Onode元数据
├─► 更新ExtentMap
├─► 更新Blob元数据
├─► 更新FreelistManager
└─► 准备RocksDB事务


_txc_submit_aio(txc)
└─► 提交所有AIO到块设备


_txc_aio_wait(txc)
└─► 等待所有AIO完成


_txc_commit_kv(txc)
└─► 提交RocksDB事务


_txc_finish(txc)
├─► 更新缓存状态
├─► 释放资源
└─► 更新统计信息

空间管理

8.1 分配流程

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
需要分配空间


Allocator::allocate()

├─► 查找空闲extent
│ ├─► StupidAllocator: 线性搜索
│ ├─► BitmapAllocator: 位图查找
│ ├─► AvlAllocator: AVL树查找
│ └─► BtreeAllocator: B树查找

├─► 选择最佳extent
│ └─► 考虑hint、连续性等

└─► 返回分配的extent列表


FreelistManager::allocate()
└─► 更新RocksDB中的空闲列表

8.2 释放流程

1
2
3
4
5
6
7
8
9
10
11
12
13
需要释放空间


Allocator::release()

├─► 将extent标记为空闲
│ └─► 更新分配器内部结构

└─► 合并相邻空闲extent(可选)


FreelistManager::release()
└─► 更新RocksDB中的空闲列表

8.3 碎片整理

BlueStore支持在线碎片整理:

  1. 识别碎片:扫描Blob,识别碎片化区域
  2. 迁移数据:将数据迁移到连续区域
  3. 更新映射:更新ExtentMap
  4. 释放旧空间:释放原来的extent

缓存机制

8.1 Onode缓存

目的:缓存对象的元数据(Onode)

实现

  • 使用LRU策略
  • 分片缓存(减少锁竞争)
  • 支持固定(pinned)Onode

缓存结构

1
2
3
4
5
struct OnodeCacheShard {
ceph::mutex lock; // 锁
LRUOnodeCache cache; // LRU缓存
map<ghobject_t, OnodeRef> pinned; // 固定的Onode
};

8.2 Buffer缓存

目的:缓存对象数据

实现

  • 使用PriorityCache
  • 支持多级缓存
  • 支持写回缓存

缓存状态

  • CLEAN:数据与磁盘一致
  • WRITING:正在写入
  • DIRTY:数据已修改,待写入

8.3 缓存策略

  1. 读缓存:读取的数据加入缓存
  2. 写缓存:小写先写入缓存
  3. 预读:预测性读取
  4. 淘汰:LRU淘汰策略

压缩和校验

9.1 压缩

支持的算法

  • snappy:快速压缩
  • zlib:标准压缩
  • lz4:快速压缩
  • zstd:高性能压缩

压缩流程

1
2
3
4
5
6
7
8
9
10
11
12
13
14
写入数据


尝试压缩

├─► 压缩率检查
│ └─► 如果压缩率 < threshold,使用压缩

├─► 压缩数据
│ └─► 使用选择的压缩算法

└─► 存储压缩数据
├─► 设置compressed_length
└─► 设置compression算法

压缩配置

  • bluestore_compression_algorithm:压缩算法
  • bluestore_compression_min_blob_size:最小压缩大小
  • bluestore_compression_max_blob_size:最大压缩大小
  • bluestore_compression_required_ratio:压缩率要求

9.2 校验和

支持的算法

  • crc32c:CRC32校验
  • xxhash32:XXHash校验
  • xxhash64:XXHash64校验

校验和流程

1
2
3
4
5
6
7
8
9
10
写入数据


计算校验和

├─► 按块计算(每64KB一个)
│ └─► 块大小由csum_chunk_order决定

└─► 存储校验和
└─► 存储在Blob元数据中

读取验证

1
2
3
4
5
6
7
8
9
读取数据


验证校验和

├─► 读取数据块
├─► 计算校验和
├─► 与存储的校验和比较
└─► 如果不匹配,返回EIO错误

总结

核心优势

  1. 高性能:直接操作块设备,避免文件系统开销
  2. 灵活性:支持多种分配算法和压缩算法
  3. 可靠性:支持校验和、事务、写时复制
  4. 可扩展性:支持多设备、大容量

关键设计决策

  1. 元数据与数据分离:元数据存储在RocksDB,数据直接存储在块设备
  2. 写时分配:写入时分配空间,支持延迟分配
  3. 多级缓存:Onode缓存和Buffer缓存
  4. 异步I/O:使用AIO提高并发性能

适用场景

  • SSD存储:针对SSD优化
  • 高性能要求:需要低延迟、高IOPS
  • 大容量存储:支持PB级存储
  • 云存储:适合云环境部署

参考资料

  1. Ceph源码:cephMain/src/os/bluestore/
  2. BlueStore文档:Ceph官方文档
  3. RocksDB文档:https://rocksdb.org/