< Back

Application practice of Raft in CubeFS erasure code system

2023-12-22Zongchao Hu

Introduction

Raft algorithm is a distributed consensus algorithm used to maintain consistent state in distributed systems. It ensures that all nodes in the system reach a consistent state through an election mechanism and log replication. Compared to other consensus algorithms such as Paxos, Raft algorithm is easier to understand and implement, and therefore has been widely used in distributed systems.

The given text describes the CubeFS erasure coding service’s metadata management module (Blobstore/ClusterMgr), Raft algorithm practice, and daily operation recommendations.

Background

CubeFS’s erasure coding storage subsystem (BlobStore) is an independent Blob storage system that is highly reliable, highly available, low-cost, and supports EB scale. Through the previous series of articles, you have gained a deeper understanding of erasure coding. This article will continue to explore more internal implementation details, focusing on how to implement strong consistency metadata management services based on Raft.

图片

Figure1 . Blobstore architecture

As shown in Figure 1, the metadata center (Cluster Manager, CM) plays a crucial role in BlobStore, which is used for managing cluster resources (such as disks, nodes, and storage space units), service registration and discovery, background task storage, and cluster configuration management. CM adopts a distributed architecture and sharding design, providing horizontal scalability and theoretically unlimited performance. With thousands of clients accessing simultaneously, metadata operations will not become a bottleneck.

Model Design

The CM design fully integrates the layered design concept, which mainly consists of service interfaces, state machines, data protocols, persistence, and network layers. The overall architecture can be referred to in Figure 2, and the specific functions of each layer are described as follows:

Service layer: Provides CM’s interface services to the outside world.

State machine: Includes six modules: VolumeMgr, DiskMgr, ScopeMgr, ServiceMgr, ConfigMgr, and KvMgr.

  • VolumeMgr: Volume management module, responsible for creating, renewing, updating volumes, and mapping them to disks.
  • DiskMgr: Responsible for disk and space management, including space chunk allocation, disk registration, and heartbeat updates of disk information.
  • ScopeMgr: Manages the allocation of scope segments such as bid, vid, and diskid.
  • ServierMgr: A general-purpose service registration and discovery management module, responsible for registering and discovering services such as allocator and mqproxy.
  • ConfigMgr: Accesses cluster configuration, task switches, etc., using KvMgr as the storage base.
  • KvMgr: A general-purpose configuration K-V access management module that stores background tasks, configuration information, etc. Data protocol layer: Based on the Raft protocol to ensure data consistency, using RPC for data synchronization communication.

Persistence layer: Uses Rocksdb to persist data, including four basic DBs according to data types: normal, volume, raft, and kv.

Network layer: Responsible for data synchronization, message transmission, etc.

图片

Figure2. CM architecture

Component Introduction

etcd/raft

  1. etcd/raft is a pure state machine model, and all events are processed in the form of messages. It mainly has two structures:
    • Raft: State machine, triggers state transitions based on messages.
    • Node: Represents a node in the raft cluster, providing an interface for interaction between the upper layer and the raft state machine.
  2. The Raft structure maintains the current node’s Raft state, such as Term, Vote, Lead, etc., which is not thread-safe and provides methods to trigger state changes. The main methods are as follows:
    • tick: Each call increases the internal logical clock, used to trigger timeout events such as leader election and heartbeat.
    • step: Type is type stepFunc func(r *raft, m pb.Message) error, used to process messages.
    • There are other methods for handling membership changes, obtaining status, etc.
  3. The application layer handles Storage and Transportation. The decoupling methods of etcd/raft for these two modules are not very similar. Since Storage and Raft are closely related, etcd/raft provides a Storage interface, which is implemented by the user and then set to Config at startup. The Storage interface only involves query operations, and when and how to persist is determined by the user. The design of Transportation is completely customized by the user.
// 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() 
}

RaftServer

RaftServer encapsulates etcd/raft and implements WAL log writing and storage, snapshot mechanism. RaftServer provides upper-layer applications (such as ClusterMgr) with capabilities such as upload, read, member management, WAL log pruning, and raft state query.

