您好,登錄后才能下訂單哦!
今天就跟大家聊聊有關golang etcd raft協議是怎樣的,可能很多人都不太了解,為了讓大家更加了解,小編給大家總結了以下內容,希望大家根據這篇文章可以有所收獲。
分布式存儲系統通常會通過維護多個副本來進行容錯, 以提高系統的可用性。 這就引出了分布式存儲系統的核心問題——如何保證多個副本的一致性? Raft算法把問題分解成了四個子問題: 1. 領袖選舉(leader election)、 2. 日志復制(log replication)、 3. 安全性(safety) 4. 成員關系變化(membership changes) 這幾個子問題。 源碼gitee地址: https://gitee.com/ioly/learning.gooop
根據raft協議,實現高可用分布式強一致的kv存儲
雖然Leader State還有細節沒處理完,但應該能啟動并提供基本服務了
添加外圍功能,為首次“點火”做準備:
config/tRaftConfig:從本地json文件讀取集群節點配置,提供IRaftConfig/IRaftNodeConfig的實現
lsm/tRaftLSMImplement: 提供對頂層接口IRaftLSM的實現,將“配置/kv存儲/節點通訊”三大塊粘合起來
server/IRaftKVServer:server啟動器接口
server/tRaftKVServer: server啟動器的實現,監聽raft rpc和kv rpc
從本地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 }
提供對頂層接口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啟動器接口
package server type IRaftKVServer interface { BeginServeTCP(port int) error }
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 }
(未完待續)
看完上述內容,你們對golang etcd raft協議是怎樣的有進一步的了解嗎?如果還想了解更多知識或者相關內容,請關注億速云行業資訊頻道,感謝大家的支持。
免責聲明:本站發布的內容(圖片、視頻和文字)以原創、轉載和分享為主,文章觀點不代表本網站立場,如果涉及侵權請聯系站長郵箱:is@yisu.com進行舉報,并提供相關證據,一經查實,將立刻刪除涉嫌侵權內容。