redisx

package module
v1.3.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Jul 21, 2026 License: MIT Imports: 14 Imported by: 0

README

redisx

基于 redis/go-redis v9 的生产级 Redis 客户端封装。

设计原则:错误全部通过返回值传递,库内不产生任何日志;不内置熔断、降级、重试编排等策略,只提供扩展点;核心 Redis 能力基于 go-redis,JSON 助手使用 gtkit/json/v2。

特性

  • 无全局变量,多实例(多服务器/多套配置)共存
  • 全局 key 前缀透明拼接,业务层无感知;支持每库(per-DB)独立前缀和自定义前缀连接符
  • 完整连接池 / 超时配置(PoolSize、MinIdleConns、DialTimeout、ReadTimeout 等)
  • 支持 TLS 连接(WithTLSConfig,适配云厂商强制加密实例)
  • 可选降级初始化(WithAllowPartialInit),失败 DB 经 InitError 按编号提取
  • Proxy 代理 String / Hash / List / Set / ZSet / Pub-Sub / Stream / Lua / Pipeline 的常用命令子集(非全量 Redis 命令),key 自动加前缀;完整命令面经 RawClient() / GetClient(db) 直接使用 go-redis;Client 内嵌默认 DB 代理,全部代理命令直接可用
  • Consume / ConsumePattern 受管订阅消费,ConsumeStream 消费组可靠消费(at-least-once)
  • 单实例分布式锁:TryLock / Release / Refresh / TTL,Lua 校验 token 杜绝误删且对底层自动重试幂等(不产生"幽灵锁");FencedLock 额外提供单调递增 fencing token,供下游资源拒绝旧持有者
  • DelByPattern:SCAN + UNLINK 批量删除,异步释放内存不阻塞 Redis 主线程
  • HealthCheck 健康检查、PoolStats 连接池统计透传,监控出口齐备
  • AddHook 一次调用为所有 DB 挂接 go-redis Hook,外接熔断 / 限流 / metrics / tracing
  • 包名 redisx 与 go-redis 的 redis 包不冲突,业务代码无需 import 别名

安装

go get github.com/gtkit/redisx@latest

要求 Go 1.26+、Redis >= 6.2

本库重度依赖需要 6.2 的能力(GetSetSET..GETZRANGEREV/BYSCOREXAUTOCLAIM 等),故最低承诺版本统一为 6.2;CI 在 Redis 6.2 / 7 / 8 上验证。

快速开始

import "github.com/gtkit/redisx"

c, err := redisx.NewClient(
    redisx.WithAddr("127.0.0.1:6379"),
    redisx.WithPassword(os.Getenv("REDIS_PASSWORD")), // 真实值从环境变量读取
    redisx.WithKeyPrefix("myapp"),                    // 全局前缀
    redisx.WithInitDBs(0, 1),                         // 初始化多个 DB
    redisx.WithInitDBPrefix(2, "session"),            // DB2 使用独立前缀
)
if err != nil {
    log.Fatal(err)
}
defer c.Close()

ctx := context.Background()

// 默认 DB 操作,实际 key 为 "myapp:user:1"
c.Set(ctx, "user:1", "hello", time.Hour)
val, err := c.Get(ctx, "user:1").Result()

// 切换 DB 链式调用,DB2 的 key 前缀为 "session:"
token, err := c.MustSelectDB(2).Get(ctx, "token:abc").Result()

初始化与配置

全部配置项(Functional Options)
Option 说明 默认值
WithAddr(addr) 服务器地址host:port(必填)
WithUsername(u) Redis 6+ ACL 用户名 空(不使用)
WithPassword(p) 认证密码 空(不认证)
WithDB(db) 默认 DB 编号 0
WithInitDBs(dbs...) 初始化多个 DB(共享全局前缀) 仅默认 DB
WithInitDBPrefix(db, prefix) 初始化单个 DB 并指定独立前缀
WithDBConfig(db, prefix) 已废弃兼容别名,等价于 WithInitDBPrefix
WithKeyPrefix(prefix) 全局 key 前缀 空(不加前缀)
WithKeyPrefixSeparator(s) key 前缀与原始 key 之间的连接符 :
WithChannelPrefix(prefix) Pub/Sub channel 前缀(独立于 key 前缀) 空(不加前缀)
WithChannelPrefixSeparator(s) channel 前缀与原始 channel 之间的连接符 :
WithPoolSize(n) 每个 DB 的最大连接数 10
WithMinIdleConns(n) 最小空闲连接数 3
WithMaxRetries(n) 命令失败重试次数(非幂等命令慎用,见下文) 3
WithDialTimeout(d) 建连超时 5s
WithReadTimeout(d) 读超时 3s
WithWriteTimeout(d) 写超时 3s
WithIdleTimeout(d) 空闲连接回收时间 5m
WithTLSConfig(cfg) TLS 配置 nil(不启用)
WithAllowPartialInit() 降级模式:失败 DB 缺席集合,错误聚合返回(DefaultDB 仍须成功) 关闭(全有或全无)

关于 WithMaxRetries:对非幂等命令(如 INCRLPUSH),读超时后的自动重试可能导致命令被重复执行。对此敏感的场景请设置为 0 关闭重试,或在业务层用 Lua 脚本保证幂等。

初始化语义
  • NewClient 会对每个声明的 DB 建立独立连接池并执行 PING 验证,DefaultDB 优先拨号:它承载 Client 级快捷方法,失败时立即整体失败,其余 DB 不再拨号;其余 DB 按编号升序拨号,错误信息确定有序。
  • 默认全有或全无:任一 DB 验证失败,整体返回错误并回收已建连接。
  • clients / proxies 集合在构建完成后只读,后续并发使用无需加锁。
降级初始化与 InitError

启用 WithAllowPartialInit() 后,非默认 DB 初始化失败不再阻断整体:失败的 DB 缺席集合,错误以 *InitError可用的 Client 一同返回(两者可同时非 nil):

c, err := redisx.NewClient(
    redisx.WithAddr("127.0.0.1:6379"),
    redisx.WithInitDBs(0, 1, 2),
    redisx.WithAllowPartialInit(),
)
// ⚠️ 降级模式下 c 与 err 可同时非 nil!
// 不要写惯用的 `if err != nil { return nil, err }`——那会把可用的 Client 丢掉。
// 先用 errors.As 判断是否为部分失败,再决定降级使用还是拒绝启动。
var ie *redisx.InitError
if errors.As(err, &ie) {
    for db, cause := range ie.Failed {
        log.Printf("db=%d 初始化失败: %v", db, cause) // 按业务决定降级还是拒绝启动
    }
}
if c == nil {
    log.Fatal(err) // DefaultDB 失败时 Client 为 nil,整体不可用
}

InitError.Unwrap() 返回各 DB 的底层错误,errors.Is 可穿透到网络层错误。对缺席 DB 调用 SelectDB 返回错误。

TLS 连接
c, err := redisx.NewClient(
    redisx.WithAddr("redis.example.com:6380"),
    redisx.WithPassword(os.Getenv("REDIS_PASSWORD")),
    redisx.WithTLSConfig(&tls.Config{MinVersion: tls.VersionTLS12}),
)

Key 前缀机制

设置 WithKeyPrefix("myapp") 后,所有带 key 参数的命令默认自动拼接为 myapp:{key},业务层无感知。可通过 WithKeyPrefixSeparator(".") 改为 myapp.{key} 等自定义格式。

前缀优先级WithInitDBPrefix 的 per-DB 前缀 > 全局 WithKeyPrefix > 不加前缀。

Pub/Sub channel 是 Redis 实例级全局命名空间,不随 DB 切换;WithKeyPrefix 不影响 channel。如需隔离 topic,请使用 WithChannelPrefix("myapp")

不拼前缀的出口(需要自己负责完整 key):

  • GetClient(db) / DefaultClient() / Proxy.RawClient() 返回原生 *redis.Client
  • Pipeline() / TxPipeline() 内的命令(见 Pipeline 与事务);
  • Lua 脚本体内自行构造的 key(Eval 系列只对 keys 参数列表拼前缀)。

陷阱Scan 返回的 key 已含前缀,直接回传给本库其他带前缀方法(如 Del)会二次拼接。遍历 key 请优先用 ScanKeys——自动翻页且产出已剥前缀的 key,可直接回传:

for key, err := range c.ScanKeys(ctx, "user:*") {
    if err != nil {
        return err
    }
    c.Del(ctx, key) // key 不含前缀,不会二次拼接
}

批量删除直接用 DelByPattern

手动拼接工具(Pipeline 等场景):

p := c.MustSelectDB(0)
full := p.Key("user:1")           // "myapp:user:1"
fulls := p.Keys("k1", "k2")       // ["myapp:k1", "myapp:k2"]
prefix := c.Prefix()              // "myapp"
同一 DB 多命名空间

同一个 DB 里可以复用同一连接池派生多套 key 命名空间;Pub/Sub channel 仍按独立 channel 前缀处理:

db2 := c.MustSelectDB(2)

cache, err := db2.WithPrefix("cache")
if err != nil {
    return err
}
triggers, err := db2.WithPrefix("trigger")
if err != nil {
    return err
}

cache.Set(ctx, "user:42", payload, time.Minute)      // Redis key: cache:user:42
triggers.Set(ctx, "user:42", "refresh", time.Minute) // Redis key: trigger:user:42

// channel 不随 key prefix 变化;需要 topic 隔离请在 NewClient 时配置 WithChannelPrefix。
db2.Publish(ctx, "events:user", "changed")

多 DB 使用

p, err := c.SelectDB(1)   // DB 未初始化时返回错误(错误信息列出可用 DB)
p := c.MustSelectDB(1)    // 未初始化时 panic,仅用于启动阶段或确定存在的场景

rdb, ok := c.GetClient(1) // 原生 *redis.Client,不拼前缀
rdb := c.DefaultClient()  // 默认 DB 的原生客户端

Client 内嵌默认 DB 的 ProxyProxy 的全部命令方法在 Client 上直接可用,等价于 c.MustSelectDB(defaultDB).Xxx(...)。每个 DB 的 Proxy 在初始化时缓存,重复获取无额外分配。

需要连接多个 Redis 服务器时,分别 NewClient 即可,实例间完全独立(无全局变量)。


常用命令

所有命令返回 go-redis 原生的 *redis.XxxCmd,用 .Result() / .Val() / .Err() 取值;key 一律自动拼前缀。

错误处理

key 不存在时 GetHGet 等返回 redis.Nil,这是"未命中"而非故障,务必区分:

val, err := c.Get(ctx, "user:1").Result()
switch {
case errors.Is(err, redis.Nil):
    // 缓存未命中,回源
case err != nil:
    // 真实错误(网络/超时),按业务决定降级或上报
default:
    // 命中
}

库内不产生日志,所有失败都经返回值传出,日志策略完全由调用方决定。

String / 计数器
c.Set(ctx, "k", "v", time.Hour)        // expiration 为 0 表示不过期
c.SetEX(ctx, "k", "v", time.Hour)      // SET 带过期(SETEX 现代等价)
ok, _ := c.SetNX(ctx, "k", "v", ttl).Result() // 不存在才写,true=写入成功
old, _ := c.GetSet(ctx, "k", "new").Result()  // 写新值返回旧值(SET..GET,Redis>=6.2)
v, _ := c.GetDel(ctx, "k").Result()    // 取值并删除

c.Incr(ctx, "counter")                 // +1
c.IncrBy(ctx, "counter", 10)
c.IncrByFloat(ctx, "price", 0.5)
c.Decr(ctx, "counter")
c.DecrBy(ctx, "counter", 10)

vals, _ := c.MGet(ctx, "k1", "k2").Result()       // 批量取,缺失项为 nil
c.MSet(ctx, "k1", "v1", "k2", "v2")               // 交替 key-value,长度必须为偶数且 key 位必须为 string

结构体值用 JSON 泛型助手(包级函数,Client 内嵌 Proxy 直接传):

redisx.SetJSON(ctx, c.Proxy, "user:1", User{Name: "alice"}, time.Hour)

u, err := redisx.GetJSON[User](ctx, c.Proxy, "user:1")
if errors.Is(err, redis.Nil) { /* key 不存在,u 为零值 */ }

