本文目录导读:

我将为您提供一个基于Etcd的完整实现案例,展示其核心功能和应用场景。
项目结构
etcd-demo/
├── go.mod
├── main.go
├── config/
│ └── config.go
├── service/
│ ├── registry.go
│ ├── discovery.go
│ └── leader_election.go
├── service/
│ └── distributed_lock.go
└── watch/
└── config_watcher.go
服务注册与发现
1 服务注册
// service/registry.go
package service
import (
"context"
"encoding/json"
"fmt"
"time"
"go.etcd.io/etcd/client/v3"
)
type ServiceInfo struct {
Name string `json:"name"`
Address string `json:"address"`
Port int `json:"port"`
Meta map[string]string `json:"meta"`
TTL int64 `json:"ttl"`
}
type ServiceRegistry struct {
client *clientv3.Client
leaseID clientv3.LeaseID
serviceKey string
stopChan chan struct{}
}
func NewServiceRegistry(etcdEndpoints []string) (*ServiceRegistry, error) {
client, err := clientv3.New(clientv3.Config{
Endpoints: etcdEndpoints,
DialTimeout: 5 * time.Second,
})
if err != nil {
return nil, fmt.Errorf("创建etcd客户端失败: %v", err)
}
return &ServiceRegistry{
client: client,
stopChan: make(chan struct{}),
}, nil
}
// 注册服务
func (r *ServiceRegistry) RegisterService(info ServiceInfo) error {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
// 创建租约
leaseResp, err := r.client.Grant(ctx, info.TTL)
if err != nil {
return fmt.Errorf("创建租约失败: %v", err)
}
r.leaseID = leaseResp.ID
r.serviceKey = fmt.Sprintf("/services/%s/%s:%d",
info.Name, info.Address, info.Port)
// 序列化服务信息
data, err := json.Marshal(info)
if err != nil {
return fmt.Errorf("序列化服务信息失败: %v", err)
}
// 将服务信息写入etcd,绑定租约
_, err = r.client.Put(ctx, r.serviceKey, string(data),
clientv3.WithLease(r.leaseID))
if err != nil {
return fmt.Errorf("写入服务信息失败: %v", err)
}
// 启动租约保活
go r.keepAlive()
fmt.Printf("服务 %s 注册成功,地址: %s:%d\n",
info.Name, info.Address, info.Port)
return nil
}
// 租约续期
func (r *ServiceRegistry) keepAlive() {
keepAliveChan, err := r.client.KeepAlive(
context.Background(), r.leaseID)
if err != nil {
fmt.Printf("租约保活启动失败: %v\n", err)
return
}
for {
select {
case <-r.stopChan:
return
case resp, ok := <-keepAliveChan:
if !ok {
fmt.Println("租约保活通道关闭")
return
}
if resp != nil {
// 续期成功,可以添加日志
}
}
}
}
// 注销服务
func (r *ServiceRegistry) UnregisterService() error {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
// 删除服务注册信息
if _, err := r.client.Delete(ctx, r.serviceKey); err != nil {
return err
}
// 撤销租约
if r.leaseID != 0 {
_, err := r.client.Revoke(ctx, r.leaseID)
if err != nil {
return err
}
}
close(r.stopChan)
fmt.Println("服务注销成功")
return nil
}
2 服务发现
// service/discovery.go
package service
import (
"context"
"encoding/json"
"fmt"
"sync"
"time"
"go.etcd.io/etcd/client/v3"
)
type ServiceDiscovery struct {
client *clientv3.Client
services map[string][]ServiceInfo
watchChan chan string
mutex sync.RWMutex
stopChan chan struct{}
}
func NewServiceDiscovery(etcdEndpoints []string) (*ServiceDiscovery, error) {
client, err := clientv3.New(clientv3.Config{
Endpoints: etcdEndpoints,
DialTimeout: 5 * time.Second,
})
if err != nil {
return nil, fmt.Errorf("创建etcd客户端失败: %v", err)
}
return &ServiceDiscovery{
client: client,
services: make(map[string][]ServiceInfo),
watchChan: make(chan string, 100),
stopChan: make(chan struct{}),
}, nil
}
// 发现服务
func (d *ServiceDiscovery) DiscoverService(serviceName string) ([]ServiceInfo, error) {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
// 获取所有服务实例
resp, err := d.client.Get(ctx, fmt.Sprintf("/services/%s", serviceName),
clientv3.WithPrefix())
if err != nil {
return nil, err
}
var instances []ServiceInfo
for _, kv := range resp.Kvs {
var info ServiceInfo
if err := json.Unmarshal(kv.Value, &info); err != nil {
continue
}
instances = append(instances, info)
}
d.mutex.Lock()
d.services[serviceName] = instances
d.mutex.Unlock()
return instances, nil
}
// 监听服务变化
func (d *ServiceDiscovery) WatchService(serviceName string) {
prefix := fmt.Sprintf("/services/%s", serviceName)
// 监听指定的服务前缀
watchChan := d.client.Watch(context.Background(), prefix,
clientv3.WithPrefix())
for {
select {
case <-d.stopChan:
return
case resp, ok := <-watchChan:
if !ok {
return
}
// 服务注册信息变化时重新发现
for _, ev := range resp.Events {
fmt.Printf("服务 %s 发生变化: %s\n",
serviceName, ev.Type)
// 重新获取服务实例列表
newInstances, err := d.DiscoverService(serviceName)
if err != nil {
fmt.Printf("重新发现服务失败: %v\n", err)
continue
}
// 通知订阅者
d.notifyWatchers(serviceName, newInstances)
}
}
}
}
// 通知服务变更
func (d *ServiceDiscovery) notifyWatchers(serviceName string,
instances []ServiceInfo) {
select {
case d.watchChan <- serviceName:
// 通知成功
default:
// 通道满,丢弃通知
}
fmt.Printf("服务 %s 当前实例数: %d\n", serviceName, len(instances))
}
分布式锁实现
// service/distributed_lock.go
package service
import (
"context"
"fmt"
"time"
"go.etcd.io/etcd/client/v3"
"go.etcd.io/etcd/client/v3/concurrency"
)
type DistributedLock struct {
client *clientv3.Client
session *concurrency.Session
mutex *concurrency.Mutex
lockKey string
}
func NewDistributedLock(etcdEndpoints []string, lockKey string) (*DistributedLock, error) {
client, err := clientv3.New(clientv3.Config{
Endpoints: etcdEndpoints,
DialTimeout: 5 * time.Second,
})
if err != nil {
return nil, fmt.Errorf("创建etcd客户端失败: %v", err)
}
// 创建会话
session, err := concurrency.NewSession(client,
concurrency.WithTTL(10))
if err != nil {
return nil, fmt.Errorf("创建会话失败: %v", err)
}
return &DistributedLock{
client: client,
session: session,
mutex: concurrency.NewMutex(session, lockKey),
lockKey: lockKey,
}, nil
}
// 获取锁
func (l *DistributedLock) TryLock(ctx context.Context, timeout time.Duration) error {
lockCtx, cancel := context.WithTimeout(ctx, timeout)
defer cancel()
// 尝试获取锁
if err := l.mutex.Lock(lockCtx); err != nil {
return fmt.Errorf("获取分布式锁失败: %v", err)
}
fmt.Printf("成功获取分布式锁: %s\n", l.lockKey)
return nil
}
// 释放锁
func (l *DistributedLock) Unlock() error {
if err := l.mutex.Unlock(context.Background()); err != nil {
return fmt.Errorf("释放分布式锁失败: %v", err)
}
// 关闭会话
if err := l.session.Close(); err != nil {
return fmt.Errorf("关闭会话失败: %v", err)
}
fmt.Printf("成功释放分布式锁: %s\n", l.lockKey)
return nil
}
// 带锁执行任务
func (l *DistributedLock) ExecuteWithLock(ctx context.Context,
timeout time.Duration, task func() error) error {
// 获取锁
if err := l.TryLock(ctx, timeout); err != nil {
return err
}
// 确保释放锁
defer l.Unlock()
// 执行任务
return task()
}
分布式配置管理
// watch/config_watcher.go
package watch
import (
"context"
"encoding/json"
"fmt"
"sync"
"time"
"go.etcd.io/etcd/client/v3"
)
type ConfigWatcher struct {
client *clientv3.Client
configs map[string]interface{}
watchers map[string][]func(string, interface{})
mutex sync.RWMutex
stopChan chan struct{}
}
func NewConfigWatcher(etcdEndpoints []string) (*ConfigWatcher, error) {
client, err := clientv3.New(clientv3.Config{
Endpoints: etcdEndpoints,
DialTimeout: 5 * time.Second,
})
if err != nil {
return nil, fmt.Errorf("创建etcd客户端失败: %v", err)
}
return &ConfigWatcher{
client: client,
configs: make(map[string]interface{}),
watchers: make(map[string][]func(string, interface{})),
stopChan: make(chan struct{}),
}, nil
}
// 获取配置
func (cw *ConfigWatcher) GetConfig(key string, defaultValue interface{}) interface{} {
cw.mutex.RLock()
defer cw.mutex.RUnlock()
if value, ok := cw.configs[key]; ok {
return value
}
return defaultValue
}
// 加载配置到本地缓存
func (cw *ConfigWatcher) LoadConfig(ctx context.Context, prefix string) error {
resp, err := cw.client.Get(ctx, prefix, clientv3.WithPrefix())
if err != nil {
return fmt.Errorf("加载配置失败: %v", err)
}
cw.mutex.Lock()
defer cw.mutex.Unlock()
for _, kv := range resp.Kvs {
cw.configs[string(kv.Key)] = string(kv.Value)
}
fmt.Printf("加载配置成功,共 %d 条配置\n", len(resp.Kvs))
return nil
}
// 监听配置变化
func (cw *ConfigWatcher) WatchConfig(prefix string) {
watchChan := cw.client.Watch(context.Background(), prefix,
clientv3.WithPrefix())
for {
select {
case <-cw.stopChan:
return
case resp, ok := <-watchChan:
if !ok {
return
}
for _, ev := range resp.Events {
key := string(ev.Kv.Key)
switch ev.Type {
case clientv3.EventTypePut:
cw.mutex.Lock()
cw.configs[key] = string(ev.Kv.Value)
cw.mutex.Unlock()
fmt.Printf("配置更新: %s = %s\n", key, ev.Kv.Value)
// 通知订阅者
cw.notifyWatchers(key, string(ev.Kv.Value))
case clientv3.EventTypeDelete:
cw.mutex.Lock()
delete(cw.configs, key)
cw.mutex.Unlock()
fmt.Printf("配置删除: %s\n", key)
cw.notifyWatchers(key, nil)
}
}
}
}
}
// 注册配置变更回调
func (cw *ConfigWatcher) AddConfigWatcher(key string,
callback func(string, interface{})) {
cw.mutex.Lock()
defer cw.mutex.Unlock()
cw.watchers[key] = append(cw.watchers[key], callback)
}
// 通知配置变更
func (cw *ConfigWatcher) notifyWatchers(key string, value interface{}) {
cw.mutex.RLock()
watchers := cw.watchers[key]
cw.mutex.RUnlock()
for _, watcher := range watchers {
go watcher(key, value)
}
}
主程序示例
// main.go
package main
import (
"context"
"fmt"
"log"
"time"
"etcd-demo/service"
"etcd-demo/watch"
)
func main() {
// etcd配置
etcdEndpoints := []string{"localhost:2379"}
// 1. 服务注册示例
registry, err := service.NewServiceRegistry(etcdEndpoints)
if err != nil {
log.Fatalf("创建服务注册器失败: %v", err)
}
serviceInfo := service.ServiceInfo{
Name: "user-service",
Address: "192.168.1.100",
Port: 8080,
Meta: map[string]string{
"version": "1.0.0",
"env": "prod",
},
TTL: 10,
}
if err := registry.RegisterService(serviceInfo); err != nil {
log.Fatalf("服务注册失败: %v", err)
}
// 2. 服务发现示例
discovery, err := service.NewServiceDiscovery(etcdEndpoints)
if err != nil {
log.Fatalf("创建服务发现器失败: %v", err)
}
instances, err := discovery.DiscoverService("user-service")
if err != nil {
log.Printf("服务发现失败: %v", err)
} else {
fmt.Printf("发现服务实例: %+v\n", instances)
}
// 启动服务监听
go discovery.WatchService("user-service")
// 3. 分布式锁示例
lock, err := service.NewDistributedLock(etcdEndpoints,
"/locks/order-process")
if err != nil {
log.Fatalf("创建分布式锁失败: %v", err)
}
// 带锁执行任务
err = lock.ExecuteWithLock(context.Background(), 3*time.Second, func() error {
fmt.Println("执行需要加锁的业务逻辑...")
time.Sleep(2 * time.Second)
return nil
})
if err != nil {
log.Printf("执行加锁任务失败: %v", err)
}
// 4. 配置管理示例
configWatcher, err := watch.NewConfigWatcher(etcdEndpoints)
if err != nil {
log.Fatalf("创建配置监听器失败: %v", err)
}
// 注册配置变更回调
configWatcher.AddConfigWatcher("/config/database", func(key string, value interface{}) {
fmt.Printf("数据库配置已更新: %v\n", value)
})
// 加载配置
if err := configWatcher.LoadConfig(context.Background(), "/config"); err != nil {
log.Printf("加载配置失败: %v", err)
}
// 监听配置变化
go configWatcher.WatchConfig("/config")
// 保持程序运行
fmt.Println("服务运行中,按 Ctrl+C 退出...")
select {}
}
部署和测试
1 启动Etcd服务
# 启动单节点etcd
etcd --listen-client-urls http://localhost:2379 \
--advertise-client-urls http://localhost:2379
# 使用docker启动
docker run -d --name etcd \
-p 2379:2379 \
-e ETCDCTL_API=3 \
quay.io/coreos/etcd:v3.5.0 \
/usr/local/bin/etcd \
--listen-client-urls http://0.0.0.0:2379 \
--advertise-client-urls http://0.0.0.0:2379
2 测试脚本
# 测试服务注册和发现
curl -v http://localhost:2379/v3/kv/put \
-X POST \
-d '{"key": "L3NlcnZpY2VzL3VzZXItc2VydmljZS8xOTIuMTY4LjEuMTAwOjgwODA="}'
# 查看服务列表
etcdctl get /services --prefix
# 测试分布式锁
etcdctl get /locks --prefix
# 测试配置管理
etcdctl put /config/database '{"host": "localhost", "port": 3306, "user": "root"}'
使用场景总结
1 微服务注册与发现
- 服务启动时自动注册
- 服务故障时自动剔除
- 动态感知服务状态变化
- 支持负载均衡和故障转移
2 分布式协调
- 分布式锁确保资源互斥访问
- Leader选举实现高可用
- 分布式任务调度
3 配置管理
- 集中式配置存储和分发
- 配置变更实时通知
- 支持动态调整配置
4 服务治理
- 服务健康检查
- 流量控制
- 版本控制
- 灰度发布
生产环境注意事项
1 高可用配置
# 集群模式配置 etcd --name etcd-node-1 \ --data-dir /data/etcd \ --listen-client-urls http://0.0.0.0:2379 \ --advertise-client-urls http://192.168.1.1:2379 \ --listen-peer-urls http://0.0.0.0:2380 \ --initial-advertise-peer-urls http://192.168.1.1:2380 \ --initial-cluster "etcd-node-1=http://192.168.1.1:2380,etcd-node-2=http://192.168.1.2:2380,etcd-node-3=http://192.168.1.3:2380" \ --initial-cluster-state new
2 安全配置
- 启用TLS加密通信
- 配置访问认证
- 设置权限控制
- 监控告警
3 性能优化
- 合理设置连接池大小
- 使用批量操作
- 监控延迟和吞吐量
- 设置合适的租约TTL
这个案例涵盖了Etcd最核心的应用场景,您可以根据具体需求进行扩展和调整。