基于Etcd实现案例

wen java案例 1

本文目录导读:

基于Etcd实现案例

  1. 项目结构
  2. 服务注册与发现
  3. 分布式锁实现
  4. 分布式配置管理
  5. 主程序示例
  6. 部署和测试
  7. 使用场景总结
  8. 生产环境注意事项

我将为您提供一个基于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最核心的应用场景,您可以根据具体需求进行扩展和调整。

抱歉,评论功能暂时关闭!