非 JSON 编码可显式存取 bytes,或提供业务自己的 Codec(例如 protobuf/gob/msgpack 适配器;redisx 不引入这些依赖):

_ = redisx.SetBytes(ctx, c.Proxy, "payload:1", data, time.Hour)
data, err := redisx.GetBytes(ctx, c.Proxy, "payload:1")

_ = redisx.SetCodec(ctx, c.Proxy, "user:pb:1", protoCodec, user, time.Hour)
user, err := redisx.GetCodec[User](ctx, c.Proxy, "user:pb:1", protoCodec)
Key 管理
c.Del(ctx, "k1", "k2")                  // 返回删除数量
n, _ := c.Exists(ctx, "k1", "k2").Result() // 返回存在数量
c.Expire(ctx, "k", time.Hour)
c.ExpireAt(ctx, "k", deadline)
c.Persist(ctx, "k")                     // 移除过期时间
d, _ := c.TTL(ctx, "k").Result()        // -2=不存在 -1=无过期
c.PTTL(ctx, "k")                        // 毫秒精度
c.Type(ctx, "k")
c.Rename(ctx, "old", "new")             // 两个 key 都拼前缀
Hash
c.HSet(ctx, "user:1", "name", "alice", "age", 18) // field-value 对
v, _ := c.HGet(ctx, "user:1", "name").Result()
all, _ := c.HGetAll(ctx, "user:1").Result()        // map[string]string
vals, _ := c.HMGet(ctx, "user:1", "name", "age").Result()
c.HDel(ctx, "user:1", "age")
ok, _ := c.HExists(ctx, "user:1", "name").Result()
n, _ := c.HLen(ctx, "user:1").Result()
c.HIncrBy(ctx, "user:1", "score", 10)
c.HIncrByFloat(ctx, "user:1", "balance", 1.5)
List
c.LPush(ctx, "list", "a", "b")
c.RPush(ctx, "list", "c")
v, _ := c.LPop(ctx, "list").Result()
v, _ := c.RPop(ctx, "list").Result()
items, _ := c.LRange(ctx, "list", 0, -1).Result() // -1 表示最后一个
n, _ := c.LLen(ctx, "list").Result()
c.LRem(ctx, "list", 1, "a")            // 移除 1 个等于 "a" 的元素
v, _ := c.LIndex(ctx, "list", 0).Result()
c.LTrim(ctx, "list", 0, 99)            // 只保留前 100 个
Set
c.SAdd(ctx, "tags", "go", "redis")
members, _ := c.SMembers(ctx, "tags").Result()
ok, _ := c.SIsMember(ctx, "tags", "go").Result()
c.SRem(ctx, "tags", "redis")
n, _ := c.SCard(ctx, "tags").Result()
v, _ := c.SRandMember(ctx, "tags").Result() // 随机取(不删)
v, _ := c.SPop(ctx, "tags").Result()        // 随机弹出(删)
Sorted Set
c.ZAdd(ctx, "rank", redis.Z{Score: 100, Member: "alice"})
score, _ := c.ZScore(ctx, "rank", "alice").Result()
i, _ := c.ZRank(ctx, "rank", "alice").Result()       // 升序排名,从 0 开始
top, _ := c.ZRange(ctx, "rank", 0, 9).Result()       // 按排名升序
top, _ = c.ZRevRange(ctx, "rank", 0, 9).Result()     // 逆序(ZRANGE REV,Redis>=6.2)
hits, _ := c.ZRangeByScore(ctx, "rank", &redis.ZRangeBy{
    Min: "60", Max: "100", Offset: 0, Count: 10,     // 按分值范围(Redis>=6.2)
}).Result()
c.ZRem(ctx, "rank", "alice")
c.ZRemRangeByScore(ctx, "rank", "0", "59")
n, _ := c.ZCard(ctx, "rank").Result()
n, _ = c.ZCount(ctx, "rank", "60", "100").Result()
c.ZIncrBy(ctx, "rank", 5, "alice")
批量删除与遍历
// SCAN + UNLINK 批量删除:UNLINK 后台异步释放内存,大 value 不阻塞主线程。
// pattern 自动拼前缀("user:*" 实际匹配 "myapp:user:*"),返回删除总数。
// 建议传带超时的 ctx 控制执行时间。
deleted, err := c.DelByPattern(ctx, "user:*")

// 裸 SCAN:match 拼前缀;注意返回的 key 已含前缀(见前缀机制一节)
keys, cursor, err := c.Scan(ctx, 0, "user:*", 100).Result()

Pipeline 与事务

Pipeline()(批量打包)与 TxPipeline()(MULTI/EXEC 事务)返回 go-redis 原生 redis.PipelinerPipeline 内的命令不会自动拼前缀,用 Key / Keys 手动拼:

p := c.MustSelectDB(0)
pipe := p.Pipeline()
get := pipe.Get(ctx, p.Key("user:1"))
pipe.Del(ctx, p.Keys("k1", "k2")...)
if _, err := pipe.Exec(ctx); err != nil && !errors.Is(err, redis.Nil) {
    return err
}
val, _ := get.Result()

Lua 脚本

keys 参数列表中的每个 key 自动拼前缀;脚本体内自行构造的 key 不受管:

// 直接执行
n, err := c.Eval(ctx, `return redis.call("incrby", KEYS[1], ARGV[1])`,
    []string{"counter"}, 10).Int64()

// 按 SHA 执行已缓存脚本
n, err = c.EvalSha(ctx, sha1, []string{"counter"}, 10).Int64()

// 推荐:*redis.Script 自动处理 NOSCRIPT 回退
var script = redis.NewScript(`return redis.call("get", KEYS[1])`)
v, err := c.EvalScript(ctx, script, []string{"k"}).Result()

Pub/Sub

Redis Pub/Sub 为 at-most-once:不落盘、不重投,订阅者断线期间的消息永久丢失。适合在线通知、缓存失效广播等"丢了无所谓"的场景;需要可靠投递请用 Stream

Pub/Sub channel 默认不使用 key 前缀。需要 topic 命名空间隔离时,在初始化时显式设置 WithChannelPrefix("myapp"),或配合 WithChannelPrefixSeparator(".") 自定义拼接格式。

// 发布(默认使用原始 channel;不会套用 WithKeyPrefix)
c.Publish(ctx, "events:user", "user_created")

// 受管消费(推荐):订阅确认、连接关闭、ctx 退出、panic 恢复均由库内管理
err := c.Consume(ctx, func(m *redis.Message) {
    fmt.Println(m.Channel, m.Payload)
}, "events:user")
// ctx 取消时返回 context.Canceled,可据此判断优雅退出

// 模式订阅的受管消费,生命周期语义与 Consume 一致
err = c.ConsumePattern(ctx, func(m *redis.Message) {
    fmt.Println(m.Pattern, m.Channel, m.Payload)
}, "events:*")

// 裸订阅 / 模式订阅:返回 *redis.PubSub,调用方负责 Close
sub := c.Subscribe(ctx, "events:user")
defer sub.Close()
sub = c.PSubscribe(ctx, "events:*")
defer sub.Close()

Consume / ConsumePattern 的 handler 串行执行以保证单频道顺序,耗时处理请投递到业务自有 worker;handler panic 会被恢复,消费终止并以 error 返回。

如果配置了 WithChannelPrefix("myapp"),上面的 "events:user" 实际发布 / 订阅到 "myapp:events:user";模式 "events:*" 实际订阅到 "myapp:events:*"


分布式锁

单 Redis 实例锁:SET NX + 随机 token,Lua 校验 token 后原子释放/续期,杜绝误删他人锁。

lock, err := c.TryLock(ctx, "job:daily-report", 30*time.Second)
if errors.Is(err, redisx.ErrLockNotObtained) {
    return nil // 他人持有,按业务节奏稍后重试
}
if err != nil {
    return err
}
defer lock.Release(ctx)

// 长任务期间显式续期;锁已失去返回 ErrLockLost
if err := lock.Refresh(ctx, 30*time.Second); errors.Is(err, redisx.ErrLockLost) {
    return nil // 锁已过期被他人接管,停止当前工作
}

更推荐的闭包模式——WithLock 封装"拿锁—执行—必释放"(含 panic 路径),消除漏 Release 风险:

err := c.WithLock(ctx, "job:daily-report", 30*time.Second, func(ctx context.Context) error {
    return doReport(ctx) // 预计耗时须显著小于 ttl
})
switch {
case errors.Is(err, redisx.ErrLockNotObtained):
    return nil // 他人持有,跳过本轮
case errors.Is(err, redisx.ErrLockLost):
    // fn 执行超过 ttl,互斥可能已被破坏——触发业务侧补偿/告警
}

要点与边界:

  • TryLock 非阻塞;需要阻塞等待时由调用方按业务节奏循环重试(库不内置轮询策略)。冲突错误用 errors.Is(err, ErrLockNotObtained) 判断,错误消息携带完整 key 便于排障。
  • ttl 必须至少为 1ms——Redis TTL 精度为毫秒,无 TTL 的锁等于死锁隐患。
  • Release / Refresh 在锁已过期或被他人重新获取时返回 ErrLockLost,且不会影响他人的锁
  • 观测自检:lock.Key() 返回完整锁 key(供日志/打点);lock.TTL(ctx) 校验 token 后原子返回剩余时长,锁已失去返回 ErrLockLost。注意 TTL 仅供观测——查询与后续操作之间锁仍可能过期,互斥正确性依赖 Release/Refresh 自身的 token 校验。
  • 非 RedLock:主从异步复制下故障切换瞬间存在双持有的理论窗口,关键互斥请在业务层做幂等兜底。
  • 不提供 watchdog 自动续期,生命周期由业务显式管理。
  • 获取对底层自动重试幂等:SET NX 因响应丢失被重发时,重发命中的锁值仍等于本次 token 亦视为获得,不会把自己已持有的锁误报为被他人持有。
  • lock.TTL(ctx) 在锁被外部置为无过期时间(契约外的永久锁)时返回 ErrLockNoExpiry,而非把死锁隐患伪装成正常。
fencing token(FencedLock)

进程暂停 / GC 卡顿 / 网络阻塞会让旧持有者的锁在 TTL 过期、新持有者已拿到锁后,旧持有者仍继续操作外部资源——随机 token 只能防误删他人锁,挡不住这种"过期后仍写"。FencedLock 在获取时原子生成一个在同一 key 上单调递增的 fencing token:

lock, err := c.FencedLock(ctx, "resource:42", 30*time.Second)
if err != nil { /* errors.Is(err, redisx.ErrLockNotObtained) */ }
defer lock.Release(ctx)

writeToResource(lock.Fence(), payload) // 把 token 传给下游资源

FencedLock 嵌入 LockKey / Release / Refresh / TTL 语义完全一致,额外提供 Fence()栅栏仅在下游被保护资源记录见过的最大 token 并拒绝更小者时才生效,本库只负责原子生成与暴露 token。fencing 计数器是一个持久(无 TTL)的 <锁key>:__fence__Release 不删除它以保证跨获取单调,请勿手动删除。


Stream 消费组(可靠消息)

Stream 是 Redis 内置的可靠消息队列(at-least-once):消息持久化、消费组分摊、ACK 确认、pending 重投、死消费者接管。

生产
id, err := c.XAdd(ctx, &redis.XAddArgs{
    Stream: "orders",                          // 自动拼前缀(不修改传入的 args)
    MaxLen: 100000, Approx: true,              // 建议设置,控制流长度
    Values: map[string]any{"order_id": 1001},
}).Result()
受管消费
err := c.ConsumeStream(ctx, redisx.StreamConfig{
    Stream:   "orders",
    Group:    "billing",
    Consumer: "worker-1",                  // 组内唯一(如实例 ID)
    BatchSize: 16,                         // 单次最多读取数,默认 16
    Block:    5 * time.Second,             // 无消息时阻塞等待,默认 5s
    AutoClaimMinIdle: 30 * time.Second,    // 可选:接管死消费者闲置超时的消息
    OnError: func(m redis.XMessage, err error) {
        log.Printf("msg %s failed: %v", m.ID, err) // 库内无日志,失败感知交给业务
    },
}, func(m redis.XMessage) error {
    return process(m.Values) // 返回 nil 即 XACK;返回 error 则留在 pending 等待重投
})

