Skip to content

heytom-labs/heytom-dlm

Folders and files

NameName
Last commit message
Last commit date

Latest commit

 

History

2 Commits
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 

Repository files navigation

heytom-dlm

通用的Go语言分布式锁管理库,提供统一的接口和多种后端实现。

特性

  • 🔒 统一接口 - 提供标准的分布式锁接口,支持多种后端实现
  • 🔐 Token机制 - 基于令牌验证锁的所有权,防止误释放
  • ⏱️ 超时控制 - 支持获取锁的等待超时和锁的自动过期
  • 🔄 锁续期 - 支持刷新锁的过期时间,适用于长时间任务
  • 🎯 自动管理 - 提供回调模式自动管理锁的生命周期
  • 🚀 多后端支持 - 支持Redis和etcd,易于扩展
  • 🛡️ 企业级 - 符合企业级分布式锁标准,防止死锁和锁泄露

安装

go get github.com/heytom-labs/heytom-dlm

快速开始

Redis实现

package main

import (
    "context"
    "fmt"
    "time"

    "github.com/heytom-labs/heytom-dlm"
    "github.com/heytom-labs/heytom-dlm/redis"
    goredis "github.com/redis/go-redis/v9"
)

func main() {
    // 创建Redis客户端
    client := goredis.NewClient(&goredis.Options{
        Addr: "localhost:6379",
    })
    defer client.Close()

    // 创建分布式锁
    locker := redis.NewRedisLocker(client)

    ctx := context.Background()
    key := "my-resource-lock"

    // 获取锁
    token, err := locker.Lock(ctx, key, 10*time.Second, 5*time.Second)
    if err != nil {
        panic(err)
    }
    fmt.Printf("Lock acquired: %s\n", token)

    // 执行业务逻辑
    // ...

    // 释放锁
    if err := locker.Unlock(ctx, key, token); err != nil {
        panic(err)
    }
    fmt.Println("Lock released")
}

etcd实现

package main

import (
    "context"
    "time"

    "github.com/heytom-labs/heytom-dlm/etcd"
    clientv3 "go.etcd.io/etcd/client/v3"
)

func main() {
    // 创建etcd客户端
    client, err := clientv3.New(clientv3.Config{
        Endpoints:   []string{"localhost:2379"},
        DialTimeout: 5 * time.Second,
    })
    if err != nil {
        panic(err)
    }
    defer client.Close()

    // 创建分布式锁
    locker := etcd.NewEtcdLocker(client)

    ctx := context.Background()
    key := "/locks/my-resource"

    // 获取锁
    token, err := locker.Lock(ctx, key, 10*time.Second, 5*time.Second)
    if err != nil {
        panic(err)
    }

    // 执行业务逻辑
    // ...

    // 释放锁
    locker.Unlock(ctx, key, token)
}

核心接口

DistributedLocker

type DistributedLocker interface {
    // Lock 获取锁
    Lock(ctx context.Context, key string, ttl time.Duration, waitTimeout time.Duration) (LockToken, error)

    // Unlock 释放锁
    Unlock(ctx context.Context, key string, token LockToken) error

    // Refresh 刷新锁的过期时间
    Refresh(ctx context.Context, key string, token LockToken, ttl time.Duration) error

    // IsLocked 检查锁是否被持有
    IsLocked(ctx context.Context, key string) (bool, error)

    // LockAndExecute 获取锁后自动执行回调函数,执行完成后自动释放锁
    LockAndExecute(ctx context.Context, key string, ttl time.Duration, waitTimeout time.Duration, callback func(ctx context.Context) error) error
}

参数说明

  • key: 锁的键,用于标识要锁定的资源
  • ttl: 锁的过期时间,防止死锁
  • waitTimeout: 等待获取锁的超时时间,超过此时间返回失败
  • token: 锁的令牌,用于验证锁的所有权

使用示例

1. 基本用法

locker := redis.NewRedisLocker(client)

// 获取锁,等待最多5秒,锁10秒后自动过期
token, err := locker.Lock(ctx, "order:123", 10*time.Second, 5*time.Second)
if err != nil {
    if errors.Is(err, dlm.ErrLockFailed) {
        // 获取锁超时
    }
    return err
}

// 执行业务逻辑
processOrder()

// 释放锁
locker.Unlock(ctx, "order:123", token)

2. 自动管理锁(推荐)

err := locker.LockAndExecute(
    ctx,
    "order:123",
    10*time.Second,  // 锁的TTL
    5*time.Second,   // 等待超时
    func(ctx context.Context) error {
        // 在这里执行需要加锁保护的业务逻辑
        return processOrder()
    },
)
// 锁会自动释放,即使发生panic或错误

3. 长时间任务 - 锁续期

token, err := locker.Lock(ctx, "long-task", 10*time.Second, 5*time.Second)
if err != nil {
    return err
}
defer locker.Unlock(ctx, "long-task", token)

// 启动一个goroutine定期刷新锁
go func() {
    ticker := time.NewTicker(5 * time.Second)
    defer ticker.Stop()
    
    for {
        select {
        case <-ticker.C:
            if err := locker.Refresh(ctx, "long-task", token, 10*time.Second); err != nil {
                log.Printf("Failed to refresh lock: %v", err)
                return
            }
        case <-ctx.Done():
            return
        }
    }
}()

// 执行长时间任务
performLongTask()

4. 检查锁状态

locked, err := locker.IsLocked(ctx, "order:123")
if err != nil {
    return err
}

if locked {
    fmt.Println("Resource is locked")
} else {
    fmt.Println("Resource is available")
}

5. 错误处理

token, err := locker.Lock(ctx, "resource", 10*time.Second, 5*time.Second)
if err != nil {
    switch {
    case errors.Is(err, dlm.ErrLockFailed):
        // 获取锁超时
        log.Println("Failed to acquire lock within timeout")
    case errors.Is(err, context.Canceled):
        // 上下文被取消
        log.Println("Context canceled")
    default:
        // 其他错误
        log.Printf("Unexpected error: %v", err)
    }
    return err
}

// 释放锁时的错误处理
if err := locker.Unlock(ctx, "resource", token); err != nil {
    if errors.Is(err, dlm.ErrInvalidToken) {
        // Token无效或锁已被释放
        log.Println("Invalid token or lock already released")
    }
}

错误类型

var (
    ErrLockFailed    = errors.New("failed to acquire lock")      // 获取锁失败
    ErrLockNotHeld   = errors.New("lock not held")               // 锁未被持有
    ErrInvalidToken  = errors.New("invalid lock token")          // 无效的锁令牌
    ErrRefreshFailed = errors.New("failed to refresh lock")      // 刷新锁失败
    ErrUnlockFailed  = errors.New("failed to unlock")            // 释放锁失败
)

最佳实践

1. 始终设置合理的TTL

// ✅ 好的做法:设置合理的TTL防止死锁
token, err := locker.Lock(ctx, "key", 30*time.Second, 5*time.Second)

// ❌ 不好的做法:TTL过长可能导致资源长时间被锁定
token, err := locker.Lock(ctx, "key", 1*time.Hour, 5*time.Second)

2. 使用defer确保锁被释放

token, err := locker.Lock(ctx, "key", 10*time.Second, 5*time.Second)
if err != nil {
    return err
}
defer locker.Unlock(context.Background(), "key", token)

// 执行业务逻辑

3. 优先使用LockAndExecute

// ✅ 推荐:自动管理锁的生命周期
err := locker.LockAndExecute(ctx, "key", 10*time.Second, 5*time.Second, 
    func(ctx context.Context) error {
        return doWork()
    })

// ❌ 不推荐:手动管理容易忘记释放
token, _ := locker.Lock(ctx, "key", 10*time.Second, 5*time.Second)
doWork()
locker.Unlock(ctx, "key", token)

4. 合理设置等待超时

// 根据业务场景设置合理的等待超时
// 快速失败场景
token, err := locker.Lock(ctx, "key", 10*time.Second, 100*time.Millisecond)

// 可以等待的场景
token, err := locker.Lock(ctx, "key", 10*time.Second, 30*time.Second)

5. 长时间任务使用锁续期

// 对于执行时间不确定的任务,使用锁续期机制
token, err := locker.Lock(ctx, "key", 10*time.Second, 5*time.Second)
if err != nil {
    return err
}
defer locker.Unlock(ctx, "key", token)

// 定期刷新锁
go refreshLockPeriodically(ctx, locker, "key", token, 5*time.Second)

// 执行长时间任务
performLongRunningTask()

架构设计

Token机制

每次获取锁时会生成一个唯一的Token,只有持有正确Token的客户端才能释放或刷新锁。这防止了以下问题:

  • 客户端A获取锁后因网络问题超时,锁自动过期
  • 客户端B获取到锁并开始工作
  • 客户端A恢复后尝试释放锁,但因Token不匹配而失败,不会影响客户端B

原子性保证

  • Redis实现: 使用Lua脚本确保检查Token和操作的原子性
  • etcd实现: 使用etcd的事务和lease机制保证原子性

性能考虑

Redis vs etcd

特性 Redis etcd
性能 更高 较高
一致性 最终一致性 强一致性
适用场景 高并发、低延迟 强一致性要求
部署复杂度 简单 较复杂

性能优化建议

  1. 合理设置重试间隔: Redis实现默认50ms重试间隔
  2. 使用连接池: 复用Redis/etcd客户端连接
  3. 避免频繁刷新: 刷新间隔应为TTL的1/2到2/3

贡献

欢迎提交Issue和Pull Request!

许可证

MIT License

About

通用的Go语言分布式锁管理库,提供统一的接口和多种后端实现。

Resources

Stars

Watchers

Forks

Releases

Packages

Contributors

Languages