一. golang 操作 etcd 基础示例

  1. 当前使用"go.etcd.io/etcd/clientv3" 库,进行演示
import (
	"context"
	"fmt"
	"log"
	"time"
	"go.etcd.io/etcd/clientv3"
)

func main() {
	//1.连接etcd,拿到连接句柄
	cli, err := clientv3.New(clientv3.Config{
		//etcd地址
		Endpoints: []string{"127.0.0.1:2379"},
		//连接超时时间
		DialTimeout: 5 * time.Second,
	})
	if err != nil {
		fmt.Printf("connect to etcd failed, err:%v\n", err)
		return
	}

	//最后不要忘记关闭连接
	defer cli.Close()

	//创建带取消的context
	ctx, cancel := context.WithTimeout(context.Background(), time.Second)

	//2.put 添加数据
	_, err = cli.Put(ctx, "q1mi", "dsb")

	//context的取消函数执行
	cancel()
	if err != nil {
		fmt.Printf("put to etcd failed, err:%v\n", err)
		return
	}

	ctx, cancel = context.WithTimeout(context.Background(), time.Second)
	//3.获取数据
	resp, err := cli.Get(ctx, "q1mi")
	cancel()
	if err != nil {
		fmt.Printf("get from etcd failed, err:%v\n", err)
		return
	}
	for _, ev := range resp.Kvs {
		fmt.Printf("%s:%s\n", ev.Key, ev.Value)
	}

	//4.watch用来获取未来更改的通知
	//当执此处代码执行完毕后,此时程序就会等待etcd中q1mi这个key的变化
	//后续对该key修改、删除、设置时
	//for循环中会打印例如:
	//	Type: PUT Key:q1mi Value:dsb2
	//	Type: DELETE Key:q1mi Value:
	//	Type: PUT Key:q1mi Value:dsb3
	rch := cli.Watch(context.Background(), "q1mi") // <-chan WatchResponse
	for wresp := range rch {
		for _, ev := range wresp.Events {
			fmt.Printf("Type: %s Key:%s Value:%s\n", ev.Type, ev.Kv.Key, ev.Kv.Value)
		}
	}

	//5. lease租约
	//创建一个5秒的租约
	resp1, err := cli.Grant(context.TODO(), 5)
	if err != nil {
		log.Fatal(err)
	}

	//拿到租约后,将租约绑定到/nazha/ 这个key
	//如果没有续约,5秒后这个key就会被移除
	_, err = cli.Put(context.TODO(), "/nazha/", "dsb", clientv3.WithLease(resp1.ID))
	if err != nil {
		log.Fatal(err)
	}

	//6.KeepAlive
	ch, kaerr := cli.KeepAlive(context.TODO(), resp1.ID)
	if kaerr != nil {
		log.Fatal(kaerr)
	}
	for {
		ka := <-ch
		fmt.Println("ttl:", ka.TTL)
	}
}

二. 基于etcd 实现分布式锁

1.使用示例

  1. github.com/coreos/etcd/clientv3/concurrency 库中对基于etcd提供了分布式锁功能
import (
	"context"
	"fmt"
	"github.com/coreos/etcd/clientv3"
	"github.com/coreos/etcd/clientv3/concurrency"
	"log"
	"time"
)

func main() {
	//1.连接etcd
	cli, err := clientv3.New(clientv3.Config{
		Endpoints:   []string{"127.0.0.1:2379"},
		DialTimeout: time.Second * 5,
	})
	if err != nil {
		log.Fatal(err)
	}
	defer cli.Close()
	//2.开启一个会话
	s1, err := concurrency.NewSession(cli, concurrency.WithTTL(5))
	if err != nil {
		log.Fatal(err)
	}
	defer s1.Close()

	//3.执行concurrency.NewMutex()基于会话1上锁成功,返回Mutex,
	m1 := concurrency.NewMutex(s1, "mylock")
	//m1上锁
	if err := m1.Lock(context.TODO()); err != nil {
		log.Fatal(err)
	}
	fmt.Printf("session1 上锁成功。 time:%d \n", time.Now().Unix())

	g2 := make(chan struct{})
	go func() {
		defer close(g2)
		//会话2
		s2, err := concurrency.NewSession(cli, concurrency.WithTTL(5))
		if err != nil {
			log.Fatal(err)
		}
		defer s2.Close()
		//基于会话2获取到Mutex
		m2 := concurrency.NewMutex(s2, "mylock")
		//m2上锁
		if err := m2.Lock(context.TODO()); err != nil {
			log.Fatal(err)
		}
		fmt.Printf("session2 上锁成功。 time:%d \n", time.Now().Unix())
		//释放锁
		if err := m2.Unlock(context.TODO()); err != nil {
			log.Fatal(err)
		}
		fmt.Printf("session2 解锁。 time:%d \n", time.Now().Unix())

	}()

	time.Sleep(5 * time.Second)
	//释放锁
	if err := m1.Unlock(context.TODO()); err != nil {
		log.Fatal(err)
	}
	fmt.Printf("session1 解锁。 time:%d \n", time.Now().Unix())
	<-g2
}

2. 加锁底层

  1. Lock方法中:
  1. 调用tryAcquire
  2. 如果已经加锁成功,或者已经加过锁(可重入),则直接返回
  3. 调用waitDeletes方法,等待所有小于当前Revsion的Key删除
  1. tryAcquire(),是加锁核心