StreamConfig 字段:

字段 说明 默认
Stream 流名,自动拼前缀(必填)
Group 消费组名,不存在时自动创建(MKSTREAM)(必填)
Consumer 消费者名,组内应唯一(必填)
BatchSize 单次 XREADGROUP 最多读取数 16
Block 无消息时阻塞等待时长(0 或负值取默认,无法表达无限阻塞) 5s
AutoClaimMinIdle >0 时周期接管闲置超过该时长的他人 pending 消息(Redis>=6.2);实际接管周期下限受 Block 钳制 0(关)
MaxDeliver >0 时启用死信策略:投递次数超过该值的消息转入死信流(即 handler 最多尝试 MaxDeliver 次) 0(关)
DeadLetterStream 死信流名,自动拼前缀;MaxDeliver > 0 时必填且不得与 Stream 同名
OnError handler 业务失败时的同步回调,应保持轻量;回调 panic 终止消费 nil

生命周期语义:

  • 消费组不存在时自动创建;启动时先续传本消费者的 pending(崩溃重启不丢已读未确认的消息),再消费新消息。
  • handler 串行执行保证顺序;panic 被恢复,消费终止并以 error 返回。
  • ctx 取消时返回 ctx 的错误(errors.Is(err, context.Canceled) 判断优雅退出)。
死信队列(毒消息隔离)

默认失败消息无限重投。配置 MaxDeliver + DeadLetterStream 后,投递次数超限的毒消息经 Lua(XADD + XACK)转入死信流,不再阻塞消费。该 Lua 只保证服务端两步不被其他命令穿插;死信投递语义为 at-least-once——脚本执行成功但客户端响应丢失后的重发可能重复写死信(_redisx_origin_id 相同),去重由下游按 _redisx_origin_id 负责:

err := c.ConsumeStream(ctx, redisx.StreamConfig{
    Stream: "orders", Group: "billing", Consumer: "worker-1",
    MaxDeliver:       5,            // handler 最多尝试 5 次
    DeadLetterStream: "orders:dlq", // 第 6 次投递前转入死信流
    OnError: func(m redis.XMessage, err error) {
        if errors.Is(err, redisx.ErrMessageDeadLettered) {
            alert("毒消息已隔离", m.ID) // 死信化事件经 OnError 通知
        }
    },
}, handler)

死信消息保留原字段,并附加元数据:_redisx_origin_stream(原流)、_redisx_origin_id(原消息 ID)、_redisx_deliveries(投递次数)、_redisx_dead_at(死信时间,RFC3339)。

边界:

  • 投递次数来自 Redis pending entry 的 delivery counter(重投/接管自增);仅 pending 续传与 AutoClaim 路径检查,新消息热路径零额外开销。
  • 死信流的裁剪、监控、重放由业务负责,库内不做 MAXLEN 限制;重放时可按 _redisx_origin_id 幂等。
  • 业务侧手工 XCLAIM ... JUSTID 不自增计数,会让该次接管不计入 MaxDeliver。
失败处理与监控

handler 返回 error 的消息不会 ACK,留在 pending:重启续传或被 AutoClaim 接管时重投(启用死信策略时超限即隔离)。持续失败的消息会堆积,务必监控:

// pending 摘要:总数、最小/最大 ID、各消费者的数量
summary, err := c.XPending(ctx, "orders", "billing").Result()

// pending 明细:按条件过滤(不修改传入的 args)
details, err := c.XPendingExt(ctx, &redis.XPendingExtArgs{
    Stream: "orders", Group: "billing",
    Start: "-", End: "+", Count: 10,
    Idle: time.Minute, // 只看闲置超过 1 分钟的
}).Result()
管理命令
n, _ := c.XLen(ctx, "orders").Result()                  // 流长度
msgs, _ := c.XRange(ctx, "orders", "-", "+").Result()   // 按 ID 区间读取
c.XAck(ctx, "orders", "billing", "1-0")                 // 手动确认
c.XDel(ctx, "orders", "1-0")                            // 删除消息
c.XTrimMaxLen(ctx, "orders", 100000)                    // 裁剪流长度

// 组管理
c.XGroupCreateMkStream(ctx, "orders", "billing", "$")   // 提前建组(只消费新消息)
groups, _ := c.XInfoGroups(ctx, "orders").Result()      // 组状态
cs, _ := c.XInfoConsumers(ctx, "orders", "billing").Result() // 消费者状态(发现死条目)
c.XGroupDelConsumer(ctx, "orders", "billing", "pod-old") // 清理死亡消费者
c.XGroupDestroy(ctx, "orders", "billing")               // 销毁组

消费者条目运维:以实例 ID 做 Consumer 名时,滚动发布会在组内累积死亡消费者条目(AutoClaim 只接管消息、不删条目)。请在实例下线钩子或定期任务中用 XInfoConsumers 找出长期闲置者、XGroupDelConsumer 清理;其未确认消息会被丢弃,清理前确认 pending 已被接管。


监控与运维

// 健康检查:PING 所有已初始化 DB,不短路,聚合返回全部失败;nil 表示全部健康
if err := c.HealthCheck(ctx); err != nil {
    log.Printf("redis unhealthy: %v", err)
}

// 连接池统计:按 DB 编号透传 go-redis PoolStats(命中/未命中/超时/连接数)
for db, s := range c.PoolStats() {
    metrics.Report(db, s.Hits, s.Misses, s.Timeouts, s.TotalConns, s.IdleConns)
}

// 优雅关闭:关闭所有 DB 连接池,多 DB 失败用 errors.Join 合并返回
defer c.Close()

当前库版本经 redisx.Version 常量获取。


Hook 扩展点(外接熔断 / 限流 / 观测)

库本身不内置熔断 / 降级策略——这类决策属于业务层(回源数据库、返回兜底值还是直接报错,只有调用方知道)。库提供的是扩展点:AddHook 把任意 go-redis redis.Hook 一次性安装到所有已初始化 DB 的底层客户端上,可用于接入 gobreaker、sentinel-golang 等熔断库,或挂接限流、metrics、tracing 中间件。

// breakerHook 用任意熔断器包装 Redis 命令执行链(以 gobreaker 风格为例)
type breakerHook struct {
    cb *gobreaker.CircuitBreaker // 业务自选的熔断器实现
}

func (h *breakerHook) DialHook(next redis.DialHook) redis.DialHook { return next }

func (h *breakerHook) ProcessHook(next redis.ProcessHook) redis.ProcessHook {
    return func(ctx context.Context, cmd redis.Cmder) error {
        _, err := h.cb.Execute(func() (any, error) {
            return nil, next(ctx, cmd) // 熔断打开时直接快速失败,不再打到 Redis
        })
        return err
    }
}

func (h *breakerHook) ProcessPipelineHook(next redis.ProcessPipelineHook) redis.ProcessPipelineHook {
    return func(ctx context.Context, cmds []redis.Cmder) error {
        _, err := h.cb.Execute(func() (any, error) {
            return nil, next(ctx, cmds)
        })
        return err
    }
}

// 初始化阶段安装,对所有 DB 生效
c, err := redisx.NewClient(redisx.WithAddr("127.0.0.1:6379"), redisx.WithInitDBs(0, 1))
if err != nil {
    log.Fatal(err)
}
c.AddHook(&breakerHook{cb: newBreaker()})

注意事项:

  • 调用时机:与 go-redis 原生 AddHook 约束一致,必须在 NewClient 之后、开始并发执行命令之前安装完成,运行中途添加不保证并发安全。
  • 只想对个别 DB 安装时,改用 GetClient(db) 拿到原生客户端后自行 AddHook
  • 熔断快速失败时业务拿到的是熔断器返回的错误(如 gobreaker 的 ErrOpenState),可据此走降级分支(回源、兜底值等)。

并发安全

ClientProxyNewClient 返回后内部状态只读,可在任意 goroutine 并发使用;底层 *redis.Client 由 go-redis 连接池保证并发安全。

已知边界

  • 仅支持单实例 Redis,不支持 Cluster / Sentinel 拓扑(Cluster 协议无多 DB 概念,如有需求属独立特性)。
  • Pub/Sub 为 at-most-once;可靠投递请用 Stream。
  • 分布式锁为单实例语义(非 RedLock)。

License

MIT,详见 LICENSE

Documentation

Overview

Package redisx 提供基于 github.com/redis/go-redis/v9 的生产级 Redis 客户端封装。

redisx 是 github.com/gtkit/redis(v1/v2)的后继包,能力完整覆盖两者, 新项目请直接使用本包。主要特性:

  • 无全局变量,支持多实例(多服务器/多配置)共存
  • 全局 Key 前缀透明封装,业务层无感知
  • 支持 per-DB 独立前缀
  • Pub/Sub channel 前缀独立配置,默认不复用 key 前缀
  • 完整连接池/超时配置(PoolSize、DialTimeout、ReadTimeout 等)
  • 支持 TLS 连接(WithTLSConfig)
  • 可选降级初始化(WithAllowPartialInit,失败 DB 缺席集合、错误聚合返回)
  • SCAN + UNLINK 批量删除,异步释放内存不阻塞 Redis 主线程
  • 无内部日志:所有失败经返回值传递,由调用方决定日志策略;零日志框架依赖
  • HealthCheck 健康检查、PoolStats 连接池统计透传
  • AddHook 一次调用为所有 DB 挂接 go-redis Hook(外接熔断/限流/观测)
  • Client 内嵌默认 DB 的 Proxy,全部命令方法直接可用
  • 包名 redisx 与 go-redis 的 redis 包不冲突,业务层无需 import 别名

快速使用:

c, err := redisx.NewClient(
    redisx.WithAddr("127.0.0.1:6379"),
    redisx.WithPassword("123456"),
    redisx.WithKeyPrefix("app:demo"),
    redisx.WithInitDBs(0, 1, 2),
)
if err != nil {
    log.Fatal(err)
}
defer c.Close()

ctx := context.Background()

// 默认 DB 操作,key 自动添加前缀 "app:demo:user:1"
c.Set(ctx, "user:1", "hello", 0)

// 切换 DB1 链式调用
val, _ := c.MustSelectDB(1).Get(ctx, "config").Result()

// SCAN 安全批量删除
deleted, _ := c.DelByPattern(ctx, "user:*")
Example
package main

import (
	"context"
	"fmt"
	"log"
	"time"

	"github.com/gtkit/redisx"

	goredis "github.com/redis/go-redis/v9"
)

