go 进阶 注册中心: 二. etcd 与 golang
·
一. golang 操作 etcd 基础示例
- 当前使用"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.使用示例
- 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. 加锁底层
- Lock方法中:
- 调用tryAcquire
- 如果已经加锁成功,或者已经加过锁(可重入),则直接返回
- 调用waitDeletes方法,等待所有小于当前Revsion的Key删除
- 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
}
- tryAcquire()方法中通过事务来执行加锁逻辑:
- 判断当前Key是否为空,即代码中Revision为0
- 如果为空,使用Put设置并附加Lease
- 如果不为空,获取当前锁的所有者,即最先加锁的对象,避免惊群效应
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
}
服务发现
- 思路
- 服务消费方提供时, 根据持有的服务提供方名称,在etcd上查询,如果查询不到启动报错
- 如果查询成功,将获取到的服务调用地址列表保存到到本地,后续在本地获取指定服务的调用列表,负载均衡选择其中一个地址发起调用
- 防止服务提供方下线更新, 服务消费方启动一个定时任务,定时通过etcd查询更新服务地址,或者使用watch机制,监听服务提供方在etcd上存储的key,如果地址更新时,会触发监听,同步更新本地
更多推荐
所有评论(0)