< Back

Raft在CubeFS 纠删码系统中的应用实践

2023-12-22Zongchao Hu

导语

Raft算法是一种分布式共识算法,用于在分布式系统中维护一致性状态。它通过选举机制和日志复制来确保系统中的各个节点达成一致的状态。Raft算法相对于其他一致性算法(如Paxos)更易于理解和实现,因此在分布式系统中得到了广泛的应用。

本文主要介绍CubeFS纠删码(Erasure Coding)服务的元数据管理模块(Blobstore/ClusterMgr)、Raft算法实践以及日常运维建议。

背景

CubeFS的纠删码存储子系统(BlobStore),是一个高可靠、高可用、低成本、支持EB规模的独立Blob存储系统。通过之前的系列分享文章,大家对纠删码已经有了较深入的理解,本文将继续带领大家探讨更多内部实现细节,将着重于介绍如何基于Raft实现强一致的元数据管理服务。

图片

图1. Blobstore架构图

如图1所示,元数据模块(Cluster Manager, CM)在BlobStore中扮演着至关重要的角色,用于集群资源(如磁盘、节点、存储空间单元)管理、服务注册与发现、后台任务存储以及集群配置管理等。

模型设计

CM设计实现充分融合分层设计理念,其中主要由服务接口、状态机、协议层以及持久化以及网络层组成。整体架构可参考图2, 每层具体功能描述如下:

服务层:对外提供CM的接口服务

状态机:包括 VolumeMgr, DiskMgr, ScopeMgr, ServiceMgr, ConfigMgr, KvMgr 六个模块

  • VolumeMgr:卷管理模块,负责卷的创建、续租、更新卷与磁盘的映射等
  • DiskMgr:负责磁盘及空间管理,包括空间 chunk 分配、磁盘注册、心跳更新磁盘信息等
  • ScopeMgr: bid 、vid 及 diskid 等分配 scope 段管理模块
  • ServierMgr:通用的服务注册发现管理模块,负责 allocator 、 mqproxy 等服务的注册与发现
  • ConfigMgr:存取集群配置、任务开关等,以KvMgr作为存储底座
  • KvMgr:通用的配置 K-V 存取管理模块, 保存后台任务、配置信息等

协议层:基于 Raft 协议保证数据一致性,采用rpc进行数据同步通信

持久化层:使用 Rocksdb 持久化数据,按数据类型包含四种基础表: normal、volume、raft、kv

网络层: 负责数据同步、消息传输等

图片

图2. CM设计架构图

组件介绍

etcd/raft

  1. etcd/raft 为纯粹的状态机模型,所有的事件都以消息的形式处理。它主要有两个结构:
    • Raft: 状态机,根据消息触发状态转换
    • Node: 代表raft 集群中的一个节点,提供上层与 raft 状态机交互的接口
  2. Raft 结构体维护了当前节点的 Raft 状态,如 Term, Vote, Lead 等,是非线程安全的,提供方法触发状态变化。主要的方法如下:
    • tick: 每次调用会增加内部的逻辑时钟,用于触发超时事件如 leader election 和 heartbeat;
    • step: 类型为 type stepFunc func(r *raft, m pb.Message) error,用于处理消息;
    • 还有其他方法用于处理成员变更、获取状态等。
  3. 应用层处理 Storage 和 Transportation。etcd/raft 对这两个模块的解耦方式不太相同,由于 Storage 和 Raft 联系紧密,所以 etcd/raft 提供了 Storage 接口,由用户实现,然后在启动时设置到 Config 中,Storage 接口只涉及查询操作,何时持久化和如何持久化由用户决定;而 Transportation 的设计完全由用户自定义 。
// Node represents a node in a raft cluster.
type Node interface {
	Step(ctx context.Context, msg pb.Message) error //提供消息的写入
	Ready() <-chan Ready // 读取返回结果
	ApplyConfChange(cc pb.ConfChangeI) *pb.ConfState // 成员配置变更
	TransferLeadership(ctx context.Context, lead, transferee uint64) // 切换节点主从状态
	ReadIndex(ctx context.Context, rctx []byte) error // 读
	Status() Status // 返回状态值
	Stop() // 停止raft
}

RaftServer

RaftServer 对 etcd/raft 进行封装、并实现 wal 日志写入和存储、快照机制。RaftServer 为上层应用(如 ClusterMgr)提供上传、读取、成员管理、wal 日志裁剪以及 raft 状态查询等能力。

type RaftServer interface {
	Stop() // 停止raft节点服务,其中主要关闭raft、rpc和快照所占用的系统资源
	Propose(ctx context.Context, data []byte) error // 上报消息
	ReadIndex(ctx context.Context) error  // 基于ReadIndex的方式读取
	TransferLeadership(ctx context.Context, leader, transferee uint64) // 切主
	AddMember(ctx context.Context, member Member) error // 添加raft成员
	RemoveMember(ctx context.Context, nodeID uint64) error // 移除raft成员
	IsLeader() bool // 判断节点是否是主节点
	Status() Status // 当前节点的状态

	// In order to prevent log expansion, the application needs to call this method.
	Truncate(index uint64) error // wal日志裁剪
}