func main() {
	// ─── 1. 极简用法(仅 Addr 必填)───
	c, err := redisx.NewClient(
		redisx.WithAddr("127.0.0.1:6379"),
	)
	if err != nil {
		log.Fatal(err)
	}
	defer c.Close()

	ctx := context.Background()
	_ = ctx

	// ─── 2. 完整配置 ───
	c2, err := redisx.NewClient(
		redisx.WithAddr("127.0.0.1:6379"),
		redisx.WithUsername("default"), // Redis 6+ ACL
		redisx.WithPassword("123456"),
		redisx.WithDB(0),                      // 默认 DB
		redisx.WithKeyPrefix("app:demo"),      // 全局前缀
		redisx.WithChannelPrefix("app:topic"), // Pub/Sub channel 前缀(独立于 key 前缀)
		redisx.WithInitDBs(0, 1, 2),           // 初始化多个 DB(共享全局前缀)
		redisx.WithInitDBPrefix(3, "session"), // DB3 使用独立前缀 "session"
		redisx.WithPoolSize(20),
		redisx.WithMinIdleConns(5),
		redisx.WithMaxRetries(3),
		redisx.WithDialTimeout(5*time.Second),
		redisx.WithReadTimeout(3*time.Second),
		redisx.WithWriteTimeout(3*time.Second),
		redisx.WithIdleTimeout(5*time.Minute),
	)
	if err != nil {
		log.Fatal(err)
	}
	defer c2.Close()

	// ─── 3. 默认 DB 操作(key 自动加前缀 "app:demo:user:1")───
	c2.Set(ctx, "user:1", "hello", time.Hour)
	val, _ := c2.Get(ctx, "user:1").Result()
	fmt.Println(val) // "hello"

	// ─── 4. 多 DB 切换(链式调用)───
	db1 := c2.MustSelectDB(1)
	db1.Set(ctx, "config:timeout", "30s", 0)
	cfg, _ := db1.Get(ctx, "config:timeout").Result()
	fmt.Println(cfg) // "30s"

	// DB3 使用独立前缀 "session":实际 key 为 "session:token:abc"
	db3 := c2.MustSelectDB(3)
	db3.Set(ctx, "token:abc", "user_123", 30*time.Minute)

	// 安全获取(带错误处理)
	db2, err := c2.SelectDB(2)
	if err != nil {
		log.Fatal(err)
	}
	db2.Set(ctx, "cache:hot", "data", 10*time.Minute)

	// GetClient 获取原生 *redis.Client(不带前缀)
	if rdb, ok := c2.GetClient(1); ok {
		rdb.Set(ctx, "raw:key", "value", 0)
	}

	// ─── 5. Hash 操作 ───
	c2.HSet(ctx, "user:100", "name", "alice", "age", 18)
	name, _ := c2.HGet(ctx, "user:100", "name").Result()
	fmt.Println(name) // "alice"

	all, _ := c2.HGetAll(ctx, "user:100").Result()
	fmt.Println(all) // map[name:alice age:18]

	// ─── 6. List 操作 ───
	c2.LPush(ctx, "queue:tasks", "task1", "task2", "task3")
	task, _ := c2.RPop(ctx, "queue:tasks").Result()
	fmt.Println(task) // "task1"

	// ─── 7. Set 操作 ───
	c2.SAdd(ctx, "tags:user:1", "go", "redis", "docker")
	isMember, _ := c2.SIsMember(ctx, "tags:user:1", "go").Result()
	fmt.Println(isMember) // true

	// ─── 8. Sorted Set 操作 ───
	c2.ZAdd(ctx, "leaderboard",
		goredis.Z{Score: 100, Member: "alice"},
		goredis.Z{Score: 200, Member: "bob"},
	)
	score, _ := c2.ZScore(ctx, "leaderboard", "bob").Result()
	fmt.Println(score) // 200

	// ─── 9. SCAN 安全批量删除 ───
	deleted, _ := c2.DelByPattern(ctx, "user:*")
	fmt.Printf("deleted %d keys\n", deleted)

	// ─── 10. 健康检查 ───
	if err := c2.HealthCheck(ctx); err != nil {
		log.Printf("redis unhealthy: %v", err)
	}

	// ─── 11. Lua 脚本(keys 自动加前缀)───
	script := `return redis.call("GET", KEYS[1])`
	result, _ := c2.Eval(ctx, script, []string{"user:1"}).Result()
	fmt.Println(result)

	// ─── 12. Pipeline(需通过 Proxy.Key 手动拼前缀)───
	proxy := c2.MustSelectDB(0)
	pipe := proxy.Pipeline()
	pipe.Set(ctx, proxy.Key("batch:1"), "v1", time.Hour)
	pipe.Set(ctx, proxy.Key("batch:2"), "v2", time.Hour)
	pipe.Get(ctx, proxy.Key("batch:1"))
	pipe.Del(ctx, proxy.Keys("batch:1", "batch:2")...) // 多 key 批量拼前缀
	_, _ = pipe.Exec(ctx)

	// ─── 13. Pub/Sub(channel 使用独立前缀,不复用 key 前缀)───
	pubsub := c2.Subscribe(ctx, "events:user")
	defer pubsub.Close()

	c2.Publish(ctx, "events:user", "user_created")
}
Example (Migration)

Example_migration 展示从 github.com/gtkit/redis(v1)迁移的等价用法。

v1 用法:

conn, err := redis.NewCollection(
    redis.WithAddr("127.0.0.1:6379"),
    redis.WithDB(0, "test"),
    redis.WithDB(1),
    redis.WithDB(2, "prefix:test2"),
)
rdb := redis.Select(2)
rdb.Client().Set(ctx, rdb.Prefix()+"key:2", "value:2", 0)  // 手动拼前缀

redisx 等价:.

package main

import (
	"context"
	"log"

	"github.com/gtkit/redisx"
)

func main() {
	c, err := redisx.NewClient(
		redisx.WithAddr("127.0.0.1:6379"),
		redisx.WithInitDBPrefix(0, "test"),
		redisx.WithInitDBs(1),
		redisx.WithInitDBPrefix(2, "prefix:test2"),
	)
	if err != nil {
		log.Fatal(err)
	}
	defer c.Close()

	ctx := context.Background()

	// redisx: 前缀自动拼接,业务层无感知
	db2 := c.MustSelectDB(2)
	db2.Set(ctx, "key:2", "value:2", 0) // 实际 key: "prefix:test2:key:2"
}

Index

Examples

Constants

View Source
const Version = "v1.3.0"

Version 是当前库版本号。

Variables

View Source
var ErrLockLost = errors.New("redisx: lock lost")

ErrLockLost 表示锁已不再被当前持有者持有(已过期或被他人重新获取)。

View Source
var ErrLockNoExpiry = errors.New("redisx: lock has no expiry")

ErrLockNoExpiry 表示锁存在但没有设置过期时间(契约外的永久锁)。

本库自身的获取/续期总会带 TTL,出现此错误通常意味着该 key 被外部 PERSIST 或以无过期方式写入,属死锁隐患,调用方应显式处理而非当作正常。

View Source
var ErrLockNotObtained = errors.New("redisx: lock not obtained")

ErrLockNotObtained 表示锁正被他人持有,本次未获得。

View Source
var ErrMessageDeadLettered = errors.New("redisx: message dead-lettered")

ErrMessageDeadLettered 表示消息投递次数超过 MaxDeliver,已被转入死信流。 经 StreamConfig.OnError 通知业务时可 errors.Is 判别。

Functions

func GetBytes

func GetBytes(ctx context.Context, p *Proxy, key string) ([]byte, error)

GetBytes reads raw bytes from key with automatic key prefixing.

Missing keys preserve redis.Nil through error wrapping.

func GetCodec

func GetCodec[T any](ctx context.Context, p *Proxy, key string, codec Codec) (T, error)

GetCodec reads key and unmarshals the stored bytes into T with codec.

Missing keys preserve redis.Nil through error wrapping.

func GetJSON

func GetJSON[T any](ctx context.Context, p *Proxy, key string) (T, error)

GetJSON 读取 key 的值并 JSON 反序列化为 T。key 自动拼接前缀。

key 不存在时返回 T 零值与 redis.Nil(errors.Is 判断);值不是合法 JSON 时返回携带 key 的错误。Go 方法不支持类型参数,因此为包级函数; Client 内嵌 Proxy,可直接作为 p 传入:

user, err := redisx.GetJSON[User](ctx, c.Proxy, "user:1")
if errors.Is(err, redis.Nil) { /* miss */ }
Example

ExampleGetJSON 演示结构体值的 JSON 读写助手与新增的人体工学 API。

package main

import (
	"context"
	"errors"
	"fmt"
	"log"
	"time"

	"github.com/gtkit/redisx"

	goredis "github.com/redis/go-redis/v9"
)

func main() {
	c, err := redisx.NewClient(redisx.WithAddr("127.0.0.1:6379"), redisx.WithKeyPrefix("app"))
	if err != nil {
		log.Fatal(err)
	}
	defer c.Close()

	ctx := context.Background()

	type User struct {
		Name string `json:"name"`
		Age  int    `json:"age"`
	}

	// 结构体读写:Client 内嵌 Proxy,直接传 c.Proxy
	_ = redisx.SetJSON(ctx, c.Proxy, "user:1", User{Name: "alice", Age: 18}, time.Hour)
	u, err := redisx.GetJSON[User](ctx, c.Proxy, "user:1")
	if errors.Is(err, goredis.Nil) {
		fmt.Println("user not found")
	}
	_ = u

	// 闭包持锁执行:拿锁—执行—必释放
	err = c.WithLock(ctx, "job:report", 30*time.Second, func(context.Context) error {
		return nil // 业务逻辑
	})
	if errors.Is(err, redisx.ErrLockNotObtained) {
		fmt.Println("他人持有,跳过本轮")
	}

	// 遍历 key:产出已剥前缀,可直接回传
	for key, err := range c.ScanKeys(ctx, "user:*") {
		if err != nil {
			log.Fatal(err)
		}
		c.Del(ctx, key)
	}
}

func SetBytes

func SetBytes(ctx context.Context, p *Proxy, key string, data []byte, expiration time.Duration) error

SetBytes writes data to key with automatic key prefixing.

expiration follows Redis SET semantics: 0 means no expiration.

func SetCodec

func SetCodec(ctx context.Context, p *Proxy, key string, codec Codec, value any, expiration time.Duration) error

SetCodec marshals value with codec and writes the resulting bytes to key.

The codec is caller-provided; redisx does not depend on protobuf, gob, msgpack, or other encoding packages.

func SetJSON

func SetJSON(ctx context.Context, p *Proxy, key string, value any, expiration time.Duration) error

SetJSON 将 value JSON 序列化后写入 key。key 自动拼接前缀, expiration 为 0 表示不设置过期时间。

序列化失败返回携带 key 的错误且不发送命令。

Types

type Client

type Client struct {
	*Proxy // 默认 DB 的命令代理,方法提升为 Client 级快捷方式
	// contains filtered or unexported fields
}

Client 是 Redis 多 DB 客户端封装。

嵌入的 *Proxy 是默认 DB 的命令代理,其全部命令方法(Get/Set/HSet/ ConsumeStream/TryLock 等)经方法提升直接在 Client 上可用,等价于 MustSelectDB(defaultDB) 上的同名调用,key 自动拼接前缀。

并发安全:内部 clients / proxies map 在 NewClient 中一次性构建完成后即为只读, 不再修改,因此并发读取无需加锁——这是比 sync.RWMutex 更轻量的方案。 底层每个 *redis.Client 本身也是并发安全的(go-redis 原生连接池)。

func NewClient

func NewClient(opts ...Option) (*Client, error)

NewClient 使用 Functional Options 创建 Redis 客户端。

等价于以 context.Background() 调用 NewClientContext;需要约束或取消 初始化耗时的场景请直接使用 NewClientContext

用法:

c, err := redisx.NewClient(
    redisx.WithAddr("127.0.0.1:6379"),
    redisx.WithPassword("secret"),
    redisx.WithKeyPrefix("myapp"),
    redisx.WithInitDBs(0, 1, 2),
)

func NewClientContext

func NewClientContext(ctx context.Context, opts ...Option) (*Client, error)

NewClientContext 与 NewClient 相同,但初始化拨号与 Ping 均受传入 ctx 约束: ctx 取消或超时立即中止初始化并回收已建连接。

初始化对每个 DB 执行 Ping 检查连通性:DefaultDB 优先拨号,失败立即整体失败且 不再拨号其余 DB;其余 DB 并发拨号以缩短启动时间。非降级模式下任一失败即返回 确定性错误(编号最小的失败 DB)并清理所有已创建的连接。启用 WithAllowPartialInit 后改为降级语义:失败的 DB 缺席集合,错误以 *InitError 返回(与可用的 Client 可同时非 nil,可经 errors.As 按 DB 提取失败原因,且与 DB 编号确定对应、不受并发 完成顺序影响),但 DefaultDB 仍必须初始化成功,否则整体失败返回 nil。 clients/proxies map 在构建完成后不再修改,后续并发读取无需加锁。

func (*Client) AddHook

func (c *Client) AddHook(hook redis.Hook)

AddHook 将 hook 安装到所有已初始化 DB 的底层 *redis.Client 上。

用于一次性挂接熔断、限流、metrics、tracing 等中间件,库本身不内置任何策略。 hook 为 nil 时静默忽略。部分初始化(WithAllowPartialInit)模式下, 初始化失败而缺席的 DB 自然跳过。

调用时机:与 go-redis 原生 AddHook 的约束一致,必须在 NewClient 之后、 开始并发执行命令之前完成安装,运行中途添加不保证并发安全。 如需只对个别 DB 安装,请改用 Client.GetClient 自行处理。

func (*Client) Close

func (c *Client) Close() error

Close 优雅关闭所有 DB 客户端连接。

如有多个 DB 关闭失败,使用 errors.Join 合并返回。