type RaftServer interface {
	Stop() 
	Propose(ctx context.Context, data []byte) error
	ReadIndex(ctx context.Context) error
	TransferLeadership(ctx context.Context, leader, transferee uint64)
	AddMember(ctx context.Context, member Member) error
	RemoveMember(ctx context.Context, nodeID uint64) error
	IsLeader() bool 
	Status() Status 

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

WAL

WAL log (Write-Ahead Log) abstracts each state update as a command and appends it to a log, which is written sequentially to ensure IO performance. In Blobstore, the wal module implements the persistent interface reserved by etcd/raft, providing log storage, query, and pruning functions.

type Storage interface {
	InitialState() (pb.HardState, pb.ConfState, error) 
	Entries(lo, hi, maxSize uint64) ([]pb.Entry, error) 
	Term(i uint64) (uint64, error) 
	LastIndex() (uint64, error) 
	FirstIndex() (uint64, error)
	Snapshot() (pb.Snapshot, error)
}
// Storage the storage
type Wal struct {
	// Log Entry
	logfiles    []logName 
	last        *logFile
	cache       *logFileCache 
	hs          pb.HardState
	st          Snapshot 
	mt          *meta      
}

Consistency

When a disk registration process is initiated, CM ensures the consistency of disk information across the cluster. When CM receives a request from blobnode (the storage unit that manages physical disks) to add a disk, it first checks if it is the leader (if not, it forwards the request to the leader). After confirming that the disk meets the conditions for adding a new disk to the cluster, CM generates a structure information as a Raft proposal through StateMachine. This structure information includes the name of the module involved in the request operation, the specific operation type within the module, specific data, and other information.

ModuleNameOpTypeDiskInfoOthers

Specifically, in the code

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

StateMachine submits the proposal to raftServer, which is further packaged into a message that raft can recognize.

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

Raft forwards this message to other raft nodes based on RPC through the raft port and returns the message to StateMachine after confirming that most nodes have completed synchronization. According to the state machine safety of Raft, it is known that the instructions executed with the same log index must be the same, although each node may execute instructions at different times, and the final state machine result is consistent. The master node issues a log synchronization request, the slave node receives and synchronizes the log, and each CM node’s state machine executes instructions according to the log order, ultimately ensuring the consistency of the entire cluster state.

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

Specifically, StateMachine first distributes the message to the specific module based on the module name in the message. After receiving the message, the module changes the memory state and persists it to the DB according to the operation type

// 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
  	 }
  	 ...

After the entire process is completed, the blocked request is notified to return success or specific error information to blobnode, completing the entire registration and disk addition process. The above is from the perspective of the master node. Switching to the perspective of the slave node, the slave node has already received the packaged raft message, so it can directly enter the raft state machine rotation until the master node tells that this message can be applied to StateMachine. Like the master node, the slave node changes its memory state and persists it to the 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
}

Availability

Previously, we discussed how CM ensures consistency and briefly demonstrated how the cluster ensures consistency through disk registration information. How is availability achieved?

WAL+DB

In the CM cluster, each write operation goes through WAL first and then persists to the database. WAL records every change operation in the metadata cluster, fully ensuring data security, and applies to accurately recover to the state before the node crashes.

// 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)
		}
	}

Under normal circumstances, there is no problem recovering node data using WAL logs. However, as CM write events accumulate, the WAL file gradually increases in size. Each time the system starts, it needs to recover the cluster status from the complete WAL log, which increases the time required and reduces availability.

图片

Figure 3. WAL+DB

To reduce the system’s dependence on WAL logs, CM uses WAL+DB (storage interface based on rocksdb) to record the entire cluster state and data, as shown in Figure 3. The advantage is that when the cluster state at time t is completely preserved in the DB, time t+1 is equivalent to time 0, so WAL data before time t can be selectively retained to reduce the pressure on storage media.

图片

Figure4. WAL+Mem+Snapshot

