手撸golangetcdraft协议之11

ioly      2022-02-12     295

关键词:

手撸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
}

(未完待续)

手撸golanggo与微服务chatserver之1

缘起最近阅读<<Go微服务实战>>(刘金亮,2021.1)本系列笔记拟采用golang练习之案例需求(聊天服务器)用户可以连接到服务器。用户可以设定自己的用户名。用户可以向服务器发送消息,同时服务器也会向其他用户广播该消息... 查看详情

手撸golangspringioc/aop之2

手撸golangspringioc/aop之2缘起最近阅读[Offer来了:Java面试核心知识点精讲(框架篇)](王磊,2020.6)本系列笔记拟采用golang练习之Talkischeap,showmethecode.SpringSpring基于J2EE技术实现了一套轻量级的JavaWebService系统应用框架。它有很多优秀的... 查看详情

手撸golang仿springioc/aop之4蓝图

手撸golang仿springioc/aop之4蓝图缘起最近阅读[SpringBoot技术内幕:架构设计与实现原理](朱智胜,2020.6)本系列笔记拟采用golang练习之Talkischeap,showmethecode.SpringSpring的主要特性:1.控制反转(InversionofControl,IoC)2.面向容器3.面向切面(Aspect... 查看详情

手撸golang仿springioc/aop之12增强3

手撸golang仿springioc/aop之12增强3缘起最近阅读[SpringBoot技术内幕:架构设计与实现原理](朱智胜,2020.6)本系列笔记拟采用golang练习之Talkischeap,showmethecode.SpringSpring的主要特性:1.控制反转(InversionofControl,IoC)2.面向容器3.面向切面(Aspe... 查看详情

手撸golang仿springioc/aop之5如何扫描

手撸golang仿springioc/aop之5如何扫描缘起最近阅读[SpringBoot技术内幕:架构设计与实现原理](朱智胜,2020.6)本系列笔记拟采用golang练习之Talkischeap,showmethecode.SpringSpring的主要特性:1.控制反转(InversionofControl,IoC)2.面向容器3.面向切面(... 查看详情

手撸golangspringioc/aop之2

参考技术A手撸golangspringioc/aop之2最近阅读[Offer来了:Java面试核心知识点精讲(框架篇)](王磊,2020.6)本系列笔记拟采用golang练习之Talkischeap,showmethecode.配置接口指令接口指令构建器接口指令执行上下文接口保存配置另存配置添加... 查看详情

手撸golanggo与微服务saga模式之1

缘起最近阅读<<Go微服务实战>>(刘金亮,2021.1)本系列笔记拟采用golang练习之Saga模式saga模式将分布式长事务切分为一系列独立短事务每个短事务是可通过补偿动作进行撤销的事务动作和补偿动作都是幂等的,允许重复执行而... 查看详情

手撸golanggo与微服务saga模式之7

缘起最近阅读<<Go微服务实战>>(刘金亮,2021.1)本系列笔记拟采用golang练习之Saga模式saga模式将分布式长事务切分为一系列独立短事务每个短事务是可通过补偿动作进行撤销的事务动作和补动作偿都是幂等的,允许重复执行而... 查看详情

手撸系列之——orm(对象关系映射)(代码片段)

ORM:对象关系映射类》》》数据库的一张表对象》》》表的一条记录对象点属性》》》记录某一个字段对应的值废话不多少,先上代码:#orm.pyfrommysql_singletionimportMysql#设置表字段类,通常需要的属性为字段名,字段类型,是否为... 查看详情

手撸golanggo与微服务chatserver之3压测与诊断

缘起最近阅读<<Go微服务实战>>(刘金亮,2021.1)本系列笔记拟采用golang练习之案例需求(聊天服务器)用户可以连接到服务器。用户可以设定自己的用户名。用户可以向服务器发送消息,同时服务器也会向其他用户广播该消息... 查看详情

goroutine并发调度模型深度解析之手撸一个协程池(代码片段)

golanggoroutine协程池GroutinePool高并发并发(并行),一直以来都是一个编程语言里的核心主题之一,也是被开发者关注最多的话题;Go语言作为一个出道以来就自带『高并发』光环的富二代编程语言,它的并发(并行)编程肯定是... 查看详情

手撸orm

本文目录ORM简介Python中常用ORM框架 原生操作数据库模块pymysqlORM框架之SQLAlchemy手把手带你写一个自己的ORM框架回到目录ORM简介ORM即ObjectRelationalMapping,全称对象关系映射。当我们需要对数据库进行操作时,势必需要通过连接... 查看详情

2021-11-27wpf上位机100-西门子s7协议之modbus读取数据

文章目录前言一、西门子S7协议之modbus读取数据二、使用步骤1.modbus读取数据代码2.读取数据总结前言随着人工智能的不断发展,物联网这门技术也越来越重要,很多人都开启了物联网学习,本文就介绍了物联网的S7报文协议。提... 查看详情

2021-11-27wpf上位机101-西门子s7协议之s7.net

文章目录前言一、西门子S7协议之S7.NET读取数据二、使用步骤1.S7.NET2.读取数据总结前言随着人工智能的不断发展,物联网这门技术也越来越重要,很多人都开启了物联网学习,本文就介绍了物联网的S7报文协议。提示:以下是本... 查看详情

2021-11-27wpf上位机102-西门子s7协议之i区读写封装

文章目录前言一、西门子S7协议之I区读写封装二、使用步骤1.modbus读取数据代码2.读取数据总结前言随着人工智能的不断发展,物联网这门技术也越来越重要,很多人都开启了物联网学习,本文就介绍了物联网的S7报文协议。提示... 查看详情

2021-12-11wpf上位机112-欧姆龙协议之finstcp协议

FinsTCP协议1、Fins是一个公开的协议网口(Fins-》UDPFinsTCP)FinsTCP在Fins的基础上添加一个FinsTCP的HeadFins官方文档:https://www.fa.omron.com.cn/data_pdf/mnu/w342-e1-17_cs1_cj1_cp1_com_cmd.pdf?id=16382、欧姆龙常用协议关系Hostlink(C-Mode(串口)、Fins(... 查看详情

使用javasocket手撸一个http服务器(代码片段)

原文连接:使用JavaSocket手撸一个http服务器作为一个java后端,提供http服务可以说是基本技能之一了,但是你真的了解http协议么?你知道知道如何手撸一个http服务器么?tomcat的底层是怎么支持http服务的呢?大名鼎鼎的Servlet又是... 查看详情

2021-12-11wpf上位机111-欧姆龙协议之omronfinstcp.net(代码片段)

OmronFinsTCP.Net的使用OmronFinsTCP.Net.EtherNetPLCetherNetPLC=newOmronFinsTCP.Net.EtherNetPLC();//建立连接//1、TCP三次握手//2、FincTCP建立通信etherNetPLC.Link("192.168.151.132", 查看详情