func (*Client) DefaultClient

func (c *Client) DefaultClient() *redis.Client

DefaultClient 返回默认 DB 的底层 *redis.Client

func (*Client) GetClient

func (c *Client) GetClient(db int) (*redis.Client, bool)

GetClient 安全获取指定 DB 的底层 *redis.Client

返回 false 表示该 DB 未初始化。用于需要直接操作 go-redis 原生 API 的场景。

func (*Client) HealthCheck

func (c *Client) HealthCheck(ctx context.Context) error

HealthCheck 对所有已初始化的 DB 并发执行 PING 健康检查。

不会短路:即使某个 DB 失败也会检查其余 DB,最终以 errors.Join 返回所有失败的 聚合错误(按 DB 编号确定有序,不受并发完成顺序影响)。返回 nil 表示全部健康。

func (*Client) MustSelectDB

func (c *Client) MustSelectDB(db int) *Proxy

MustSelectDB 返回指定 DB 编号上的命令代理 Proxy

如果 DB 未初始化,直接 panic。仅用于程序启动阶段或确定 DB 存在的场景。

func (*Client) PoolStats

func (c *Client) PoolStats() map[int]*redis.PoolStats

PoolStats 返回每个已初始化 DB 的连接池统计,key 为 DB 编号。

透传 go-redis 的连接池统计(命中/未命中/超时/连接数等),供业务接入 监控拉取;库本身不做任何聚合、阈值或告警判断。

func (*Client) Prefix

func (c *Client) Prefix() string

Prefix 返回当前配置的全局 key 前缀。

func (*Client) SelectDB

func (c *Client) SelectDB(db int) (*Proxy, error)

SelectDB 返回指定 DB 编号上的命令代理 Proxy,支持链式调用。

如果指定的 DB 未在 WithInitDBs / WithInitDBPrefix 中初始化,返回错误。

type Codec

type Codec interface {
	Marshal(v any) ([]byte, error)
	Unmarshal(data []byte, v any) error
}

Codec defines caller-provided value encoding for SetCodec and GetCodec.

It is intentionally small so protobuf, gob, msgpack, encrypted payloads, or application-specific codecs can be adapted without redisx depending on them.

type Config

type Config struct {
	// Addr 是 Redis 服务器地址,格式为 host:port。必填。
	Addr string

	// Username 用于 Redis 6+ ACL 认证。空字符串表示不使用。
	Username string

	// Password 是 Redis 认证密码。空字符串表示无需认证。
	Password string

	// DefaultDB 是默认使用的数据库编号,默认 0。
	// 上限取决于服务端 databases 配置(默认 16 个库即 0~15),越界在初始化 Ping 时报错。
	DefaultDB int

	// InitDBs 是初始化时需要创建连接的 DB 配置列表。
	// DefaultDB 会自动包含,无需重复添加。
	InitDBs []DBConfig

	// PoolSize 是每个 DB 客户端的最大连接池大小。默认 10。
	PoolSize int

	// MinIdleConns 是连接池中保持的最小空闲连接数。默认 3。
	MinIdleConns int

	// MaxRetries 是命令失败后最大重试次数。默认 3。
	MaxRetries int

	// DialTimeout 是建立 TCP 连接的超时时间。默认 5s。
	DialTimeout time.Duration

	// ReadTimeout 是 socket 读操作超时时间。默认 3s。
	ReadTimeout time.Duration

	// WriteTimeout 是 socket 写操作超时时间。默认 3s。
	WriteTimeout time.Duration

	// IdleTimeout 是空闲连接被回收前的最大存活时间。默认 5m。
	IdleTimeout time.Duration

	// KeyPrefix 是全局 key 前缀。
	// 设置后所有命令的 key 会自动拼接为 "{KeyPrefix}{KeyPrefixSeparator}{key}"。
	// 可被 DBConfig.Prefix 覆盖。空字符串表示不使用前缀。
	KeyPrefix string

	// KeyPrefixSeparator 是 key 前缀与原始 key 之间的连接符。
	// 默认 ":";仅当前缀非空时参与拼接。
	KeyPrefixSeparator string

	// ChannelPrefix 是 Pub/Sub channel 前缀。
	// Redis Pub/Sub channel 是实例级全局命名空间,不随 DB 切换;默认不使用 key 前缀。
	ChannelPrefix string

	// ChannelPrefixSeparator 是 channel 前缀与原始 channel 之间的连接符。
	// 默认 ":";仅当 ChannelPrefix 非空时参与拼接。
	ChannelPrefixSeparator string

	// TLSConfig 是连接 Redis 的 TLS 配置。nil 表示不启用 TLS。
	TLSConfig *tls.Config

	// AllowPartialInit 允许部分 DB 初始化失败(降级模式)。
	// 默认 false:任一 DB 失败则整体失败(全有或全无)。
	AllowPartialInit bool
}

Config 定义 Redis 客户端的完整配置。 所有字段均有合理默认值,仅 Addr 为必填项。

type DBConfig

type DBConfig struct {
	DB     int
	Prefix string // 可选,非空时覆盖全局 KeyPrefix
}

DBConfig 定义单个 DB 的配置(编号 + 可选独立前缀)。

当 Prefix 非空时,该 DB 使用独立前缀替代全局 KeyPrefix。 兼容 v1 的 WithDB(0, "myprefix") 用法。

type FencedLock

type FencedLock struct {
	Lock
	// contains filtered or unexported fields
}

FencedLock 表示一把已持有、带单调递增 fencing token 的分布式锁,由 Proxy.FencedLock 获取。它嵌入 Lock,故 Key / Release / Refresh / TTL 语义与 Lock 完全一致(token 校验、防误删、过期返回 ErrLockLost);额外通过 FencedLock.Fence 暴露本次获取的 fencing token。

与 Lock 相同,这是单 Redis 实例锁(非 RedLock)。

func (*FencedLock) Fence

func (l *FencedLock) Fence() int64

Fence 返回本次获取的 fencing token。同一 key 上后获得者的 token 严格大于 先获得者(可能有间隙,但不回退、不重复),供下游被保护资源拒绝较旧持有者。

type InitError

type InitError struct {
	// Failed 记录每个初始化失败的 DB 编号及其失败原因。
	Failed map[int]error
}

InitError 聚合降级初始化(WithAllowPartialInit)下各 DB 的失败详情。

通过 errors.As 提取后可按 DB 编程决策:

var ie *redisx.InitError
if errors.As(err, &ie) {
    if _, bad := ie.Failed[2]; bad {
        // DB2 不可用,按业务决定降级或拒绝启动
    }
}

func (*InitError) Error

func (e *InitError) Error() string

Error 按 DB 编号升序稳定输出所有失败原因。

func (*InitError) Unwrap

func (e *InitError) Unwrap() []error

Unwrap 按 DB 编号升序返回各失败的底层错误, 使 errors.Is / errors.As 可穿透到网络层错误。

type Lock

type Lock struct {
	// contains filtered or unexported fields
}

Lock 表示一把已持有的单实例分布式锁,由 Proxy.TryLock 获取。

注意语义边界:这是单 Redis 实例锁(非 RedLock),主从异步复制下 故障切换瞬间存在双持有的理论窗口,关键互斥请在业务层做幂等兜底。

func (*Lock) Key

func (l *Lock) Key() string

Key 返回已拼前缀的完整锁 key,供日志、打点等观测场景使用。

func (*Lock) Refresh

func (l *Lock) Refresh(ctx context.Context, ttl time.Duration) error

Refresh 将锁的 TTL 重设为 ttl(校验 token 后原子执行)。

长任务在持有期间显式调用本方法续期;ttl 必须至少为 1ms;锁已失去返回 ErrLockLost

func (*Lock) Release

func (l *Lock) Release(ctx context.Context) error

Release 释放锁。仅当锁仍被当前持有者持有(token 匹配)时删除; 锁已过期或被他人重新获取时返回 ErrLockLost,且不会影响他人的锁。

func (*Lock) TTL

func (l *Lock) TTL(ctx context.Context) (time.Duration, error)

TTL 返回锁的剩余存活时间(校验 token 后原子读取)。

返回 nil error 即表示锁仍被当前持有者持有;锁已过期或被他人重新获取 返回 ErrLockLost;token 匹配但锁无过期时间(契约外的永久锁)返回 ErrLockNoExpiry。仅供观测与业务自检:查询与后续操作之间锁仍可能 过期,互斥正确性依赖 Release/Refresh 自身的 token 校验,而非本方法。

type Option

type Option func(*Config)

Option 是 Functional Options 模式的配置函数。

func WithAddr

func WithAddr(addr string) Option

WithAddr 设置 Redis 服务器地址(必填)。

格式: "host:port",例如 "127.0.0.1:6379"。

func WithAllowPartialInit

func WithAllowPartialInit() Option

WithAllowPartialInit 允许部分 DB 初始化失败(降级模式)。

启用后,初始化失败的 DB 不进入集合,所有失败聚合后与可用的 Client 一同返回(两者可同时非 nil),对缺席 DB 调用 Client.SelectDB 返回错误; DefaultDB 承载 Client 级快捷方法,仍必须初始化成功,否则整体失败。 默认关闭,即任一 DB 失败则整体失败(全有或全无)。

func WithChannelPrefix

func WithChannelPrefix(prefix string) Option

WithChannelPrefix 设置 Pub/Sub channel 前缀。

Redis Pub/Sub channel 是实例级全局命名空间,不随 DB 切换;默认不复用 key 前缀。 设置后 Pub/Sub 方法会将 channel 拼接为 "{prefix}{separator}{channel}"。

func WithChannelPrefixSeparator

func WithChannelPrefixSeparator(separator string) Option

WithChannelPrefixSeparator 设置 Pub/Sub channel 前缀与原始 channel 之间的连接符。

默认 ":"。仅当 channel 前缀非空时参与拼接;传入空字符串表示直接连接前缀和 channel。

func WithDB

func WithDB(db int) Option

WithDB 设置默认使用的数据库编号。

编号必须 >= 0;上限取决于服务端 databases 配置(默认 16 个库即 0~15), 越界的编号在 NewClient 初始化 Ping 时由 Redis 报错。

func WithDBConfig deprecated

func WithDBConfig(db int, prefix string) Option

WithDBConfig 添加一个带独立前缀的 DB 配置。

Deprecated: use WithInitDBPrefix.

func WithDialTimeout

func WithDialTimeout(d time.Duration) Option

WithDialTimeout 设置建立 TCP 连接的超时时间。

func WithIdleTimeout

func WithIdleTimeout(d time.Duration) Option

WithIdleTimeout 设置空闲连接被回收前的最大存活时间。

func WithInitDBPrefix

func WithInitDBPrefix(db int, prefix string) Option

WithInitDBPrefix 初始化指定 DB,并为该 DB 设置独立 key 前缀。

当 prefix 非空时,该 DB 使用独立前缀替代全局 KeyPrefix。 兼容 v1 的 WithDB(db, "prefix") 语义。

示例: WithInitDBPrefix(2, "session") 使 DB2 的 key 前缀为 "session:" 而非全局前缀。 前缀连接符可通过 WithKeyPrefixSeparator 全局配置。

func WithInitDBs

func WithInitDBs(dbs ...int) Option

WithInitDBs 设置需要初始化的多个 DB 编号。

DefaultDB 会自动包含在列表中,无需重复添加。 所有 DB 共享全局 KeyPrefix。如需 per-DB 前缀,请使用 WithInitDBPrefix

示例: WithInitDBs(0, 1, 2) 将同时初始化 DB0、DB1、DB2。

func WithKeyPrefix

func WithKeyPrefix(prefix string) Option

WithKeyPrefix 设置全局 key 前缀。

设置后所有带 key 参数的命令会自动拼接为 "{prefix}{separator}{key}", 对业务层完全透明。可被 per-DB 前缀覆盖(见 WithInitDBPrefix)。

func WithKeyPrefixSeparator

func WithKeyPrefixSeparator(separator string) Option

WithKeyPrefixSeparator 设置 key 前缀与原始 key 之间的连接符。

默认 ":"。仅当前缀非空时参与拼接;传入空字符串表示直接连接前缀和 key。

func WithMaxRetries