In addition, as shown in Figure 4, etcd uses memory+snapshot+WAL log mode. Under normal circumstances, there is a certain advantage in performance because it is a full-memory operation. However, there are corresponding shortcomings. When the cluster is large, there is a lot of memory pressure, and memory media is generally more expensive than storage media, which increases storage costs. Secondly, based on the periodic snapshot mechanism to avoid the risk of data loss caused by cluster node crashes, the reliability cannot be guaranteed due to the minute-level time granularity, so it cannot be separated from the complete WAL log. Compared with the DB method, it is more elegant, although it will consume some performance, but this part can be reduced to negligible by improving the quality of storage media and optimizing rocksdb parameters, and can be fully compensated in terms of reliability and cost.

ReadIndex

ReadIndex eliminates the overhead of synchronizing logs, which can greatly improve the throughput of reads and reduce the latency of reads to some extent. The general process is as follows:

  1. When the primary node receives a read request from the client, it records the current commitIndex as readIndex.
  2. The primary node initiates a heartbeat to all secondary nodes (to ensure that the current primary node is valid and to avoid the minority leader still processing requests during network partition).
  3. Wait for the state machine to apply to at least the read index (that is, apply index >= read index).
  4. Execute the read request and return the result from the state machine to the client. Since the Leader records commitIndex when the read request is initiated, as long as the content read by the client is after that commitIndex, the results will definitely satisfy linear consistency. In practice, to improve the performance of reads, the raftserver aggregates read requests over a period of time (the time interval of a read request), and each batch of read requests reads one result and returns it.

Snapshot+WAL

When a new node joins or an old node restarts and recovers, the cost of sending full log data over the network is too high, and the long-term accumulation of logs consumes disk space, which also affects service availability. Snapshots are the simplest way to optimize this problem. As shown in Figure 5, in the snapshot, the current entire system state is written into a snapshot in persistent storage, and then all logs before this point are discarded.

图片

Figure5 . Snapshot schematic diagram

The process of snapshot is as follows:

  1. After the primary node applies tk-1, the current StateMatchine data and Entry information are used to generate snapshot data.
  2. When a node recovers or a new node joins, the State Matchine is restored from the snapshot, which is equivalent to all Entries that Raft has applied up to tk-1, and the efficiency is significantly improved.
  3. Set ApplyIndex to tk-1, and then continue to apply logs from tk. Advantages of snapshots:
  • Reduce the time required for node joining or recovery: restore through snapshot +raft log, no need to start from the first entry.
  • Save space: After the snapshot is completed, the Raft logs before the snapshot point can be deleted.

Cluster Maintenance Practice

Metadata Disaster Recovery

Disaster recovery is a necessary feature of distributed systems. When an uncontrollable event (such as a fire, earthquake, or city power outage) causes a node to become unavailable, the entire node can still continue to provide services. Therefore, in Blobstore, the metadata cluster generally adopts a 3-node form, where any node can still provide services if one node fails, ensuring service availability.

Learner Node

Learner nodes generally keep data consistent with follower nodes. When a follower node fails, the learner node can be quickly upgraded to a follower node, which takes less time to restore the cluster to a healthy level than redeploying a new node.

Upgrade Suggestions

  1. The metadata cluster should be upgraded for compatibility as much as possible.
  2. Before upgrading, you can back up the data first.
  3. Upgrade the slave node first, then upgrade the master node after the slave node is normal. The command for switching the master can refer to the official website operation and maintenance documentopen in new window
  4. You can switch the master first when upgrading the master node, and then switch back after the upgrade is complete.

Fault Recovery

When the cluster is unavailable due to a fault, if there are healthy nodes (or learner nodes) in the cluster, the cluster can be restored based on the data of this node. If there are no healthy nodes (learner nodes), the cluster can be restored to the state at the time of backup based on the backup data.

Summary

This article mainly introduces the distributed metadata management system designed in the cubefs erasure code subsystem based on engineering practice, and briefly introduces the overall architecture, consistency design, and system maintenance.

Reference

[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

Introduction of the author

Zongchao Hu, one of the CubeFS Contributors, is responsible for the design and development of the CubeFS Erasurectifying code storage engine