共计 7106 个字符,预计需要花费 18 分钟才能阅读完成。
手撸 golang etcd raft 协定之 11
缘起
最近浏览 [云原生分布式存储基石:etcd 深刻解析] (杜军 , 2019.1)
本系列笔记拟采纳 golang 练习之
raft 分布式一致性算法
分布式存储系统通常会通过保护多个副原本进行容错,以进步零碎的可用性。这就引出了分布式存储系统的外围问题——如何保障多个正本的一致性?Raft 算法把问题分解成了四个子问题:1. 首领选举(leader election)、2. 日志复制(log replication)、3. 安全性(safety)4. 成员关系变动(membership changes)这几个子问题。源码 gitee 地址:https://gitee.com/ioly/learning.gooop
指标
- 依据 raft 协定,实现高可用分布式强统一的 kv 存储
子目标(Day 11)
- 尽管 Leader State 还有细节没解决完,但应该能启动并提供根本服务了
-
增加外围性能,为首次“点火”做筹备:
- config/tRaftConfig:从本地 json 文件读取集群节点配置,提供 IRaftConfig/IRaftNodeConfig 的实现
- lsm/tRaftLSMImplement:提供对顶层接口 IRaftLSM 的实现,将“配置 /kv 存储 / 节点通信”三大块粘合起来
- server/IRaftKVServer:server 启动器接口
- server/tRaftKVServer: server 启动器的实现,监听 raft rpc 和 kv rpc
config/tRaftConfig.go
从本地 json 文件读取集群节点配置,提供 IRaftConfig/IRaftNodeConfig 的实现
package config | |
import ( | |
"encoding/json" | |
"os" | |
) | |
type tRaftConfig struct { | |
ID string | |
Nodes []*tRaftNodeConfig} | |
type tRaftNodeConfig struct { | |
ID string | |
Endpoint string | |
} | |
func (me *tRaftConfig) GetID() string {return me.ID} | |
func (me *tRaftConfig) GetNodes() []IRaftNodeConfig {a := make([]IRaftNodeConfig, len(me.Nodes)) | |
for i,it := range me.Nodes {a[i] = it | |
} | |
return a | |
} | |
func (me *tRaftNodeConfig) GetID() string {return me.ID} | |
func (me *tRaftNodeConfig) GetEndpoint() string {return me.Endpoint} | |
func LoadJSONFile(file string) IRaftConfig {data, err := os.ReadFile(file) | |
if err != nil {panic(err) | |
} | |
c := new(tRaftConfig) | |
err = json.Unmarshal(data, c) | |
if err != nil {panic(err) | |
} | |
return c | |
} |
lsm/tRaftLSMImplement.go
提供对顶层接口 IRaftLSM 的实现,将“配置 /kv 存储 / 节点通信”三大块粘合起来,并增加诊断日志。
package lsm | |
import ( | |
"learning/gooop/etcd/raft/common" | |
"learning/gooop/etcd/raft/config" | |
"learning/gooop/etcd/raft/logger" | |
"learning/gooop/etcd/raft/rpc" | |
"learning/gooop/etcd/raft/rpc/client" | |
"learning/gooop/etcd/raft/store" | |
"sync" | |
) | |
type tRaftLSMImplement struct { | |
tEventDrivenModel | |
mInitOnce sync.Once | |
mConfig config.IRaftConfig | |
mStore store.ILogStore | |
mClientService client.IRaftClientService | |
mState IRaftState | |
} | |
// trigger: init() | |
// args: empty | |
const meInit = "lsm.Init" | |
// trigger: HandleStateChanged() | |
// args: IRaftState | |
const meStateChanged = "lsm.StateChnaged" | |
func (me *tRaftLSMImplement) init() {me.mInitOnce.Do(func() {me.initEventHandlers() | |
me.raise(meInit) | |
}) | |
} | |
func (me *tRaftLSMImplement) initEventHandlers() { | |
// write only | |
me.hookEventsForConfig() | |
me.hookEventsForStore() | |
me.hookEventsForPeerService() | |
me.hookEventsForState()} | |
func (me *tRaftLSMImplement) hookEventsForConfig() {me.hook(meInit, func(e string, args ...interface{}) {logger.Logf("tRaftLSMImplement.init, ConfigFile = %v", common.ConfigFile) | |
me.mConfig = config.LoadJSONFile(common.ConfigFile) | |
}) | |
} | |
func (me *tRaftLSMImplement) hookEventsForStore() {me.hook(meInit, func(e string, args ...interface{}) {logger.Logf("tRaftLSMImplement.init, DataFile = %v", common.DataFile) | |
err, db := store.NewBoltStore(common.DataFile) | |
if err != nil {panic(err) | |
} | |
me.mStore = db | |
}) | |
} | |
func (me *tRaftLSMImplement) hookEventsForPeerService() {me.hook(meInit, func(e string, args ...interface{}) {me.mClientService = client.NewRaftClientService(me.mConfig) | |
}) | |
} | |
func (me *tRaftLSMImplement) hookEventsForState() {me.hook(meInit, func(e string, args ...interface{}) {me.mState = newFollowerState(me, me.mStore.LastCommittedTerm()) | |
me.mState.Start()}) | |
me.hook(meStateChanged, func(e string, args ...interface{}) {state := args[0].(IRaftState) | |
logger.Logf("tRaftLSMImplement.StateChanged, %v", state.Role()) | |
me.mState = state | |
state.Start()}) | |
} | |
func (me *tRaftLSMImplement) Config() config.IRaftConfig {return me.mConfig} | |
func (me *tRaftLSMImplement) Store() store.ILogStore {return me.mStore} | |
func (me *tRaftLSMImplement) HandleStateChanged(state IRaftState) {me.raise(meStateChanged, state) | |
} | |
func (me *tRaftLSMImplement) RaftClientService() client.IRaftClientService {return me.mClientService} | |
func (me *tRaftLSMImplement) Heartbeat(cmd *rpc.HeartbeatCmd, ret *rpc.HeartbeatRet) error { | |
state := me.mState | |
e := state.Heartbeat(cmd, ret) | |
logger.Logf("tRaftLSMImplement.Heartbeat, state=%v, cmd=%v, ret=%v, err=%v", cmd, ret, e) | |
return e | |
} | |
func (me *tRaftLSMImplement) AppendLog(cmd *rpc.AppendLogCmd, ret *rpc.AppendLogRet) error { | |
state := me.mState | |
e := state.AppendLog(cmd, ret) | |
logger.Logf("tRaftLSMImplement.AppendLog, state=%v, cmd=%v, ret=%v, err=%v", cmd, ret, e) | |
return e | |
} | |
func (me *tRaftLSMImplement) CommitLog(cmd *rpc.CommitLogCmd, ret *rpc.CommitLogRet) error { | |
state := me.mState | |
e := state.CommitLog(cmd, ret) | |
logger.Logf("tRaftLSMImplement.CommitLog, state=%v, cmd=%v, ret=%v, err=%v", cmd, ret, e) | |
return e | |
} | |
func (me *tRaftLSMImplement) RequestVote(cmd *rpc.RequestVoteCmd, ret *rpc.RequestVoteRet) error { | |
state := me.mState | |
e := state.RequestVote(cmd, ret) | |
logger.Logf("tRaftLSMImplement.RequestVote, state=%v, cmd=%v, ret=%v, err=%v", cmd, ret, e) | |
return e | |
} | |
func (me *tRaftLSMImplement) ExecuteKVCmd(cmd *rpc.KVCmd, ret *rpc.KVRet) error { | |
state := me.mState | |
e := state.ExecuteKVCmd(cmd, ret) | |
logger.Logf("tRaftLSMImplement.ExecuteKVCmd, state=%v, cmd=%v, ret=%v, err=%v", cmd, ret, e) | |
return e | |
} | |
func (me *tRaftLSMImplement) State() IRaftState {return me.mState} | |
func NewRaftLSM() IRaftLSM {it := new(tRaftLSMImplement) | |
it.init() | |
return it | |
} |
server/IRaftKVServer.go
server 启动器接口
package server | |
type IRaftKVServer interface {BeginServeTCP(port int) error | |
} |
server/tRaftKVServer.go
server 启动器的实现,监听 raft rpc 和 kv rpc
package server | |
import ( | |
"fmt" | |
"learning/gooop/etcd/raft/lsm" | |
rrpc "learning/gooop/etcd/raft/rpc" | |
"learning/gooop/saga/mqs/logger" | |
"net" | |
"net/rpc" | |
"time" | |
) | |
type tRaftKVServer int | |
func (me *tRaftKVServer) BeginServeTCP(port int) error {logger.Logf("tRaftKVServer.BeginServeTCP, starting, port=%v", port) | |
// resolve address | |
addy, err := net.ResolveTCPAddr("tcp", fmt.Sprintf("0.0.0.0:%d", port)) | |
if err != nil {return err} | |
// create raft lsm singleton | |
raftLSM := lsm.NewRaftLSM() | |
// register raft rpc server | |
rserver := &RaftRPCServer {mRaftLSM : raftLSM,} | |
err = rpc.Register(rserver) | |
if err != nil {return err} | |
// register kv rpc server | |
kserver := &KVStoreRPCServer{mRaftLSM: raftLSM,} | |
err = rpc.Register(kserver) | |
if err != nil {return err} | |
inbound, err := net.ListenTCP("tcp", addy) | |
if err != nil {return err} | |
go rpc.Accept(inbound) | |
logger.Logf("tRaftKVServer.BeginServeTCP, service ready at port=%v", port) | |
return nil | |
} | |
// RaftRPCServer exposes a raft rpc service | |
type RaftRPCServer struct {mRaftLSM lsm.IRaftLSM} | |
// Heartbeat leader to follower | |
func (me *RaftRPCServer) Heartbeat(cmd *rrpc.HeartbeatCmd, ret *rrpc.HeartbeatRet) error {e := me.mRaftLSM.Heartbeat(cmd, ret) | |
logger.Logf("RaftRPCServer.Heartbeat, cmd=%v, ret=%v, e=%v", cmd, ret, e) | |
return e | |
} | |
// AppendLog leader to follower | |
func (me *RaftRPCServer) AppendLog(cmd *rrpc.AppendLogCmd, ret *rrpc.AppendLogRet) error {e := me.mRaftLSM.AppendLog(cmd, ret) | |
logger.Logf("RaftRPCServer.AppendLog, cmd=%v, ret=%v, e=%v", cmd, ret, e) | |
return e | |
} | |
// CommitLog leader to follower | |
func (me *RaftRPCServer) CommitLog(cmd *rrpc.CommitLogCmd, ret *rrpc.CommitLogRet) error {e := me.mRaftLSM.CommitLog(cmd, ret) | |
logger.Logf("RaftRPCServer.CommitLog, cmd=%v, ret=%v, e=%v", cmd, ret, e) | |
return e | |
} | |
// RequestVote candidate to others | |
func (me *RaftRPCServer) RequestVote(cmd *rrpc.RequestVoteCmd, ret *rrpc.RequestVoteRet) error {e := me.mRaftLSM.RequestVote(cmd, ret) | |
logger.Logf("RaftRPCServer.RequestVote, cmd=%v, ret=%v, e=%v", cmd, ret, e) | |
return e | |
} | |
// Ping to keep alive | |
func (me *RaftRPCServer) Ping(cmd *rrpc.PingCmd, ret *rrpc.PingRet) error {ret.SenderID = me.mRaftLSM.Config().GetID() | |
ret.Timestamp = time.Now().UnixNano() | |
logger.Logf("RaftRPCServer.Ping, cmd=%v, ret=%v", cmd, ret) | |
return nil | |
} | |
// KVStoreRPCServer expose a kv storage service | |
type KVStoreRPCServer struct {mRaftLSM lsm.IRaftLSM} | |
// ExecuteKVCmd leader to follower | |
func (me *KVStoreRPCServer) ExecuteKVCmd(cmd *rrpc.KVCmd, ret *rrpc.KVRet) error {e := me.mRaftLSM.ExecuteKVCmd(cmd, ret) | |
logger.Logf("KVStoreRPCServer.ExecuteKVCmd, cmd=%v, ret=%v, e=%v", cmd, ret, e) | |
return e | |
} |
(未完待续)
正文完