func WithMaxRetries(n int) Option

WithMaxRetries 设置命令失败后最大重试次数。

n 为 0 表示关闭自动重试(库内会映射为 go-redis 的 -1 哨兵值), 负值在 NewClient 返回错误。

注意:对非幂等命令(如 INCR、LPUSH),读超时后的自动重试可能导致 命令被重复执行。对此类命令敏感的场景请设置为 0 关闭重试, 或在业务层使用 Lua 脚本保证幂等。

func WithMinIdleConns

func WithMinIdleConns(n int) Option

WithMinIdleConns 设置连接池中保持的最小空闲连接数。

func WithPassword

func WithPassword(password string) Option

WithPassword 设置 Redis 认证密码。

func WithPoolSize

func WithPoolSize(size int) Option

WithPoolSize 设置每个 DB 连接池的最大连接数。

func WithReadTimeout

func WithReadTimeout(d time.Duration) Option

WithReadTimeout 设置 socket 读操作超时时间。

func WithTLSConfig

func WithTLSConfig(tlsConfig *tls.Config) Option

WithTLSConfig 设置连接 Redis 的 TLS 配置。

传入非 nil 配置后,所有 DB 的连接均通过 TLS 建立, 适用于云厂商强制加密的 Redis 实例。nil 表示不启用 TLS。

func WithUsername

func WithUsername(username string) Option

WithUsername 设置 Redis 6+ ACL 认证用户名。

func WithWriteTimeout

func WithWriteTimeout(d time.Duration) Option

WithWriteTimeout 设置 socket 写操作超时时间。

type Proxy

type Proxy struct {
	// contains filtered or unexported fields
}

Proxy 是对单个 DB 的 *redis.Client 命令代理。

所有带 key 参数的方法会自动拼接前缀(全局或 per-DB),对业务层透明。 Pub/Sub channel 使用独立 channel 前缀;默认不复用 key 前缀。 通过 Client.SelectDB / Client.MustSelectDB 获取。

Proxy 在 NewClient 中按 DB 缓存,不会重复创建。

func WrapClient

func WrapClient(rdb *redis.Client, opts ...ProxyOption) (*Proxy, error)

WrapClient 基于已有 *redis.Client 构造 Proxy

返回的 Proxy 不接管 rdb 生命周期,调用方仍负责关闭传入的 client。

Example
package main

import (
	"fmt"
	"log"

	"github.com/gtkit/redisx"

	goredis "github.com/redis/go-redis/v9"
)

func main() {
	rdb := goredis.NewClient(&goredis.Options{Addr: "127.0.0.1:1"})
	defer rdb.Close()

	p, err := redisx.WrapClient(
		rdb,
		redisx.WithProxyPrefix("app"),
		redisx.WithProxyPrefixSeparator(":"),
	)
	if err != nil {
		log.Fatal(err)
	}

	fmt.Println(p.Key("user:1"))
	fmt.Println(p.Keys("a", "b"))
}
Output:
app:user:1
[app:a app:b]

func (*Proxy) Consume

func (p *Proxy) Consume(ctx context.Context, handler func(*redis.Message), channels ...string) error

Consume 以受管方式订阅频道并串行消费消息,阻塞直到 ctx 取消或出错。

订阅生命周期由库内管理:订阅确认失败立即返回错误,退出时自动关闭订阅; ctx 取消时返回 ctx 的错误(可用 errors.Is(err, context.Canceled) 判断优雅退出); handler panic 会被恢复,消费终止并以 error 返回。

handler 串行执行以保证单频道消息顺序,耗时处理请在业务侧自行分发。 注意 Redis Pub/Sub 为 at-most-once,断线期间的消息会丢失; 需要可靠投递请使用 Stream。所有 channel 使用独立 channel 前缀。

func (*Proxy) ConsumePattern

func (p *Proxy) ConsumePattern(ctx context.Context, handler func(*redis.Message), patterns ...string) error

ConsumePattern 以受管方式按模式订阅(PSUBSCRIBE)并串行消费消息, 生命周期语义与 Proxy.Consume 完全一致。所有模式使用独立 channel 前缀。

func (*Proxy) ConsumeStream

func (p *Proxy) ConsumeStream(ctx context.Context, cfg StreamConfig, handler func(redis.XMessage) error) error

ConsumeStream 以消费组方式受管消费流消息(at-least-once), 阻塞直到 ctx 取消或发生不可恢复错误。

生命周期由库内管理:消费组不存在时自动创建(MKSTREAM);启动时先续传 本消费者的 pending 消息(崩溃重启不丢已读未确认的消息),再消费新消息。

handler 返回 nil 即 XACK;返回 error 时消息留在 pending(重启续传或被 AutoClaim 接管时重投),消费继续,可经 StreamConfig.OnError 感知失败明细; handler panic 会被恢复,消费终止并以 error 返回。handler 串行执行以保证消息顺序。 ctx 取消时返回 ctx 的错误(errors.Is(err, context.Canceled) 判断优雅退出)。

持续失败的消息默认堆积在 pending(请业务侧用 Proxy.XPending 监控); 配置 StreamConfig.MaxDeliver + StreamConfig.DeadLetterStream 可在 投递超限后将毒消息原子隔离到死信流,不再阻塞重投。

库内不提供 handler 超时控制(Go 无法安全中断不配合的函数),长耗时处理 请在 handler 内自行用 context 控制;StreamConfig.OnError 回调中如需 stream/group/consumer 上下文,闭包捕获自己构造的 StreamConfig 即可, 消息 ID 在回调参数 msg.ID 中。

func (*Proxy) Decr

func (p *Proxy) Decr(ctx context.Context, key string) *redis.IntCmd

Decr 将 key 的整数值减 1。key 自动拼接前缀。

func (*Proxy) DecrBy

func (p *Proxy) DecrBy(ctx context.Context, key string, value int64) *redis.IntCmd

DecrBy 将 key 的整数值减少指定减量。key 自动拼接前缀。

func (*Proxy) Del

func (p *Proxy) Del(ctx context.Context, keys ...string) *redis.IntCmd

Del 删除一个或多个 key。所有 key 自动拼接前缀。返回成功删除的数量。

func (*Proxy) DelByPattern

func (p *Proxy) DelByPattern(ctx context.Context, pattern string) (int64, error)

DelByPattern 使用 SCAN 安全批量删除匹配 pattern 的 key。

pattern 自动拼接前缀:如 pattern="user:*",实际匹配 "{prefix}:user:*"。 内部使用 SCAN + UNLINK 批量删除:UNLINK 在后台异步释放内存, 大 value 场景不会阻塞 Redis 主线程(需要 Redis >= 4.0)。

建议调用方传入带超时的 ctx 控制执行时间。返回成功删除的 key 总数。

func (*Proxy) Eval

func (p *Proxy) Eval(ctx context.Context, script string, keys []string, args ...any) *redis.Cmd

Eval 执行 Lua 脚本。keys 列表中的每个 key 自动拼接前缀。

func (*Proxy) EvalScript

func (p *Proxy) EvalScript(ctx context.Context, script *redis.Script, keys []string, args ...any) *redis.Cmd

EvalScript 执行 *redis.Script。keys 列表中的每个 key 自动拼接前缀。

func (*Proxy) EvalSha

func (p *Proxy) EvalSha(ctx context.Context, sha1 string, keys []string, args ...any) *redis.Cmd

EvalSha 执行已缓存的 Lua 脚本。keys 列表中的每个 key 自动拼接前缀。

func (*Proxy) Exists

func (p *Proxy) Exists(ctx context.Context, keys ...string) *redis.IntCmd

Exists 检查一个或多个 key 是否存在。所有 key 自动拼接前缀。返回存在的数量。

func (*Proxy) Expire

func (p *Proxy) Expire(ctx context.Context, key string, expiration time.Duration) *redis.BoolCmd

Expire 设置 key 的过期时间。key 自动拼接前缀。

func (*Proxy) ExpireAt

func (p *Proxy) ExpireAt(ctx context.Context, key string, tm time.Time) *redis.BoolCmd

ExpireAt 设置 key 在指定时间点过期。key 自动拼接前缀。

func (*Proxy) FencedLock

func (p *Proxy) FencedLock(ctx context.Context, key string, ttl time.Duration) (*FencedLock, error)

FencedLock 尝试获取带 fencing token 的分布式锁,非阻塞。key 自动拼接前缀。

相比 Proxy.TryLock,额外返回一个在同一 key 上严格单调递增的 fencing token (见 FencedLock.Fence):进程暂停 / GC / 网络阻塞导致旧持有者的锁在 TTL 过期、新持有者已拿到锁后,旧持有者仍可能继续操作外部资源;随机 token 只能防 误删他人锁,无法阻止这种"过期后仍写"。fencing token 的能力是——由下游被保护 资源记录见过的最大 token 并拒绝更小者,从而排除旧持有者。**仅当下游确实校验 token 时栅栏才生效**,本库只负责原子生成与暴露 token。

成功返回 *FencedLock;锁被他人持有返回 ErrLockNotObtained(errors.Is 判断); ttl 必须至少为 1ms。获取路径与 TryLock 一样对底层自动重试幂等。

实现说明:fencing token 由一个持久(无 TTL)的计数器 key 提供(锁 key 追加 ":__fence__" 后缀),Release 不删除该计数器以保证跨获取单调。该计数器 key 会 长期存在,且可能被 Proxy.DelByPattern 的宽 pattern 扫到,请勿手动删除, 否则单调性重置。

Example

ExampleProxy_FencedLock 演示带 fencing token 的锁:token 单调递增, 供下游被保护资源拒绝较旧持有者(仅当下游校验 token 时栅栏才生效)。

package main

import (
	"context"
	"errors"
	"fmt"
	"log"
	"time"

	"github.com/gtkit/redisx"
)

func main() {
	c, err := redisx.NewClient(redisx.WithAddr("127.0.0.1:6379"), redisx.WithKeyPrefix("app"))
	if err != nil {
		log.Fatal(err)
	}
	defer c.Close()

	ctx := context.Background()

	lock, err := c.FencedLock(ctx, "resource:42", 30*time.Second)
	if errors.Is(err, redisx.ErrLockNotObtained) {
		fmt.Println("他人持有,跳过本轮")
		return
	}
	if err != nil {
		log.Fatal(err)
	}
	defer lock.Release(ctx)

	// 把 fence 传给下游资源;下游记录见过的最大 token 并拒绝更小者
	fence := lock.Fence()
	_ = fence
}

func (*Proxy) Get

func (p *Proxy) Get(ctx context.Context, key string) *redis.StringCmd

Get 获取 key 的值。key 自动拼接前缀。

key 不存在返回 redis.Nil,可通过 errors.Is(err, redis.Nil) 判断。

func (*Proxy) GetDel

func (p *Proxy) GetDel(ctx context.Context, key string) *redis.StringCmd

GetDel 获取 key 的值并删除该 key。key 自动拼接前缀。

func (*Proxy) GetSet

func (p *Proxy) GetSet(ctx context.Context, key string, value any) *redis.StringCmd

GetSet 设置新值并返回旧值。key 自动拼接前缀。

内部使用 SET ... GET 实现(GETSET 的现代等价形式,需要 Redis >= 6.2); 旧值不存在时返回 redis.Nil,新值仍会写入。

func (*Proxy) HDel

func (p *Proxy) HDel(ctx context.Context, key string, fields ...string) *redis.IntCmd

HDel 删除 hash 中一个或多个字段。key 自动拼接前缀。

func (*Proxy) HExists

func (p *Proxy) HExists(ctx context.Context, key, field string) *redis.BoolCmd

HExists 检查 hash 中指定字段是否存在。key 自动拼接前缀。

func (*Proxy) HGet

func (p *Proxy) HGet(ctx context.Context, key, field string) *redis.StringCmd

HGet 获取 hash 中指定字段的值。key 自动拼接前缀。

func (*Proxy) HGetAll

func (p *Proxy) HGetAll(ctx context.Context, key string) *redis.MapStringStringCmd

HGetAll 获取 hash 中所有字段和值。key 自动拼接前缀。

func (*Proxy) HIncrBy

func (p *Proxy) HIncrBy(ctx context.Context, key, field string, incr int64) *redis.IntCmd