func (m *Mutex) Lock(ctx context.Context) error {
	resp, err := m.tryAcquire(ctx)
	if err != nil {
		return err
	}
	// if no key on prefix / the minimum rev is key, already hold the lock
	ownerKey := resp.Responses[1].GetResponseRange().Kvs
	if len(ownerKey) == 0 || ownerKey[0].CreateRevision == m.myRev {
		m.hdr = resp.Header
		return nil
	}
	client := m.s.Client()
	_, werr := waitDeletes(ctx, client, m.pfx, m.myRev-1)
	// release lock key if wait failed
	if werr != nil {
		m.Unlock(client.Ctx())
		return werr
	}

	// make sure the session is not expired, and the owner key still exists.
	gresp, werr := client.Get(ctx, m.myKey)

	return nil
}

  1. tryAcquire()方法中通过事务来执行加锁逻辑:
  1. 判断当前Key是否为空,即代码中Revision为0
  2. 如果为空,使用Put设置并附加Lease
  3. 如果不为空,获取当前锁的所有者,即最先加锁的对象,避免惊群效应
func (m *Mutex) tryAcquire(ctx context.Context) (*v3.TxnResponse, error) {
	s := m.s
	client := m.s.Client()

	m.myKey = fmt.Sprintf("%s%x", m.pfx, s.Lease())
	cmp := v3.Compare(v3.CreateRevision(m.myKey), "=", 0)
	// put self in lock waiters via myKey; oldest waiter holds lock
	put := v3.OpPut(m.myKey, "", v3.WithLease(s.Lease()))
	// reuse key in case this session already holds the lock
	get := v3.OpGet(m.myKey)
	// fetch current holder to complete uncontended path with only one RPC
	getOwner := v3.OpGet(m.pfx, v3.WithFirstCreate()...)
	resp, err := client.Txn(ctx).If(cmp).Then(put, getOwner).Else(get, getOwner).Commit()
	if err != nil {
		return nil, err
	}
	m.myRev = resp.Header.Revision
	if !resp.Succeeded {
		m.myRev = resp.Responses[0].GetResponseRange().Kvs[0].CreateRevision
	}
	return resp, nil
}

三. 基于etcd 实现注册中心

服务注册

import (
	"context"
	"fmt"
	//没有获取到这两个包,不知道什么原因
	//"go.etcd.io/etcd/client/v3"
	//"go.etcd.io/etcd/client/v3/naming/endpoints"
	"log"
	"time"
)

//Etcd服务注册上下文
type EtcdRegister struct {
	cli       *clientv3.Client
	em        endpoints.Manager
	srvName   string
	srvAddr   string
	ttl       int64
	leaseID   clientv3.LeaseID
	leaseChan <-chan *clientv3.LeaseKeepAliveResponse
}

//1.初始化
func New(etcdAddr []string) (*EtcdRegister, error) {
	log.Printf("开始连接etcd注册中心... \n")
	cli, err := clientv3.New(clientv3.Config{
		Endpoints:   etcdAddr,
		DialTimeout: 5 * time.Second,
	})
	if err != nil {
		log.Printf("连接etcd注册中心失败,失败原因是: %v \n", err)
		return nil, err
	}
	log.Printf("连接etcd注册中心成功! \n")
	return &EtcdRegister{
		cli: cli,
	}, nil
}

//2.注册服务
func (etcdRegister *EtcdRegister) Register(srvName string, srvAddr string, ttl int64) error {
	etcdRegister.srvName = srvName
	etcdRegister.srvAddr = srvAddr
	log.Printf("开始创建etcd端点管理器... \n")
	//1.获取etcd连接
	em, err := endpoints.NewManager(etcdRegister.cli, srvName)
	if err != nil {
		log.Printf("创建etcd端点管理器失败,失败原因是: %v \n", err)
		return err
	}
	etcdRegister.em = em
	log.Printf("创建etcd端点管理器成功! \n")
	log.Printf("开始创建服务租期... \n")
	//2.创建一个租约(指定续约时间)
	lease, err := etcdRegister.cli.Grant(context.TODO(), ttl)
	if err != nil {
		log.Printf("创建服务租期失败,失败原因是: %v", err)
	}
	etcdRegister.ttl = ttl
	etcdRegister.leaseID = lease.ID
	log.Printf("创建服务租期成功! \n")
	log.Printf("开始注册服务,服务名: %s,服务地址: %s \n", srvName, srvAddr)
	//3.注册服务,并绑定租约
	em.AddEndpoint(context.TODO(), fmt.Sprintf("%v/%v", srvName, srvAddr), endpoints.Endpoint{Addr: srvAddr}, clientv3.WithLease(lease.ID))
	if err != nil {
		log.Printf("注册服务失败! \n")
		return err
	}
	log.Printf("注册服务成功! \n")
	//4.续约
	log.Printf("开始服务定期续租! \n")
	leaseChan, err := etcdRegister.cli.KeepAlive(context.TODO(), etcdRegister.leaseID)
	if err != nil {
		log.Printf("服务定期续租失败! \n")
		return err
	}
	etcdRegister.leaseChan = leaseChan
	log.Printf("服务定期续租成功! \n")
	return nil
}

func (etcdRegister *EtcdRegister) UnRegister() error {
	//注册失败返回
	err := etcdRegister.em.DeleteEndpoint(context.TODO(), fmt.Sprintf("%v/%v", etcdRegister.srvName, etcdRegister.srvAddr))
	if err != nil {
		log.Printf("注销服务失败,失败原因是: %v \n", err)
	}
	return nil
}

服务发现

  1. 思路
  1. 服务消费方提供时, 根据持有的服务提供方名称,在etcd上查询,如果查询不到启动报错
  2. 如果查询成功,将获取到的服务调用地址列表保存到到本地,后续在本地获取指定服务的调用列表,负载均衡选择其中一个地址发起调用
  3. 防止服务提供方下线更新, 服务消费方启动一个定时任务,定时通过etcd查询更新服务地址,或者使用watch机制,监听服务提供方在etcd上存储的key,如果地址更新时,会触发监听,同步更新本地
Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