WAL

WAL日志(Write-Ahead Log,预写日志)将每次状态更新抽象为一个命令并追加写入一个日志中,也就是顺序写入,保证 IO 性能。在 Blobstore 中,wal 模块实现了etcd/raft 所预留的持久化接口,提供日志存储、查询以及裁剪功能。

type Storage interface {
	InitialState() (pb.HardState, pb.ConfState, error) // 初始状态
	Entries(lo, hi, maxSize uint64) ([]pb.Entry, error) // 获取范围内的entry
	Term(i uint64) (uint64, error) // 根据log号得到Term
	LastIndex() (uint64, error) // 
	FirstIndex() (uint64, error) //
	Snapshot() (pb.Snapshot, error)// 快照元数据
}
// Storage the storage
type Wal struct {
	// Log Entry
	logfiles    []logName // 日志按固定文件大小,多文件存储,便于日志压缩裁剪
	last        *logFile  // 最新的log文件句柄
	cache       *logFileCache // 日志文件缓存
	hs          pb.HardState  // raft转态元数据
	st          Snapshot      // 快照元数据的日志
	mt          *meta         // wal 元数据
}

一致性

下面通过磁盘注册过程,来讲解CM是如何保证磁盘信息在集群的一致性的。

当 CM 收到来自 blobnode(管理物理磁盘的存储单元)添加磁盘的请求后,首先判断自己是否是Leader(若不是会将请求转发给Leader),然后确认该磁盘是集群内的新的磁盘等符合添加新盘的条件后,由于CM 提供了多个模块的管理,StateMachine 会生成形如以下的结构信息作为一条 raft 提议,具体包含请求操作涉及的模块名、模块内特定的操作类型、具体数据以及其他信息。

ModuleNameOpTypeDiskInfoOthers

具体到代码表现为

// cubefs/blobstore/clustermgr/disk.go
base.EncodeProposeInfo(s.DiskMgr.GetModuleName(), diskmgr.OperTypeAddDisk, data, base.ProposeContext{ReqID: span.TraceID()})

StateMachine 将提议提交给 raftServer,进一步包装成一条 raft 能识别的消息体

// cubefs/blobstore/clustermgr/disk.go
s.raftNode.Propose(ctx, proposeInfo)// 阻塞等待
...
pb.Message{Type: pb.MsgProp, Entries: []pb.Entry{{Type: pr.entryType, Data: pr.b}}}

raft 将这条消息基于 RPC 通过 raft 端口转发给其他 raft 节点,并确认大多数节点已同步完成后,将这条消息返回给 StateMachine。根据 Raft 的状态机安全性(State Machine Safety)可知,相同的日志索引执行的指令一定是相同的,尽管每个节点可能在不同时间执行指令,最终状态机结果一致即可。 主节点发出日志同步请求,从节点接收并同步日志,每个CM节点状态机根据日志顺序执行指令,最终保证整个集群状态的一致性。

// cubefs/blobstore/common/raftserver/server.go
case rd := <-s.n.Ready():

具体来说,StateMachine 首先会根据消息中的模块名,将这个消息分发到具体的模块,模块收到后根据操作类型进行内存状态变更以及持久化到DB操作

// cubefs/blobstore/clustermgr/diskmgr/applier.go
case OperTypeAddDisk:
  	 diskInfo := &blobnode.DiskInfo{}
  	 err := json.Unmarshal(datas[idx], diskInfo)
  	 if err != nil {
  	 	errs[idx] = errors.Info(err, "json unmarshal failed, data: ", datas[idx]).Detail(err)
  	 	wg.Done()
  	 	continue
  	 }
  	 ...

整个步骤完成后,通知阻塞的请求返回成功或者具体错误信息给 blobnode,从而整个注册添加磁盘的流程便完成。以上是主节点视角,下面切换到从节点视角,从节点收到已是包装好的 raft 消息体,因此可以直接进入 raft 状态机轮转,直至主节点告诉这条消息可以应用至 StateMachine,同主节点一样,改变本节点的内存状态和持久化到 DB。

// cubefs/blobstore/common/raftserver/server.go
func (s *raftServer) handleMessage(msgs raftMsgs) error {
	...
	for i := 0; i < msgs.Len(); i++ {
		if err := s.n.Step(ctx, msgs[i]); err != nil {
			return err
		}
	}
	return nil
}

可用性

在上面介绍了CM如何保证一致性,以及通过注册磁盘的信息简要演示集群如何来保证一致性。那可用性又如何来实现呢?

WAL+DB

在CM集群中每次写入都会先经过WAL来持久化,WAL记录元数据集群中的每一次变更操作,在节点宕机恢复时,通过重放WAL记录可恢复到节点宕机前状态。

// cubefs/blobstore/common/raftserver/server.go	
if len(rd.Entries) > 0 {
		err := s.store.SaveEntries(rd.Entries)
		if err != nil {
			log.Panicf("save raft entries error: %v", err)
		}
	}