HIncrBy 将 hash 中指定字段的整数值增加增量。key 自动拼接前缀。

func (*Proxy) HIncrByFloat

func (p *Proxy) HIncrByFloat(ctx context.Context, key, field string, incr float64) *redis.FloatCmd

HIncrByFloat 将 hash 中指定字段的浮点数值增加增量。key 自动拼接前缀。

func (*Proxy) HLen

func (p *Proxy) HLen(ctx context.Context, key string) *redis.IntCmd

HLen 返回 hash 中字段数量。key 自动拼接前缀。

func (*Proxy) HMGet

func (p *Proxy) HMGet(ctx context.Context, key string, fields ...string) *redis.SliceCmd

HMGet 批量获取 hash 中多个字段的值。key 自动拼接前缀。

func (*Proxy) HMSet deprecated

func (p *Proxy) HMSet(ctx context.Context, key string, values ...any) *redis.BoolCmd

HMSet 批量设置 hash 字段。key 自动拼接前缀。

Deprecated: Redis 官方建议使用 Proxy.HSet,HMSet 保留用于兼容。

func (*Proxy) HSet

func (p *Proxy) HSet(ctx context.Context, key string, values ...any) *redis.IntCmd

HSet 设置 hash 中一个或多个字段。key 自动拼接前缀。

values 接受 field-value 对: HSet(ctx, "user:1", "name", "alice", "age", 18).

func (*Proxy) Incr

func (p *Proxy) Incr(ctx context.Context, key string) *redis.IntCmd

Incr 将 key 的整数值加 1。key 自动拼接前缀。

func (*Proxy) IncrBy

func (p *Proxy) IncrBy(ctx context.Context, key string, value int64) *redis.IntCmd

IncrBy 将 key 的整数值增加指定增量。key 自动拼接前缀。

func (*Proxy) IncrByFloat

func (p *Proxy) IncrByFloat(ctx context.Context, key string, value float64) *redis.FloatCmd

IncrByFloat 将 key 的浮点数值增加指定增量。key 自动拼接前缀。

func (*Proxy) Key

func (p *Proxy) Key(k string) string

Key 返回拼接了前缀的完整 key。

用于在 Pipeline 等需要手动拼接前缀的场景。

func (*Proxy) Keys

func (p *Proxy) Keys(ks ...string) []string

Keys 返回逐个拼接了前缀的完整 key 列表。

用于在 Pipeline 等需要手动拼接前缀的场景批量处理多个 key, 例如 pipe.Del(ctx, proxy.Keys("k1", "k2")...)。

func (*Proxy) LIndex

func (p *Proxy) LIndex(ctx context.Context, key string, index int64) *redis.StringCmd

LIndex 返回列表中指定索引的元素。key 自动拼接前缀。

func (*Proxy) LLen

func (p *Proxy) LLen(ctx context.Context, key string) *redis.IntCmd

LLen 返回列表长度。key 自动拼接前缀。

func (*Proxy) LPop

func (p *Proxy) LPop(ctx context.Context, key string) *redis.StringCmd

LPop 从列表左侧弹出一个元素。key 自动拼接前缀。

func (*Proxy) LPush

func (p *Proxy) LPush(ctx context.Context, key string, values ...any) *redis.IntCmd

LPush 从列表左侧推入一个或多个元素。key 自动拼接前缀。

func (*Proxy) LRange

func (p *Proxy) LRange(ctx context.Context, key string, start, stop int64) *redis.StringSliceCmd

LRange 返回列表中指定范围的元素。key 自动拼接前缀。

start 和 stop 为从零开始的索引,-1 表示最后一个元素。

func (*Proxy) LRem

func (p *Proxy) LRem(ctx context.Context, key string, count int64, value any) *redis.IntCmd

LRem 移除列表中与 value 相等的元素。key 自动拼接前缀。

func (*Proxy) LTrim

func (p *Proxy) LTrim(ctx context.Context, key string, start, stop int64) *redis.StatusCmd

LTrim 只保留列表中指定范围的元素。key 自动拼接前缀。

func (*Proxy) MGet

func (p *Proxy) MGet(ctx context.Context, keys ...string) *redis.SliceCmd

MGet 批量获取多个 key 的值。所有 key 自动拼接前缀。

func (*Proxy) MSet

func (p *Proxy) MSet(ctx context.Context, values ...any) *redis.StatusCmd

MSet 批量设置 key-value 对。

values 为交替的 key-value 序列: MSet(ctx, "k1", "v1", "k2", "v2")。 偶数位(0, 2, 4...)必须为 string,作为 key 自动拼接前缀。

values 长度必须为偶数且 key 位必须为 string,否则返回错误。

func (*Proxy) PSubscribe

func (p *Proxy) PSubscribe(ctx context.Context, patterns ...string) *redis.PubSub

PSubscribe 按模式订阅一个或多个 channel 模式(如 "events:*")。 所有模式使用独立 channel 前缀。

返回 *redis.PubSub,调用方负责关闭。

func (*Proxy) PTTL

func (p *Proxy) PTTL(ctx context.Context, key string) *redis.DurationCmd

PTTL 返回 key 的剩余生存时间(毫秒精度)。key 自动拼接前缀。

func (*Proxy) Persist

func (p *Proxy) Persist(ctx context.Context, key string) *redis.BoolCmd

Persist 移除 key 的过期时间。key 自动拼接前缀。

func (*Proxy) Pipeline

func (p *Proxy) Pipeline() redis.Pipeliner

Pipeline 返回 go-redis 原生 Pipeline。

Pipeline 内的命令需要手动拼接前缀,可通过 Proxy.Key 获取带前缀的 key。

func (*Proxy) Publish

func (p *Proxy) Publish(ctx context.Context, channel string, message any) *redis.IntCmd

Publish 向指定 channel 发布消息。channel 使用独立 channel 前缀,不复用 key 前缀。

func (*Proxy) RPop

func (p *Proxy) RPop(ctx context.Context, key string) *redis.StringCmd

RPop 从列表右侧弹出一个元素。key 自动拼接前缀。

func (*Proxy) RPush

func (p *Proxy) RPush(ctx context.Context, key string, values ...any) *redis.IntCmd

RPush 从列表右侧推入一个或多个元素。key 自动拼接前缀。

func (*Proxy) RawClient

func (p *Proxy) RawClient() *redis.Client

RawClient 返回底层 *redis.Client

注意:通过 RawClient 执行的命令不会自动添加 key 前缀。

func (*Proxy) Rename

func (p *Proxy) Rename(ctx context.Context, key, newkey string) *redis.StatusCmd

Rename 重命名 key。两个 key 均自动拼接前缀。

func (*Proxy) SAdd

func (p *Proxy) SAdd(ctx context.Context, key string, members ...any) *redis.IntCmd

SAdd 向集合添加一个或多个成员。key 自动拼接前缀。

func (*Proxy) SCard

func (p *Proxy) SCard(ctx context.Context, key string) *redis.IntCmd

SCard 返回集合中成员的数量。key 自动拼接前缀。

func (*Proxy) SIsMember

func (p *Proxy) SIsMember(ctx context.Context, key string, member any) *redis.BoolCmd

SIsMember 判断 member 是否为集合的成员。key 自动拼接前缀。

func (*Proxy) SMembers

func (p *Proxy) SMembers(ctx context.Context, key string) *redis.StringSliceCmd

SMembers 返回集合中所有成员。key 自动拼接前缀。

func (*Proxy) SPop

func (p *Proxy) SPop(ctx context.Context, key string) *redis.StringCmd

SPop 随机移除并返回集合中一个成员。key 自动拼接前缀。

func (*Proxy) SRandMember

func (p *Proxy) SRandMember(ctx context.Context, key string) *redis.StringCmd

SRandMember 随机返回集合中一个成员。key 自动拼接前缀。

func (*Proxy) SRem

func (p *Proxy) SRem(ctx context.Context, key string, members ...any) *redis.IntCmd

SRem 移除集合中一个或多个成员。key 自动拼接前缀。

func (*Proxy) Scan

func (p *Proxy) Scan(ctx context.Context, cursor uint64, match string, count int64) *redis.ScanCmd

Scan 包装 SCAN 命令,match pattern 自动拼接前缀。

注意:返回的 key 是已含前缀的完整 key,直接回传给本库其他带前缀方法 (如 Proxy.Del)会造成二次拼接。如需对结果继续操作,请改用 Proxy.RawClient 执行,或自行剥离前缀;批量删除场景请直接使用 Proxy.DelByPattern

func (*Proxy) ScanKeys

func (p *Proxy) ScanKeys(ctx context.Context, match string) iter.Seq2[string, error]

ScanKeys 返回自动翻页的 key 迭代器,match 自动拼接前缀, 产出的 key **已剥去前缀**,可直接回传本库其他带前缀方法:

for key, err := range p.ScanKeys(ctx, "user:*") {
    if err != nil {
        return err
    }
    p.Del(ctx, key) // key 不含前缀,不会二次拼接
}

SCAN 出错时迭代产出一次非 nil error 后终止;提前 break 立即停止翻页。 一致性语义跟随 Redis SCAN:迭代期间新增/删除的 key 不保证快照视图。

func (*Proxy) Set

func (p *Proxy) Set(ctx context.Context, key string, value any, expiration time.Duration) *redis.StatusCmd

Set 设置 key-value。expiration 为 0 表示不设置过期时间。key 自动拼接前缀。

func (*Proxy) SetEX

func (p *Proxy) SetEX(ctx context.Context, key string, value any, expiration time.Duration) *redis.StatusCmd

SetEX 设置 key-value 并指定过期时间。key 自动拼接前缀。

内部使用 SET 带过期实现(SETEX 的现代等价形式); expiration 为 0 时等同于无过期时间的 SET。

func (*Proxy) SetNX

func (p *Proxy) SetNX(ctx context.Context, key string, value any, expiration time.Duration) *redis.BoolCmd

SetNX 仅当 key 不存在时设置。key 自动拼接前缀。

返回 true 表示设置成功,false 表示 key 已存在。

func (*Proxy) Subscribe

func (p *Proxy) Subscribe(ctx context.Context, channels ...string) *redis.PubSub

Subscribe 订阅一个或多个 channel。所有 channel 使用独立 channel 前缀。

返回 *redis.PubSub,调用方负责关闭;如需自动管理订阅生命周期, 请使用 Proxy.Consume

func (*Proxy) TTL

func (p *Proxy) TTL(ctx context.Context, key string) *redis.DurationCmd

TTL 返回 key 的剩余生存时间。key 自动拼接前缀。

key 不存在返回 -2,key 无过期时间返回 -1。

func (*Proxy) TryLock

func (p *Proxy) TryLock(ctx context.Context, key string, ttl time.Duration) (*Lock, error)

TryLock 尝试获取分布式锁,非阻塞。key 自动拼接前缀。

成功返回 *Lock;锁被他人持有返回 ErrLockNotObtained(可 errors.Is 判断); ttl 必须至少为 1ms(Redis TTL 精度为毫秒;无 TTL 的锁等于死锁隐患)。 需要阻塞等待的场景由调用方按业务节奏循环 TryLock。

获取经单段 Lua 完成并对底层自动重试幂等:若 SET NX 因响应丢失被重发, 重发时锁值仍等于本次 token 亦视为获得,不会把自己已持有的锁误报为被他人持有。

func (*Proxy) TxPipeline

func (p *Proxy) TxPipeline() redis.Pipeliner

TxPipeline 返回 go-redis 事务 Pipeline (MULTI/EXEC)。

Pipeline 内的命令需要手动拼接前缀,可通过 Proxy.Key 获取带前缀的 key。

func (*Proxy) Type

func (p *Proxy) Type(ctx context.Context, key string) *redis.StatusCmd

Type 返回 key 存储的值的类型。key 自动拼接前缀。

func (*Proxy) WithLock

func (p *Proxy) WithLock(ctx context.Context, key string, ttl time.Duration, fn func(ctx context.Context) error) error

WithLock 获取锁后执行 fn,并保证释放(含 fn panic 路径),消除手写 TryLock + defer Release 样板与漏释放风险。key 自动拼接前缀。

