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")
}
Output:
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"
}
Output:
Index ¶
- Constants
- Variables
- func GetBytes(ctx context.Context, p *Proxy, key string) ([]byte, error)
- func GetCodec[T any](ctx context.Context, p *Proxy, key string, codec Codec) (T, error)
- func GetJSON[T any](ctx context.Context, p *Proxy, key string) (T, error)
- func SetBytes(ctx context.Context, p *Proxy, key string, data []byte, ...) error
- func SetCodec(ctx context.Context, p *Proxy, key string, codec Codec, value any, ...) error
- func SetJSON(ctx context.Context, p *Proxy, key string, value any, expiration time.Duration) error
- type Client
- func (c *Client) AddHook(hook redis.Hook)
- func (c *Client) Close() error
- func (c *Client) DefaultClient() *redis.Client
- func (c *Client) GetClient(db int) (*redis.Client, bool)
- func (c *Client) HealthCheck(ctx context.Context) error
- func (c *Client) MustSelectDB(db int) *Proxy
- func (c *Client) PoolStats() map[int]*redis.PoolStats
- func (c *Client) Prefix() string
- func (c *Client) SelectDB(db int) (*Proxy, error)
- type Codec
- type Config
- type DBConfig
- type FencedLock
- type InitError
- type Lock
- type Option
- func WithAddr(addr string) Option
- func WithAllowPartialInit() Option
- func WithChannelPrefix(prefix string) Option
- func WithChannelPrefixSeparator(separator string) Option
- func WithDB(db int) Option
- func WithDBConfig(db int, prefix string) Optiondeprecated
- func WithDialTimeout(d time.Duration) Option
- func WithIdleTimeout(d time.Duration) Option
- func WithInitDBPrefix(db int, prefix string) Option
- func WithInitDBs(dbs ...int) Option
- func WithKeyPrefix(prefix string) Option
- func WithKeyPrefixSeparator(separator string) Option
- func WithMaxRetries(n int) Option
- func WithMinIdleConns(n int) Option
- func WithPassword(password string) Option
- func WithPoolSize(size int) Option
- func WithReadTimeout(d time.Duration) Option
- func WithTLSConfig(tlsConfig *tls.Config) Option
- func WithUsername(username string) Option
- func WithWriteTimeout(d time.Duration) Option
- type Proxy
- func (p *Proxy) Consume(ctx context.Context, handler func(*redis.Message), channels ...string) error
- func (p *Proxy) ConsumePattern(ctx context.Context, handler func(*redis.Message), patterns ...string) error
- func (p *Proxy) ConsumeStream(ctx context.Context, cfg StreamConfig, handler func(redis.XMessage) error) error
- func (p *Proxy) Decr(ctx context.Context, key string) *redis.IntCmd
- func (p *Proxy) DecrBy(ctx context.Context, key string, value int64) *redis.IntCmd
- func (p *Proxy) Del(ctx context.Context, keys ...string) *redis.IntCmd
- func (p *Proxy) DelByPattern(ctx context.Context, pattern string) (int64, error)
- func (p *Proxy) Eval(ctx context.Context, script string, keys []string, args ...any) *redis.Cmd
- func (p *Proxy) EvalScript(ctx context.Context, script *redis.Script, keys []string, args ...any) *redis.Cmd
- func (p *Proxy) EvalSha(ctx context.Context, sha1 string, keys []string, args ...any) *redis.Cmd
- func (p *Proxy) Exists(ctx context.Context, keys ...string) *redis.IntCmd
- func (p *Proxy) Expire(ctx context.Context, key string, expiration time.Duration) *redis.BoolCmd
- func (p *Proxy) ExpireAt(ctx context.Context, key string, tm time.Time) *redis.BoolCmd
- func (p *Proxy) FencedLock(ctx context.Context, key string, ttl time.Duration) (*FencedLock, error)
- func (p *Proxy) Get(ctx context.Context, key string) *redis.StringCmd
- func (p *Proxy) GetDel(ctx context.Context, key string) *redis.StringCmd
- func (p *Proxy) GetSet(ctx context.Context, key string, value any) *redis.StringCmd
- func (p *Proxy) HDel(ctx context.Context, key string, fields ...string) *redis.IntCmd
- func (p *Proxy) HExists(ctx context.Context, key, field string) *redis.BoolCmd
- func (p *Proxy) HGet(ctx context.Context, key, field string) *redis.StringCmd
- func (p *Proxy) HGetAll(ctx context.Context, key string) *redis.MapStringStringCmd
- func (p *Proxy) HIncrBy(ctx context.Context, key, field string, incr int64) *redis.IntCmd
- func (p *Proxy) HIncrByFloat(ctx context.Context, key, field string, incr float64) *redis.FloatCmd
- func (p *Proxy) HLen(ctx context.Context, key string) *redis.IntCmd
- func (p *Proxy) HMGet(ctx context.Context, key string, fields ...string) *redis.SliceCmd
- func (p *Proxy) HMSet(ctx context.Context, key string, values ...any) *redis.BoolCmddeprecated
- func (p *Proxy) HSet(ctx context.Context, key string, values ...any) *redis.IntCmd
- func (p *Proxy) Incr(ctx context.Context, key string) *redis.IntCmd
- func (p *Proxy) IncrBy(ctx context.Context, key string, value int64) *redis.IntCmd
- func (p *Proxy) IncrByFloat(ctx context.Context, key string, value float64) *redis.FloatCmd
- func (p *Proxy) Key(k string) string
- func (p *Proxy) Keys(ks ...string) []string
- func (p *Proxy) LIndex(ctx context.Context, key string, index int64) *redis.StringCmd
- func (p *Proxy) LLen(ctx context.Context, key string) *redis.IntCmd
- func (p *Proxy) LPop(ctx context.Context, key string) *redis.StringCmd
- func (p *Proxy) LPush(ctx context.Context, key string, values ...any) *redis.IntCmd
- func (p *Proxy) LRange(ctx context.Context, key string, start, stop int64) *redis.StringSliceCmd
- func (p *Proxy) LRem(ctx context.Context, key string, count int64, value any) *redis.IntCmd
- func (p *Proxy) LTrim(ctx context.Context, key string, start, stop int64) *redis.StatusCmd
- func (p *Proxy) MGet(ctx context.Context, keys ...string) *redis.SliceCmd
- func (p *Proxy) MSet(ctx context.Context, values ...any) *redis.StatusCmd
- func (p *Proxy) PSubscribe(ctx context.Context, patterns ...string) *redis.PubSub
- func (p *Proxy) PTTL(ctx context.Context, key string) *redis.DurationCmd
- func (p *Proxy) Persist(ctx context.Context, key string) *redis.BoolCmd
- func (p *Proxy) Pipeline() redis.Pipeliner
- func (p *Proxy) Publish(ctx context.Context, channel string, message any) *redis.IntCmd
- func (p *Proxy) RPop(ctx context.Context, key string) *redis.StringCmd
- func (p *Proxy) RPush(ctx context.Context, key string, values ...any) *redis.IntCmd
- func (p *Proxy) RawClient() *redis.Client
- func (p *Proxy) Rename(ctx context.Context, key, newkey string) *redis.StatusCmd
- func (p *Proxy) SAdd(ctx context.Context, key string, members ...any) *redis.IntCmd
- func (p *Proxy) SCard(ctx context.Context, key string) *redis.IntCmd
- func (p *Proxy) SIsMember(ctx context.Context, key string, member any) *redis.BoolCmd
- func (p *Proxy) SMembers(ctx context.Context, key string) *redis.StringSliceCmd
- func (p *Proxy) SPop(ctx context.Context, key string) *redis.StringCmd
- func (p *Proxy) SRandMember(ctx context.Context, key string) *redis.StringCmd
- func (p *Proxy) SRem(ctx context.Context, key string, members ...any) *redis.IntCmd
- func (p *Proxy) Scan(ctx context.Context, cursor uint64, match string, count int64) *redis.ScanCmd
- func (p *Proxy) ScanKeys(ctx context.Context, match string) iter.Seq2[string, error]
- func (p *Proxy) Set(ctx context.Context, key string, value any, expiration time.Duration) *redis.StatusCmd
- func (p *Proxy) SetEX(ctx context.Context, key string, value any, expiration time.Duration) *redis.StatusCmd
- func (p *Proxy) SetNX(ctx context.Context, key string, value any, expiration time.Duration) *redis.BoolCmd
- func (p *Proxy) Subscribe(ctx context.Context, channels ...string) *redis.PubSub
- func (p *Proxy) TTL(ctx context.Context, key string) *redis.DurationCmd
- func (p *Proxy) TryLock(ctx context.Context, key string, ttl time.Duration) (*Lock, error)
- func (p *Proxy) TxPipeline() redis.Pipeliner
- func (p *Proxy) Type(ctx context.Context, key string) *redis.StatusCmd
- func (p *Proxy) WithLock(ctx context.Context, key string, ttl time.Duration, ...) error
- func (p *Proxy) WithPrefix(prefix string) (*Proxy, error)
- func (p *Proxy) XAck(ctx context.Context, stream, group string, ids ...string) *redis.IntCmd
- func (p *Proxy) XAdd(ctx context.Context, args *redis.XAddArgs) *redis.StringCmd
- func (p *Proxy) XDel(ctx context.Context, stream string, ids ...string) *redis.IntCmd
- func (p *Proxy) XGroupCreateMkStream(ctx context.Context, stream, group, start string) *redis.StatusCmd
- func (p *Proxy) XGroupDelConsumer(ctx context.Context, stream, group, consumer string) *redis.IntCmd
- func (p *Proxy) XGroupDestroy(ctx context.Context, stream, group string) *redis.IntCmd
- func (p *Proxy) XInfoConsumers(ctx context.Context, stream, group string) *redis.XInfoConsumersCmd
- func (p *Proxy) XInfoGroups(ctx context.Context, stream string) *redis.XInfoGroupsCmd
- func (p *Proxy) XLen(ctx context.Context, stream string) *redis.IntCmd
- func (p *Proxy) XPending(ctx context.Context, stream, group string) *redis.XPendingCmd
- func (p *Proxy) XPendingExt(ctx context.Context, args *redis.XPendingExtArgs) *redis.XPendingExtCmd
- func (p *Proxy) XRange(ctx context.Context, stream, start, stop string) *redis.XMessageSliceCmd
- func (p *Proxy) XTrimMaxLen(ctx context.Context, stream string, maxLen int64) *redis.IntCmd
- func (p *Proxy) ZAdd(ctx context.Context, key string, members ...redis.Z) *redis.IntCmd
- func (p *Proxy) ZCard(ctx context.Context, key string) *redis.IntCmd
- func (p *Proxy) ZCount(ctx context.Context, key, minScore, maxScore string) *redis.IntCmd
- func (p *Proxy) ZIncrBy(ctx context.Context, key string, increment float64, member string) *redis.FloatCmd
- func (p *Proxy) ZRange(ctx context.Context, key string, start, stop int64) *redis.StringSliceCmd
- func (p *Proxy) ZRangeByScore(ctx context.Context, key string, opt *redis.ZRangeBy) *redis.StringSliceCmd
- func (p *Proxy) ZRank(ctx context.Context, key, member string) *redis.IntCmd
- func (p *Proxy) ZRem(ctx context.Context, key string, members ...any) *redis.IntCmd
- func (p *Proxy) ZRemRangeByScore(ctx context.Context, key, minScore, maxScore string) *redis.IntCmd
- func (p *Proxy) ZRevRange(ctx context.Context, key string, start, stop int64) *redis.StringSliceCmd
- func (p *Proxy) ZScore(ctx context.Context, key, member string) *redis.FloatCmd
- type ProxyOption
- type StreamConfig
Examples ¶
Constants ¶
const Version = "v1.3.0"
Version 是当前库版本号。
Variables ¶
var ErrLockLost = errors.New("redisx: lock lost")
ErrLockLost 表示锁已不再被当前持有者持有(已过期或被他人重新获取)。
var ErrLockNoExpiry = errors.New("redisx: lock has no expiry")
ErrLockNoExpiry 表示锁存在但没有设置过期时间(契约外的永久锁)。
本库自身的获取/续期总会带 TTL,出现此错误通常意味着该 key 被外部 PERSIST 或以无过期方式写入,属死锁隐患,调用方应显式处理而非当作正常。
var ErrLockNotObtained = errors.New("redisx: lock not obtained")
ErrLockNotObtained 表示锁正被他人持有,本次未获得。
var ErrMessageDeadLettered = errors.New("redisx: message dead-lettered")
ErrMessageDeadLettered 表示消息投递次数超过 MaxDeliver,已被转入死信流。 经 StreamConfig.OnError 通知业务时可 errors.Is 判别。
Functions ¶
func GetBytes ¶
GetBytes reads raw bytes from key with automatic key prefixing.
Missing keys preserve redis.Nil through error wrapping.
func GetCodec ¶
GetCodec reads key and unmarshals the stored bytes into T with codec.
Missing keys preserve redis.Nil through error wrapping.
func GetJSON ¶
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)
}
}
Output:
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.
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 ¶
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 ¶
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 ¶
AddHook 将 hook 安装到所有已初始化 DB 的底层 *redis.Client 上。
用于一次性挂接熔断、限流、metrics、tracing 等中间件,库本身不内置任何策略。 hook 为 nil 时静默忽略。部分初始化(WithAllowPartialInit)模式下, 初始化失败而缺席的 DB 自然跳过。
调用时机:与 go-redis 原生 AddHook 的约束一致,必须在 NewClient 之后、 开始并发执行命令之前完成安装,运行中途添加不保证并发安全。 如需只对个别 DB 安装,请改用 Client.GetClient 自行处理。
func (*Client) DefaultClient ¶
DefaultClient 返回默认 DB 的底层 *redis.Client。
func (*Client) GetClient ¶
GetClient 安全获取指定 DB 的底层 *redis.Client。
返回 false 表示该 DB 未初始化。用于需要直接操作 go-redis 原生 API 的场景。
func (*Client) HealthCheck ¶
HealthCheck 对所有已初始化的 DB 并发执行 PING 健康检查。
不会短路:即使某个 DB 失败也会检查其余 DB,最终以 errors.Join 返回所有失败的 聚合错误(按 DB 编号确定有序,不受并发完成顺序影响)。返回 nil 表示全部健康。
func (*Client) MustSelectDB ¶
MustSelectDB 返回指定 DB 编号上的命令代理 Proxy。
如果 DB 未初始化,直接 panic。仅用于程序启动阶段或确定 DB 存在的场景。
func (*Client) PoolStats ¶
PoolStats 返回每个已初始化 DB 的连接池统计,key 为 DB 编号。
透传 go-redis 的连接池统计(命中/未命中/超时/连接数等),供业务接入 监控拉取;库本身不做任何聚合、阈值或告警判断。
func (*Client) SelectDB ¶
SelectDB 返回指定 DB 编号上的命令代理 Proxy,支持链式调用。
如果指定的 DB 未在 WithInitDBs / WithInitDBPrefix 中初始化,返回错误。
type Codec ¶
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 ¶
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 ¶
InitError 聚合降级初始化(WithAllowPartialInit)下各 DB 的失败详情。
通过 errors.As 提取后可按 DB 编程决策:
var ie *redisx.InitError
if errors.As(err, &ie) {
if _, bad := ie.Failed[2]; bad {
// DB2 不可用,按业务决定降级或拒绝启动
}
}
type Lock ¶
type Lock struct {
// contains filtered or unexported fields
}
Lock 表示一把已持有的单实例分布式锁,由 Proxy.TryLock 获取。
注意语义边界:这是单 Redis 实例锁(非 RedLock),主从异步复制下 故障切换瞬间存在双持有的理论窗口,关键互斥请在业务层做幂等兜底。
func (*Lock) Refresh ¶
Refresh 将锁的 TTL 重设为 ttl(校验 token 后原子执行)。
长任务在持有期间显式调用本方法续期;ttl 必须至少为 1ms;锁已失去返回 ErrLockLost。
func (*Lock) Release ¶
Release 释放锁。仅当锁仍被当前持有者持有(token 匹配)时删除; 锁已过期或被他人重新获取时返回 ErrLockLost,且不会影响他人的锁。
func (*Lock) TTL ¶
TTL 返回锁的剩余存活时间(校验 token 后原子读取)。
返回 nil error 即表示锁仍被当前持有者持有;锁已过期或被他人重新获取 返回 ErrLockLost;token 匹配但锁无过期时间(契约外的永久锁)返回 ErrLockNoExpiry。仅供观测与业务自检:查询与后续操作之间锁仍可能 过期,互斥正确性依赖 Release/Refresh 自身的 token 校验,而非本方法。
type Option ¶
type Option func(*Config)
Option 是 Functional Options 模式的配置函数。
func WithAllowPartialInit ¶
func WithAllowPartialInit() Option
WithAllowPartialInit 允许部分 DB 初始化失败(降级模式)。
启用后,初始化失败的 DB 不进入集合,所有失败聚合后与可用的 Client 一同返回(两者可同时非 nil),对缺席 DB 调用 Client.SelectDB 返回错误; DefaultDB 承载 Client 级快捷方法,仍必须初始化成功,否则整体失败。 默认关闭,即任一 DB 失败则整体失败(全有或全无)。
func WithChannelPrefix ¶
WithChannelPrefix 设置 Pub/Sub channel 前缀。
Redis Pub/Sub channel 是实例级全局命名空间,不随 DB 切换;默认不复用 key 前缀。 设置后 Pub/Sub 方法会将 channel 拼接为 "{prefix}{separator}{channel}"。
func WithChannelPrefixSeparator ¶
WithChannelPrefixSeparator 设置 Pub/Sub channel 前缀与原始 channel 之间的连接符。
默认 ":"。仅当 channel 前缀非空时参与拼接;传入空字符串表示直接连接前缀和 channel。
func WithDB ¶
WithDB 设置默认使用的数据库编号。
编号必须 >= 0;上限取决于服务端 databases 配置(默认 16 个库即 0~15), 越界的编号在 NewClient 初始化 Ping 时由 Redis 报错。
func WithDBConfig
deprecated
WithDBConfig 添加一个带独立前缀的 DB 配置。
Deprecated: use WithInitDBPrefix.
func WithDialTimeout ¶
WithDialTimeout 设置建立 TCP 连接的超时时间。
func WithIdleTimeout ¶
WithIdleTimeout 设置空闲连接被回收前的最大存活时间。
func WithInitDBPrefix ¶
WithInitDBPrefix 初始化指定 DB,并为该 DB 设置独立 key 前缀。
当 prefix 非空时,该 DB 使用独立前缀替代全局 KeyPrefix。 兼容 v1 的 WithDB(db, "prefix") 语义。
示例: WithInitDBPrefix(2, "session") 使 DB2 的 key 前缀为 "session:" 而非全局前缀。 前缀连接符可通过 WithKeyPrefixSeparator 全局配置。
func WithInitDBs ¶
WithInitDBs 设置需要初始化的多个 DB 编号。
DefaultDB 会自动包含在列表中,无需重复添加。 所有 DB 共享全局 KeyPrefix。如需 per-DB 前缀,请使用 WithInitDBPrefix。
示例: WithInitDBs(0, 1, 2) 将同时初始化 DB0、DB1、DB2。
func WithKeyPrefix ¶
WithKeyPrefix 设置全局 key 前缀。
设置后所有带 key 参数的命令会自动拼接为 "{prefix}{separator}{key}", 对业务层完全透明。可被 per-DB 前缀覆盖(见 WithInitDBPrefix)。
func WithKeyPrefixSeparator ¶
WithKeyPrefixSeparator 设置 key 前缀与原始 key 之间的连接符。
默认 ":"。仅当前缀非空时参与拼接;传入空字符串表示直接连接前缀和 key。
func WithMaxRetries ¶
WithMaxRetries 设置命令失败后最大重试次数。
n 为 0 表示关闭自动重试(库内会映射为 go-redis 的 -1 哨兵值), 负值在 NewClient 返回错误。
注意:对非幂等命令(如 INCR、LPUSH),读超时后的自动重试可能导致 命令被重复执行。对此类命令敏感的场景请设置为 0 关闭重试, 或在业务层使用 Lua 脚本保证幂等。
func WithReadTimeout ¶
WithReadTimeout 设置 socket 读操作超时时间。
func WithTLSConfig ¶
WithTLSConfig 设置连接 Redis 的 TLS 配置。
传入非 nil 配置后,所有 DB 的连接均通过 TLS 建立, 适用于云厂商强制加密的 Redis 实例。nil 表示不启用 TLS。
func WithWriteTimeout ¶
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) DelByPattern ¶
DelByPattern 使用 SCAN 安全批量删除匹配 pattern 的 key。
pattern 自动拼接前缀:如 pattern="user:*",实际匹配 "{prefix}:user:*"。 内部使用 SCAN + UNLINK 批量删除:UNLINK 在后台异步释放内存, 大 value 场景不会阻塞 Redis 主线程(需要 Redis >= 4.0)。
建议调用方传入带超时的 ctx 控制执行时间。返回成功删除的 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) FencedLock ¶
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
}
Output:
func (*Proxy) GetSet ¶
GetSet 设置新值并返回旧值。key 自动拼接前缀。
内部使用 SET ... GET 实现(GETSET 的现代等价形式,需要 Redis >= 6.2); 旧值不存在时返回 redis.Nil,新值仍会写入。
func (*Proxy) HIncrByFloat ¶
HIncrByFloat 将 hash 中指定字段的浮点数值增加增量。key 自动拼接前缀。
func (*Proxy) HSet ¶
HSet 设置 hash 中一个或多个字段。key 自动拼接前缀。
values 接受 field-value 对: HSet(ctx, "user:1", "name", "alice", "age", 18).
func (*Proxy) IncrByFloat ¶
IncrByFloat 将 key 的浮点数值增加指定增量。key 自动拼接前缀。
func (*Proxy) Keys ¶
Keys 返回逐个拼接了前缀的完整 key 列表。
用于在 Pipeline 等需要手动拼接前缀的场景批量处理多个 key, 例如 pipe.Del(ctx, proxy.Keys("k1", "k2")...)。
func (*Proxy) MSet ¶
MSet 批量设置 key-value 对。
values 为交替的 key-value 序列: MSet(ctx, "k1", "v1", "k2", "v2")。 偶数位(0, 2, 4...)必须为 string,作为 key 自动拼接前缀。
values 长度必须为偶数且 key 位必须为 string,否则返回错误。
func (*Proxy) PSubscribe ¶
PSubscribe 按模式订阅一个或多个 channel 模式(如 "events:*")。 所有模式使用独立 channel 前缀。
返回 *redis.PubSub,调用方负责关闭。
func (*Proxy) Pipeline ¶
Pipeline 返回 go-redis 原生 Pipeline。
Pipeline 内的命令需要手动拼接前缀,可通过 Proxy.Key 获取带前缀的 key。
func (*Proxy) SRandMember ¶
SRandMember 随机返回集合中一个成员。key 自动拼接前缀。
func (*Proxy) Scan ¶
Scan 包装 SCAN 命令,match pattern 自动拼接前缀。
注意:返回的 key 是已含前缀的完整 key,直接回传给本库其他带前缀方法 (如 Proxy.Del)会造成二次拼接。如需对结果继续操作,请改用 Proxy.RawClient 执行,或自行剥离前缀;批量删除场景请直接使用 Proxy.DelByPattern。
func (*Proxy) ScanKeys ¶
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 ¶
Subscribe 订阅一个或多个 channel。所有 channel 使用独立 channel 前缀。
返回 *redis.PubSub,调用方负责关闭;如需自动管理订阅生命周期, 请使用 Proxy.Consume。
func (*Proxy) TryLock ¶
TryLock 尝试获取分布式锁,非阻塞。key 自动拼接前缀。
成功返回 *Lock;锁被他人持有返回 ErrLockNotObtained(可 errors.Is 判断); ttl 必须至少为 1ms(Redis TTL 精度为毫秒;无 TTL 的锁等于死锁隐患)。 需要阻塞等待的场景由调用方按业务节奏循环 TryLock。
获取经单段 Lua 完成并对底层自动重试幂等:若 SET NX 因响应丢失被重发, 重发时锁值仍等于本次 token 亦视为获得,不会把自己已持有的锁误报为被他人持有。
func (*Proxy) TxPipeline ¶
TxPipeline 返回 go-redis 事务 Pipeline (MULTI/EXEC)。
Pipeline 内的命令需要手动拼接前缀,可通过 Proxy.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,互斥可能已被 破坏)时 ErrLockLost 经 errors.Join 并入返回错误——调用方应检查 errors.Is(err, ErrLockLost) 并触发业务侧补偿。fn panic 时锁仍被释放, panic 继续向上传播。
不做自动续期:fn 预计耗时必须显著小于 ttl,长任务请自行分段或在 fn 内 通过 Proxy.TryLock 返回的 Lock 显式 Refresh。
func (*Proxy) WithPrefix ¶
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) 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 ¶
XGroupDestroy 销毁消费组(含全部 pending 状态)。stream 自动拼接前缀。
func (*Proxy) XInfoConsumers ¶
XInfoConsumers 返回消费组内全部消费者的状态(pending 数、闲置时长等), 用于发现待清理的死亡消费者。stream 自动拼接前缀。
func (*Proxy) XInfoGroups ¶
XInfoGroups 返回流上全部消费组的状态(consumer 数、pending 数、last-delivered-id 等)。 stream 自动拼接前缀。
func (*Proxy) XPending ¶
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) XTrimMaxLen ¶
XTrimMaxLen 将流裁剪到最多保留 maxLen 条消息。stream 自动拼接前缀。
func (*Proxy) ZIncrBy ¶
func (p *Proxy) ZIncrBy(ctx context.Context, key string, increment float64, member string) *redis.FloatCmd
ZIncrBy 为有序集合中指定成员的分值增加增量。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) ZRemRangeByScore ¶
ZRemRangeByScore 按分值范围移除有序集合成员。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 受管消费组的配置。