正常情况下,依靠WAL日志恢复节点数据问题不大,但随着CM写入事件积累,WAL文件逐步变大,每次启动需从完整的WAL日志恢复集群状态,耗时会增加,可用性得不到保障。

图片

图3. WAL+DB

为了降低系统对WAL日志依赖,如图3所示,CM采用了WAL+DB(基于rocksdb实现存储接口)的形式实现整个集群状态和数据的记录。优势在于,当t时刻的集群状态完整保留在DB后,t+1时刻便等价于0时刻,因而t时刻以前的WAL数据便可以选择性保留以减轻存储介质的压力。

图片

图4. WAL+Mem+Snapshot

此外,如图4.所示,etcd采用内存+快照+WAL日志的方式,正常情况下,由于是全内存操作,性能方面会有一定优势,但相应存在不足。集群规模庞大时内存压力较大,而内存介质相比于存储介质普遍更为昂贵,加重了存储成本;其二,基于定期快照机制来规避集群节点宕机带来的数据丢失风险,由于时间粒度一般是分钟级别,可靠性不能得到保障,因而不能脱离完整的WAL日志。相比于DB方式就更为优雅,虽然会损耗部分性能,但这部分可以通过提升存储介质质量以及优化rocksdb参数,可以降低损耗到忽略不计,并且能够在可靠性和成本得到了充分的补偿。

ReadIndex

ReadIndex省掉了同步 log的开销,能够大幅提升读的吞吐,一定程度上降低读的时延。其大致流程为:

  1. 主节点在收到客户端读请求时,记录下当前的 commitIndex ,称之为 readIndex
  2. 主节点向所有从节点发起一次心跳(确保当前主节点是否有效,避免网络分区时少数派 leader 仍处理请求)
  3. 等待状态机至少应用到 read index (即 apply index >= read index )
  4. 执行读请求,将状态机中的结果返回给客户端 由于在该读请求发起时, Leader 将 commitIndex 记录了下来,只要使客户端读到的内容在该 commitIndex 之后,那么结果一定都满足线性一致。

在具体实践中,为了提升读的性能,raftserver将一段时间内(一次读请求的时间区间)的读请求进行聚合,每批次读请求读取一个返回结果即可。

Snapshot+WAL

在新节点加入或老节点重启恢复时,全量发送日志数据占用网络带宽的代价太高,其次日志长期累积对磁盘空间的消耗也会影响服务的可用性。快照是最简单的方法来优化这一问题。如图5所示,在快照中,当前的整个系统状态被写入一个快照在持久化存储中,然后将整个在这个点之前的日志废弃。

图片

图5. 快照示意图

快照作用流程如下:

  1. 在主节点应用完 tk-1 后,将当前StateMatchine的数据连同Entry 信息生成快照数据;
  2. 当节点恢复或新节点加入,从快照中恢复State Matchine,等价于 Raft 已经应用tk-1为止的所有Entries,效率明显提高;
  3. 将 ApplyIndex置为tk-1,之后从tk继续应用日志。 快照的优点:
  • 降低节点加入或恢复耗时:通过 snapshot +raft log 恢复,无需从第一条 entry 开始;
  • 节省空间:快照做完后即可删除快照点之前的Raft日志。

集群维护实践

元数据容灾

容灾是一个分布式系统必要特性,因不可抗力(如火灾、地震、城市供电中断等)导致某个节点不可用,但对于整个节点仍然可以继续提供服务。因而在Blobstore中,元数据集群一般采用3节点的形式,其中任意一个节点宕机,仍可继续提供服务,保证服务的可用性。

learner节点

learner节点一般和follower节点的数据保持一致,当follower节点发生故障时,可以快速地将learner节点升级为follower节点,相比于重新部署一个新的节点,将集群恢复至健康水平所需的时间更短。

升级建议

  1. 元数据集群尽可能是兼容性升级;
  2. 升级前可先备份一份数据;
  3. 优先升级从节点,从节点正常后再升级主节点,切主命令可参考官网运维文档open in new window;
  4. 升级主节点可先将主切走,升级完成后可选择切回。

故障恢复

当集群因故障不可用时,若集群中存在健康节点(或者learner节点),可基于此节点的数据恢复集群;若没有健康节点(learner节点),可根据备份数据将集群恢复至备份时的状态。

总结

本文主要结合工程实践介绍cubefs纠删码子系统中设计的分布式元数据管理系统,从总体架构、一致性设计以及系统维护方面作简要的介绍。

参考

[1] Ongaro D, Ousterhout J. In search of an understandable consensus algorithm[C]//2014 USENIX Annual Technical Conference (Usenix ATC 14). 2014: 305-319.

[2] Ongaro D, Ousterhout J. In search of an understandable consensus algorithm (extended version)[J]. 2013.

[3] https://github.com/etcd-io/etcd/tree/main/raftopen in new window

[4] https://raft.github.io/open in new window

[5] https://youjiali1995.github.io/raft/basic/open in new window

作者介绍

Zongchao Hu,CubeFS Contributor之一,负责CubeFS纠删码存储引擎的设计研发