锁被他人持有时返回 ErrLockNotObtained(errors.Is 判断),fn 不执行; fn 的错误原样传出;释放阶段发现锁已失去(fn 执行超过 ttl,互斥可能已被 破坏)时 ErrLockLosterrors.Join 并入返回错误——调用方应检查 errors.Is(err, ErrLockLost) 并触发业务侧补偿。fn panic 时锁仍被释放, panic 继续向上传播。

不做自动续期:fn 预计耗时必须显著小于 ttl,长任务请自行分段或在 fn 内 通过 Proxy.TryLock 返回的 Lock 显式 Refresh。

func (*Proxy) WithPrefix

func (p *Proxy) WithPrefix(prefix string) (*Proxy, error)

WithPrefix 返回复用同一底层 Redis client、但使用指定 key 前缀的新 Proxy。

派生 Proxy 只替换 key 前缀,保留 key 分隔符、channel 前缀和 channel 分隔符。

Example
package main

import (
	"fmt"
	"log"

	"github.com/gtkit/redisx"

	goredis "github.com/redis/go-redis/v9"
)

func main() {
	rdb := goredis.NewClient(&goredis.Options{Addr: "127.0.0.1:1"})
	defer rdb.Close()

	base, err := redisx.WrapClient(rdb, redisx.WithProxyPrefix("base"))
	if err != nil {
		log.Fatal(err)
	}
	cache, err := base.WithPrefix("cache")
	if err != nil {
		log.Fatal(err)
	}

	fmt.Println(base.Key("k"))
	fmt.Println(cache.Key("k"))
}
Output:
base:k
cache:k

func (*Proxy) XAck

func (p *Proxy) XAck(ctx context.Context, stream, group string, ids ...string) *redis.IntCmd

XAck 确认消费组内一条或多条消息。stream 自动拼接前缀。

func (*Proxy) XAdd

func (p *Proxy) XAdd(ctx context.Context, args *redis.XAddArgs) *redis.StringCmd

XAdd 追加消息到流。args.Stream 自动拼接前缀(不修改传入的 args)。

func (*Proxy) XDel

func (p *Proxy) XDel(ctx context.Context, stream string, ids ...string) *redis.IntCmd

XDel 从流中删除一条或多条消息。stream 自动拼接前缀。

func (*Proxy) XGroupCreateMkStream

func (p *Proxy) XGroupCreateMkStream(ctx context.Context, stream, group, start string) *redis.StatusCmd

XGroupCreateMkStream 创建消费组,流不存在时一并创建(MKSTREAM)。 stream 自动拼接前缀。start 为起始 ID("0" 从头,"$" 只消费新消息)。

Proxy.ConsumeStream 会自动建组,本方法用于需要提前建组或自定义 起始位置的场景。组已存在时 Redis 返回 BUSYGROUP 错误。

func (*Proxy) XGroupDelConsumer

func (p *Proxy) XGroupDelConsumer(ctx context.Context, stream, group, consumer string) *redis.IntCmd

XGroupDelConsumer 从消费组中删除指定消费者,返回其被丢弃的 pending 数。 stream 自动拼接前缀。

以实例 ID 做 Consumer 名时,滚动发布会在组内累积死亡消费者条目, 请在实例下线或定期任务中调用本方法清理;其未确认的消息会被丢弃, 清理前请确认 pending 已被 AutoClaim 接管或确认完毕(见 Proxy.XPendingExt)。

func (*Proxy) XGroupDestroy

func (p *Proxy) XGroupDestroy(ctx context.Context, stream, group string) *redis.IntCmd

XGroupDestroy 销毁消费组(含全部 pending 状态)。stream 自动拼接前缀。

func (*Proxy) XInfoConsumers

func (p *Proxy) XInfoConsumers(ctx context.Context, stream, group string) *redis.XInfoConsumersCmd

XInfoConsumers 返回消费组内全部消费者的状态(pending 数、闲置时长等), 用于发现待清理的死亡消费者。stream 自动拼接前缀。

func (*Proxy) XInfoGroups

func (p *Proxy) XInfoGroups(ctx context.Context, stream string) *redis.XInfoGroupsCmd

XInfoGroups 返回流上全部消费组的状态(consumer 数、pending 数、last-delivered-id 等)。 stream 自动拼接前缀。

func (*Proxy) XLen

func (p *Proxy) XLen(ctx context.Context, stream string) *redis.IntCmd

XLen 返回流的消息数量。stream 自动拼接前缀。

func (*Proxy) XPending

func (p *Proxy) XPending(ctx context.Context, stream, group string) *redis.XPendingCmd

XPending 返回消费组的 pending 摘要(总数、最小/最大 ID、各消费者数量)。 stream 自动拼接前缀。用于监控 handler 持续失败导致的消息堆积。

func (*Proxy) XPendingExt

func (p *Proxy) XPendingExt(ctx context.Context, args *redis.XPendingExtArgs) *redis.XPendingExtCmd

XPendingExt 按条件返回消费组 pending 消息明细。 args.Stream 自动拼接前缀(不修改传入的 args)。

func (*Proxy) XRange

func (p *Proxy) XRange(ctx context.Context, stream, start, stop string) *redis.XMessageSliceCmd

XRange 返回流中 ID 区间内的消息("-" / "+" 表示最小 / 最大 ID)。 stream 自动拼接前缀。

func (*Proxy) XTrimMaxLen

func (p *Proxy) XTrimMaxLen(ctx context.Context, stream string, maxLen int64) *redis.IntCmd

XTrimMaxLen 将流裁剪到最多保留 maxLen 条消息。stream 自动拼接前缀。

func (*Proxy) ZAdd

func (p *Proxy) ZAdd(ctx context.Context, key string, members ...redis.Z) *redis.IntCmd

ZAdd 向有序集合添加一个或多个成员。key 自动拼接前缀。

func (*Proxy) ZCard

func (p *Proxy) ZCard(ctx context.Context, key string) *redis.IntCmd

ZCard 返回有序集合的成员数量。key 自动拼接前缀。

func (*Proxy) ZCount

func (p *Proxy) ZCount(ctx context.Context, key, minScore, maxScore string) *redis.IntCmd

ZCount 返回分值在 minScore 和 maxScore 之间的成员数量。key 自动拼接前缀。

func (*Proxy) ZIncrBy

func (p *Proxy) ZIncrBy(ctx context.Context, key string, increment float64, member string) *redis.FloatCmd

ZIncrBy 为有序集合中指定成员的分值增加增量。key 自动拼接前缀。

func (*Proxy) ZRange

func (p *Proxy) ZRange(ctx context.Context, key string, start, stop int64) *redis.StringSliceCmd

ZRange 按排名范围返回有序集合成员。key 自动拼接前缀。

func (*Proxy) ZRangeByScore

func (p *Proxy) ZRangeByScore(ctx context.Context, key string, opt *redis.ZRangeBy) *redis.StringSliceCmd

ZRangeByScore 按分值范围返回有序集合成员。key 自动拼接前缀。

内部使用 ZRANGE ... BYSCORE 实现(ZRANGEBYSCORE 的现代等价形式,需要 Redis >= 6.2)。

func (*Proxy) ZRank

func (p *Proxy) ZRank(ctx context.Context, key, member string) *redis.IntCmd

ZRank 返回有序集合中指定成员的排名(升序,从 0 开始)。key 自动拼接前缀。

func (*Proxy) ZRem

func (p *Proxy) ZRem(ctx context.Context, key string, members ...any) *redis.IntCmd

ZRem 移除有序集合中一个或多个成员。key 自动拼接前缀。

func (*Proxy) ZRemRangeByScore

func (p *Proxy) ZRemRangeByScore(ctx context.Context, key, minScore, maxScore string) *redis.IntCmd

ZRemRangeByScore 按分值范围移除有序集合成员。key 自动拼接前缀。

func (*Proxy) ZRevRange

func (p *Proxy) ZRevRange(ctx context.Context, key string, start, stop int64) *redis.StringSliceCmd

ZRevRange 按排名范围逆序返回有序集合成员。key 自动拼接前缀。

内部使用 ZRANGE ... REV 实现(ZREVRANGE 的现代等价形式,需要 Redis >= 6.2)。

func (*Proxy) ZScore

func (p *Proxy) ZScore(ctx context.Context, key, member string) *redis.FloatCmd

ZScore 返回有序集合中指定成员的分值。key 自动拼接前缀。

type ProxyOption

type ProxyOption func(*proxyConfig)

ProxyOption 配置由 WrapClient 构造的 Proxy

func WithProxyChannelPrefix

func WithProxyChannelPrefix(prefix string) ProxyOption

WithProxyChannelPrefix 设置 WrapClient 返回 Proxy 的 Pub/Sub channel 前缀。

func WithProxyChannelPrefixSeparator

func WithProxyChannelPrefixSeparator(separator string) ProxyOption

WithProxyChannelPrefixSeparator 设置 WrapClient 返回 Proxy 的 Pub/Sub channel 前缀连接符。

func WithProxyPrefix

func WithProxyPrefix(prefix string) ProxyOption

WithProxyPrefix 设置 WrapClient 返回 Proxy 的 key 前缀。

func WithProxyPrefixSeparator

func WithProxyPrefixSeparator(separator string) ProxyOption

WithProxyPrefixSeparator 设置 WrapClient 返回 Proxy 的 key 前缀连接符。

type StreamConfig

type StreamConfig struct {
	// Stream 是流名,自动拼接前缀。必填。
	Stream string

	// Group 是消费组名。必填。不存在时自动创建(MKSTREAM)。
	Group string

	// Consumer 是消费者名,组内应唯一(如实例 ID)。必填。
	Consumer string

	// BatchSize 是单次 XREADGROUP 最多读取的消息数。默认 16。
	BatchSize int64

	// Block 是无新消息时 XREADGROUP 的阻塞等待时长(占用一个连接)。
	// 0 或负值取默认 5s,无法表达 Redis 的"无限阻塞"语义。
	Block time.Duration

	// AutoClaimMinIdle 大于 0 时,周期性以 XAUTOCLAIM 接管组内其他消费者
	// 闲置超过该时长的 pending 消息(需要 Redis >= 6.2)。0 表示关闭。
	//
	// 注意:无新消息时消费循环先阻塞 Block 时长才检查接管,
	// 实际接管周期下限受 Block 钳制,设置小于 Block 的值没有意义。
	AutoClaimMinIdle time.Duration

	// MaxDeliver 大于 0 时启用死信策略:投递次数(Redis pending entry 的
	// delivery counter,每次投递/接管自增)超过该值的消息不再交给 handler,
	// 原子转入 DeadLetterStream 并 ACK。0 表示关闭(失败消息无限重投)。
	//
	// 语义即"handler 最多尝试 MaxDeliver 次"。仅 pending 续传与 AutoClaim
	// 路径检查(新消息首次投递必然未超限),新消息热路径零额外开销。
	MaxDeliver int64

	// DeadLetterStream 是死信流名,自动拼接前缀。MaxDeliver > 0 时必填,
	// 且不得与 Stream 同名。死信消息保留原消息全部字段,并附加
	// _redisx_origin_stream / _redisx_origin_id / _redisx_deliveries /
	// _redisx_dead_at 元数据(同名业务字段会被元数据遮蔽)。
	//
	// 死信转移语义为 at-least-once:客户端重试可能重复写入(_redisx_origin_id
	// 相同),下游消费者应按 _redisx_origin_id 幂等去重。死信流的裁剪、监控与
	// 重放由业务负责,库内不做 MAXLEN 限制——无人消费时请业务侧自行 XTrimMaxLen。
	DeadLetterStream string

	// OnError 可选:handler 返回业务 error 时被同步调用(此时消息不 ACK,
	// 留在 pending 等待重投)。库内不产生日志,这是业务感知"哪条消息为何
	// 失败"的口子;nil 表示静默(仅靠 XPENDING 监控)。
	//
	// 消息被死信化时同样经本回调通知,err 满足
	// errors.Is(err, ErrMessageDeadLettered)。
	//
	// 回调在消费 goroutine 内执行,应保持轻量(记日志/打点);
	// 回调内 panic 会终止消费并以 error 返回。
	OnError func(msg redis.XMessage, err error)
}

StreamConfig 定义 Proxy.ConsumeStream 受管消费组的配置。

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL