目录

Redis-16 Go 客户端 go-redis 实战

1. 客户端选型

Go 生态里的 Redis 客户端:

说明
redis/go-redis 事实标准。功能最全(哨兵、集群、Stream、RESP3、OpenTelemetry),社区活跃,文档好。原路径 github.com/go-redis/redis,v9 起迁到 github.com/redis/go-redis/v9(官方接管)
gomodule/redigo 老牌库,API 更"裸"(手动管理连接和命令参数),灵活但样板代码多。老项目常见
rueian/rueidis 新兴库,主打性能(原生 RESP3、自动 pipeline、内置客户端缓存),benchmark 显著优于 go-redis。适合极致性能场景,但生态和成熟度不如 go-redis

结论:新项目直接用 github.com/redis/go-redis/v9。本文以它为主,性能敏感场景可以关注 rueidis。

go get github.com/redis/go-redis/v9

2. 连接与配置

2.1 单机连接

package main

import (
    "context"
    "time"

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

func NewClient() *redis.Client {
    return redis.NewClient(&redis.Options{
        Addr:     "127.0.0.1:6379",
        Password: "mypassword",
        DB:       0,

        // 【连接名】务必设置,让服务端 CLIENT LIST / SLOWLOG 能定位到是哪个服务
        ClientName: "order-service-" + hostname(),

        // ---------- 连接池 ----------
        PoolSize:     50,                // 最大连接数
        MinIdleConns: 10,                // 最小空闲连接(预热,避免突发流量时建连延迟)
        MaxIdleConns: 20,                // 最大空闲连接
        PoolTimeout:  1 * time.Second,   // 等待可用连接的超时

        // ---------- 超时(必须设置!)----------
        DialTimeout:  1 * time.Second,
        ReadTimeout:  500 * time.Millisecond,
        WriteTimeout: 500 * time.Millisecond,

        // ---------- 连接生命周期 ----------
        ConnMaxIdleTime: 5 * time.Minute,   // 空闲超过这么久就关闭
        ConnMaxLifetime: 0,                 // 0 = 不主动关闭(有 LB/代理时建议设,如 30 分钟)

        // ---------- 重试 ----------
        MaxRetries:      2,                          // -1 表示不重试
        MinRetryBackoff: 8 * time.Millisecond,
        MaxRetryBackoff: 512 * time.Millisecond,

        // ---------- 其他 ----------
        Protocol: 3,                       // 使用 RESP3(默认 3,设 2 强制 RESP2)
        // OnConnect: 每个新连接建立后的钩子
        OnConnect: func(ctx context.Context, cn *redis.Conn) error {
            return nil
        },
    })
}

// 用 URL 方式(更适合从配置/环境变量读取)
func NewClientFromURL() (*redis.Client, error) {
    opt, err := redis.ParseURL("redis://user:password@127.0.0.1:6379/0?protocol=3")
    if err != nil {
        return nil, err
    }
    opt.PoolSize = 50
    return redis.NewClient(opt), nil
}

2.2 启动时验证连接

func main() {
    rdb := NewClient()
    defer rdb.Close()

    ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
    defer cancel()

    // 【必须做】启动时验证连通性,让配置错误尽早暴露
    if err := rdb.Ping(ctx).Err(); err != nil {
        log.Fatalf("Redis 连接失败: %v", err)
    }
    log.Println("Redis 连接成功")
}

2.3 关键配置的选择依据

配置 建议值 理由
PoolSize 核心数 × 10QPS/1000 太小会 PoolTimeout太大会撑爆 Redis 的 maxclients 并增加服务端内存(每连接有输入输出缓冲区)
MinIdleConns PoolSize / 5 预热连接,避免突发流量时的建连延迟(建连要 TCP 握手 + AUTH,可能几毫秒)
ReadTimeout 300ms ~ 1s 正常命令应该 < 1ms。设太大会让故障时的请求堆积;设太小会误杀慢命令
DialTimeout 500ms ~ 2s 建连要握手,比读写略长
PoolTimeout ReadTimeout + 1s 等连接的时间不该超过命令本身的超时太多
MaxRetries 2 ~ 3 go-redis 默认 3。注意:不是所有命令都该重试(见第 8 节的坑)
ConnMaxLifetime 0(无 LB)/ 30min(有 LB) 云负载均衡/代理可能静默断开长连接,主动轮换更安全

2.4 监控连接池

func MonitorPool(rdb *redis.Client) {
    ticker := time.NewTicker(30 * time.Second)
    defer ticker.Stop()
    for range ticker.C {
        s := rdb.PoolStats()
        // Hits: 从池里拿到空闲连接的次数
        // Misses: 池里没有空闲连接、需要新建的次数
        // Timeouts: 【等待连接超时的次数,> 0 说明池不够用】
        // TotalConns / IdleConns / StaleConns
        metrics.Gauge("redis.pool.total", float64(s.TotalConns))
        metrics.Gauge("redis.pool.idle", float64(s.IdleConns))
        metrics.Counter("redis.pool.timeouts", float64(s.Timeouts))

        if s.Timeouts > 0 {
            log.Warnf("连接池等待超时 %d 次,考虑调大 PoolSize", s.Timeouts)
        }
    }
}

3. 基本操作

3.1 错误处理的核心:redis.Nil

// 【最重要的一点】key 不存在时返回 redis.Nil,这【不是错误】
val, err := rdb.Get(ctx, "nonexistent").Result()
switch {
case err == redis.Nil:
    // key 不存在,正常业务分支
    return nil, ErrNotFound
case err != nil:
    // 真正的错误(网络、超时、Redis 报错)
    return nil, fmt.Errorf("redis error: %w", err)
default:
    return parse(val), nil
}

// 用 errors.Is 更规范(能穿透 wrap)
if errors.Is(err, redis.Nil) { ... }

这是新手最常见的 bug:把 redis.Nil 当成错误往上抛,导致"缓存未命中"被当成"系统故障",触发不必要的告警或降级。

哪些命令会返回 redis.Nil

  • GET(key 不存在);
  • LPOP/RPOP/SPOP(集合为空);
  • BLPOP/BRPOP(阻塞超时);
  • HGET(field 不存在);
  • ZSCORE(member 不存在);
  • SET ... NX(条件不满足,注意这个容易漏)。
// SET NX 的正确判断(分布式锁的基础)
ok, err := rdb.SetNX(ctx, "lock", "v", 30*time.Second).Result()
if err != nil {
    return err          // 真正的错误
}
if !ok {
    return ErrLockFailed  // 锁已被持有(不是错误)
}

// 如果用 Set 带 NX 参数,则要判断 redis.Nil
err := rdb.Set(ctx, "k", "v", 0).Err()      // 普通 SET,不会返回 Nil

3.2 String

// 基本读写
rdb.Set(ctx, "name", "tom", time.Hour)
rdb.Set(ctx, "name", "tom", 0)                    // 0 = 不过期
rdb.Set(ctx, "name", "tom", redis.KeepTTL)        // 【保留原有 TTL】(6.0+ 的 KEEPTTL)
val, err := rdb.Get(ctx, "name").Result()

// 各种类型转换
i, err := rdb.Get(ctx, "count").Int64()
f, err := rdb.Get(ctx, "price").Float64()
b, err := rdb.Get(ctx, "flag").Bool()
bs, err := rdb.Get(ctx, "data").Bytes()

// SetNX / SetXX / GetSet / GetDel / GetEx
ok, _ := rdb.SetNX(ctx, "lock", "v", 30*time.Second).Result()
ok, _ = rdb.SetXX(ctx, "k", "v", time.Hour).Result()
old, _ := rdb.GetSet(ctx, "k", "newval").Result()      // 注意:会清除 TTL
v, _ := rdb.GetDel(ctx, "k").Result()                  // 6.2+ 取出并删除
v, _ = rdb.GetEx(ctx, "k", time.Hour).Result()         // 6.2+ 取出并设过期

// 批量
rdb.MSet(ctx, "k1", "v1", "k2", "v2")
rdb.MSet(ctx, map[string]interface{}{"k1": "v1", "k2": "v2"})
vals, _ := rdb.MGet(ctx, "k1", "k2", "k3").Result()
// 【注意】MGet 返回 []interface{},不存在的 key 对应 nil
for i, v := range vals {
    if v == nil {
        log.Printf("key %d 不存在", i)
        continue
    }
    s := v.(string)
    _ = s
}

// 计数
n, _ := rdb.Incr(ctx, "counter").Result()
rdb.IncrBy(ctx, "counter", 100)
rdb.IncrByFloat(ctx, "price", 1.5)
rdb.Decr(ctx, "counter")

3.3 Hash

rdb.HSet(ctx, "user:1001", "name", "tom", "age", 20)
rdb.HSet(ctx, "user:1001", map[string]interface{}{"name": "tom", "age": 20})
// 也支持结构体(用 redis tag)
type User struct {
    Name string `redis:"name"`
    Age  int    `redis:"age"`
}
rdb.HSet(ctx, "user:1001", User{Name: "tom", Age: 20})

name, _ := rdb.HGet(ctx, "user:1001", "name").Result()
vals, _ := rdb.HMGet(ctx, "user:1001", "name", "age").Result()

// 【HGetAll 在大 hash 上危险】(第 2 篇),字段多时用 HScan
all, _ := rdb.HGetAll(ctx, "user:1001").Result()   // map[string]string

// 【很好用】直接扫描到结构体
var u User
err := rdb.HGetAll(ctx, "user:1001").Scan(&u)

rdb.HIncrBy(ctx, "user:1001", "age", 1)
rdb.HDel(ctx, "user:1001", "age")
rdb.HExists(ctx, "user:1001", "name")
rdb.HLen(ctx, "user:1001")

3.4 List / Set / ZSet

// ---------- List ----------
rdb.LPush(ctx, "queue", "task1", "task2")
rdb.RPush(ctx, "queue", "task3")
v, _ := rdb.LPop(ctx, "queue").Result()
vs, _ := rdb.LPopCount(ctx, "queue", 10).Result()      // 6.2+ 弹多个
items, _ := rdb.LRange(ctx, "queue", 0, 99).Result()   // 【分页,不要用 0 -1】
rdb.LTrim(ctx, "queue", 0, 19)                          // 保留最新 20 条
rdb.LLen(ctx, "queue")
rdb.LRem(ctx, "queue", 1, "task1")

// 阻塞弹出(【必须用独立的 client,见第 8 节】)
res, err := blockingRdb.BRPop(ctx, 5*time.Second, "queue").Result()
if err == redis.Nil {
    // 超时,没有消息(正常情况)
} else if err == nil {
    // res[0] = key 名,res[1] = 值
    log.Printf("从 %s 取到 %s", res[0], res[1])
}
rdb.LMove(ctx, "src", "dst", "RIGHT", "LEFT")           // 6.2+
rdb.BLMove(ctx, "src", "dst", "RIGHT", "LEFT", 0)

// ---------- Set ----------
rdb.SAdd(ctx, "tags", "redis", "golang")
rdb.SRem(ctx, "tags", "redis")
ok, _ := rdb.SIsMember(ctx, "tags", "redis").Result()
oks, _ := rdb.SMIsMember(ctx, "tags", "redis", "mysql").Result()   // 6.2+
n, _ := rdb.SCard(ctx, "tags").Result()
ms, _ := rdb.SMembers(ctx, "tags").Result()             // 【大 set 上危险】
rdb.SInter(ctx, "s1", "s2")
rdb.SInterCard(ctx, 0, "s1", "s2")                      // 7.0+ 只要交集数量
rdb.SRandMemberN(ctx, "lottery", 3)
rdb.SPopN(ctx, "lottery", 3)

// ---------- ZSet ----------
rdb.ZAdd(ctx, "rank", redis.Z{Score: 100, Member: "player1"})
rdb.ZAddArgs(ctx, "rank", redis.ZAddArgs{
    GT:      true,       // 只在新分数更大时更新(6.2+)
    Members: []redis.Z{{Score: 200, Member: "player1"}},
})
rdb.ZIncrBy(ctx, "rank", 10, "player1")

// Top 10(带分数)
zs, _ := rdb.ZRevRangeWithScores(ctx, "rank", 0, 9).Result()
for _, z := range zs {
    log.Printf("%v: %v", z.Member, z.Score)
}

rank, _ := rdb.ZRevRank(ctx, "rank", "player1").Result()
score, _ := rdb.ZScore(ctx, "rank", "player1").Result()

// 按分数范围
res, _ := rdb.ZRangeByScore(ctx, "rank", &redis.ZRangeBy{
    Min:    "60",
    Max:    "(100",       // 开区间
    Offset: 0,
    Count:  10,
}).Result()

// 6.2+ 统一的 ZRangeArgs
rdb.ZRangeArgs(ctx, redis.ZRangeArgs{
    Key:     "rank",
    Start:   "(60",
    Stop:    "+inf",
    ByScore: true,
    Rev:     true,
    Offset:  0,
    Count:   10,
})

rdb.ZRemRangeByScore(ctx, "rank", "0", "60")
rdb.ZPopMax(ctx, "rank", 1)

3.5 SCAN 迭代器

// 【永远用 SCAN,不要用 KEYS】(第 2 篇)
// go-redis 提供了 Iterator 封装,自动处理游标
iter := rdb.Scan(ctx, 0, "user:*", 100).Iterator()
for iter.Next(ctx) {
    key := iter.Val()
    // 【注意】SCAN 可能返回重复的 key,业务要能容忍或自己去重
    process(key)
}
if err := iter.Err(); err != nil {
    return err
}

// 类型内的扫描
iter = rdb.HScan(ctx, "bighash", 0, "", 100).Iterator()
for iter.Next(ctx) {
    field := iter.Val()
    iter.Next(ctx)          // 【注意】HScan 返回的是 field, value 交替
    value := iter.Val()
    _ = field
    _ = value
}

iter = rdb.SScan(ctx, "bigset", 0, "", 100).Iterator()
iter = rdb.ZScan(ctx, "bigzset", 0, "", 100).Iterator()

// 【集群模式】要遍历所有节点
clusterClient.ForEachMaster(ctx, func(ctx context.Context, master *redis.Client) error {
    iter := master.Scan(ctx, 0, "user:*", 100).Iterator()
    for iter.Next(ctx) {
        process(iter.Val())
    }
    return iter.Err()
})

4. Pipeline

4.1 普通 Pipeline

// 方式一:手动创建
pipe := rdb.Pipeline()
incr := pipe.Incr(ctx, "counter")
pipe.Expire(ctx, "counter", time.Hour)
cmds, err := pipe.Exec(ctx)

// 【关键】err != nil 不代表全部失败,必须遍历检查每条命令
if err != nil && err != redis.Nil {
    log.Printf("pipeline 整体错误: %v", err)
}
for _, cmd := range cmds {
    if cmd.Err() != nil {
        log.Printf("命令 %v 失败: %v", cmd.Args(), cmd.Err())
    }
}
// 拿单条命令的结果
log.Println(incr.Val())

// 方式二:Pipelined 回调(更简洁)
cmds, err := rdb.Pipelined(ctx, func(pipe redis.Pipeliner) error {
    for i := 0; i < 100; i++ {
        pipe.Set(ctx, fmt.Sprintf("k%d", i), i, time.Hour)
    }
    return nil
})

4.2 事务 Pipeline(MULTI/EXEC)

// TxPipeline 会自动用 MULTI/EXEC 包裹
cmds, err := rdb.TxPipelined(ctx, func(pipe redis.Pipeliner) error {
    pipe.Incr(ctx, "counter")
    pipe.Expire(ctx, "counter", time.Hour)
    return nil
})
// 注意:Redis 事务不支持回滚(第 8 篇),某条命令的运行时错误
// 不影响其他命令的执行

4.3 批量导入的正确写法

// 处理 100 万条数据的导入
func BulkImport(ctx context.Context, rdb *redis.Client, items []Item) error {
    const batchSize = 500     // 【每批 100~1000,不要一次几万条】

    for i := 0; i < len(items); i += batchSize {
        end := min(i+batchSize, len(items))
        batch := items[i:end]

        _, err := rdb.Pipelined(ctx, func(pipe redis.Pipeliner) error {
            for _, item := range batch {
                data, _ := json.Marshal(item)
                pipe.Set(ctx, item.Key(), data, randomTTL(time.Hour))
            }
            return nil
        })
        if err != nil {
            return fmt.Errorf("批次 %d 失败: %w", i/batchSize, err)
        }
    }
    return nil
}

为什么不能一次塞几万条(第 5、8 篇):所有回复会堆在服务端的输出缓冲区和客户端的接收缓冲区里。1 万条命令、每条回复 1KB = 10MB 内存,可能触发 client-output-buffer-limit 或撑爆客户端内存。

4.4 WATCH 乐观锁

// go-redis 的 Watch 封装了 WATCH + MULTI/EXEC 的模板
func TransferWithWatch(ctx context.Context, rdb *redis.Client, key string, amount int64) error {
    const maxRetries = 5

    txf := func(tx *redis.Tx) error {
        // 【在 WATCH 之后读取当前值】
        balance, err := tx.Get(ctx, key).Int64()
        if err != nil && err != redis.Nil {
            return err
        }
        // 业务判断
        if balance < amount {
            return errors.New("余额不足")
        }
        // 【在 MULTI 里写入】
        _, err = tx.TxPipelined(ctx, func(pipe redis.Pipeliner) error {
            pipe.DecrBy(ctx, key, amount)
            return nil
        })
        return err
    }

    for i := 0; i < maxRetries; i++ {
        err := rdb.Watch(ctx, txf, key)
        if err == nil {
            return nil
        }
        if err == redis.TxFailedErr {
            continue        // 被监视的 key 被修改了,重试
        }
        return err
    }
    return errors.New("达到最大重试次数")
}

但如第 8、13 篇所说:绝大多数需要 WATCH 的场景,用 Lua 更简单高效。


5. Lua 脚本

5.1 用 NewScript(推荐)

// NewScript 内部自动处理 EVALSHA → NOSCRIPT → EVAL 的回退(第 8 篇)
var stockScript = redis.NewScript(`
    local stock = tonumber(redis.call('GET', KEYS[1]) or '0')
    local qty = tonumber(ARGV[1])
    if stock < qty then
        return -1
    end
    redis.call('DECRBY', KEYS[1], qty)
    return stock - qty
`)

func DeductStock(ctx context.Context, rdb *redis.Client, skuID int64, qty int) (int64, error) {
    key := fmt.Sprintf("stock:%d", skuID)
    remain, err := stockScript.Run(ctx, rdb, []string{key}, qty).Int64()
    if err != nil {
        return 0, err
    }
    if remain < 0 {
        return 0, ErrOutOfStock
    }
    return remain, nil
}

// 也可以预加载(可选,Run 会自动处理)
func init() {
    // 在启动时预加载脚本,减少首次调用的开销
    // 但注意:主从切换/重启后缓存会丢,所以 Run 的自动回退仍然是必需的
}

5.2 常用脚本封装

// ---------- 分布式锁的释放(第 13 篇)----------
var unlockScript = redis.NewScript(`
    if redis.call('GET', KEYS[1]) == ARGV[1] then
        return redis.call('DEL', KEYS[1])
    end
    return 0
`)

// ---------- 分布式锁的续期 ----------
var renewScript = redis.NewScript(`
    if redis.call('GET', KEYS[1]) == ARGV[1] then
        return redis.call('PEXPIRE', KEYS[1], ARGV[2])
    end
    return 0
`)

// ---------- 滑动窗口限流(第 2 篇)----------
var slidingWindowScript = redis.NewScript(`
    local key = KEYS[1]
    local now = tonumber(ARGV[1])
    local window = tonumber(ARGV[2])
    local limit = tonumber(ARGV[3])
    redis.call('ZREMRANGEBYSCORE', key, 0, now - window)
    local count = redis.call('ZCARD', key)
    if count >= limit then
        return 0
    end
    redis.call('ZADD', key, now, ARGV[4])
    redis.call('PEXPIRE', key, window)
    return 1
`)

func AllowRequest(ctx context.Context, rdb *redis.Client,
    userID int64, limit int, window time.Duration) (bool, error) {

    key := fmt.Sprintf("rate:%d", userID)
    now := time.Now().UnixMilli()
    member := fmt.Sprintf("%d-%d", now, rand.Int63())   // member 必须唯一

    ok, err := slidingWindowScript.Run(ctx, rdb,
        []string{key}, now, window.Milliseconds(), limit, member).Int()
    if err != nil {
        // 【降级策略】:限流器故障时选择"放行"还是"拒绝"取决于业务
        return true, err     // 这里选择放行(可用性优先)
    }
    return ok == 1, nil
}

// ---------- 原子的"取出并删除到期任务"(第 9 篇延迟队列)----------
var fetchDelayedScript = redis.NewScript(`
    local tasks = redis.call('ZRANGEBYSCORE', KEYS[1], 0, ARGV[1], 'LIMIT', 0, ARGV[2])
    if #tasks > 0 then
        redis.call('ZREM', KEYS[1], unpack(tasks))
    end
    return tasks
`)

5.3 Lua 的注意点(复习第 8 篇)

// ✅ 所有 key 通过 KEYS 传入(Cluster 路由和 ACL 需要)
script.Run(ctx, rdb, []string{"key1", "key2"}, arg1, arg2)

// ❌ 绝不能在脚本里拼 key 名
// return redis.call('GET', 'user:' .. ARGV[1])
//   → 集群下会报错,且会导致脚本缓存无限膨胀

// ⚠️ 集群模式下所有 KEYS 必须在同一个 slot,用 hash tag
keys := []string{
    fmt.Sprintf("seckill:{%d}:stock", skuID),
    fmt.Sprintf("seckill:{%d}:users", skuID),
}

// ⚠️ 返回值类型转换(第 8 篇的坑)
n, _ := script.Run(...).Int64()      // Lua number → 整数(小数会被截断!)
s, _ := script.Run(...).Text()       // Lua string
ss, _ := script.Run(...).StringSlice()  // Lua table
b, _ := script.Run(...).Bool()       // Lua false → nil → 这里会是 error redis.Nil

6. 哨兵与集群客户端

6.1 哨兵模式

rdb := redis.NewFailoverClient(&redis.FailoverOptions{
    MasterName: "mymaster",
    SentinelAddrs: []string{         // 【必须配多个】
        "192.168.1.21:26379",
        "192.168.1.22:26379",
        "192.168.1.23:26379",
    },
    Password:         "redis-password",      // Redis 节点密码
    SentinelPassword: "sentinel-password",   // 哨兵密码(如果哨兵设了 requirepass)
    DB:               0,
    PoolSize:         50,
    ClientName:       "order-service",

    // 【读写分离】谨慎使用(第 10 篇:有延迟、3.2 前会读到过期数据)
    // ReplicaOnly:    true,             // 只读从节点
    // RouteByLatency: true,             // 按延迟路由(主从中选最快的)
    // RouteRandomly:  true,             // 随机路由到从节点
})

// go-redis 内部自动完成(第 11 篇):
// 1. SENTINEL GET-MASTER-ADDR-BY-NAME 查询主节点地址
// 2. SUBSCRIBE +switch-master 订阅切换事件
// 3. 切换时自动【重建整个连接池】
// 4. 收到 READONLY 错误时重新查询地址

如果需要显式读写分离(同时用主和从):

// 这个客户端会把只读命令发到从节点、写命令发到主节点
rdb := redis.NewFailoverClusterClient(&redis.FailoverOptions{
    MasterName:     "mymaster",
    SentinelAddrs:  []string{...},
    RouteByLatency: true,
})

6.2 集群模式

rdb := redis.NewClusterClient(&redis.ClusterOptions{
    Addrs: []string{                  // 给几个种子节点即可,客户端会自动发现全部
        "127.0.0.1:7000",
        "127.0.0.1:7001",
        "127.0.0.1:7002",
    },
    Password:   "mypassword",
    ClientName: "order-service",

    // 【每个节点】的连接池大小(总连接数 = PoolSize × 节点数)
    PoolSize:     20,
    MinIdleConns: 5,

    DialTimeout:  1 * time.Second,
    ReadTimeout:  500 * time.Millisecond,
    WriteTimeout: 500 * time.Millisecond,

    // 重定向重试次数(MOVED/ASK,第 12 篇)
    MaxRedirects: 3,

    // 【读写分离】
    ReadOnly:       false,     // true = 只读命令发到从节点
    RouteByLatency: false,     // 按延迟选节点
    RouteRandomly:  false,     // 随机选节点

    // 自定义节点的 Options(比如给不同节点不同配置)
    // NewClient: func(opt *redis.Options) *redis.Client { ... },
})

// 【集群特有的操作】
// 遍历所有主节点
rdb.ForEachMaster(ctx, func(ctx context.Context, master *redis.Client) error {
    return master.FlushDB(ctx).Err()
})
// 遍历所有从节点
rdb.ForEachSlave(ctx, func(ctx context.Context, slave *redis.Client) error {
    return slave.Ping(ctx).Err()
})
// 遍历所有节点(主+从)
rdb.ForEachShard(ctx, func(ctx context.Context, shard *redis.Client) error {
    return shard.Ping(ctx).Err()
})

// 手动刷新槽映射(通常不需要,客户端会自动处理 MOVED)
rdb.ReloadState(ctx)

go-redis 集群客户端做的事(第 12 篇):

  1. 启动时执行 CLUSTER SLOTS 获取完整槽映射并缓存 → 本地计算 CRC16 直接连正确节点,实现"一跳直达"
  2. 收到 MOVED 时刷新映射表并重试;
  3. 收到 ASK不更新映射,向目标节点先发 ASKING 再发命令;
  4. 自动把跨槽的 MGET/MSET 拆分成按节点分组的多个并行请求(所以业务代码里 MGET 跨槽是能用的!);
  5. 定期刷新拓扑,感知扩缩容和故障转移。

6.3 集群模式的限制处理

// ---------- MGET/MSET:客户端会自动拆分,可以照常用 ----------
vals, err := clusterRdb.MGet(ctx, "k1", "k2", "k3").Result()   // ✅ 能用

// ---------- Pipeline:客户端会按节点分组并行执行 ----------
cmds, err := clusterRdb.Pipelined(ctx, func(pipe redis.Pipeliner) error {
    pipe.Set(ctx, "k1", "v1", 0)
    pipe.Set(ctx, "k2", "v2", 0)
    return nil
})   // ✅ 能用(内部拆成多个请求)

// ---------- TxPipeline / Lua / WATCH:【必须同 slot】 ----------
// ❌ 会报 CROSSSLOT
clusterRdb.TxPipelined(ctx, func(pipe redis.Pipeliner) error {
    pipe.Set(ctx, "k1", "v1", 0)
    pipe.Set(ctx, "k2", "v2", 0)
    return nil
})

// ✅ 用 hash tag 保证同 slot
clusterRdb.TxPipelined(ctx, func(pipe redis.Pipeliner) error {
    pipe.Set(ctx, "{user1001}:name", "tom", 0)
    pipe.Set(ctx, "{user1001}:age", "20", 0)
    return nil
})

// ---------- 查看 key 属于哪个 slot ----------
slot, _ := clusterRdb.ClusterKeySlot(ctx, "mykey").Result()

// ---------- SCAN 要遍历所有主节点 ----------
clusterRdb.ForEachMaster(ctx, func(ctx context.Context, master *redis.Client) error {
    iter := master.Scan(ctx, 0, "user:*", 100).Iterator()
    for iter.Next(ctx) {
        process(iter.Val())
    }
    return iter.Err()
})

hash tag 的粒度提醒(第 12 篇):用高基数标识{user:1001}{order:88888}),绝不要用 {user}{orders} 这种粗粒度 tag,否则所有数据挤在一个槽,无法通过扩容解决。

6.4 通用接口:redis.UniversalClient

// 一套代码同时支持单机/哨兵/集群(很实用)
func NewUniversalClient(cfg *Config) redis.UniversalClient {
    return redis.NewUniversalClient(&redis.UniversalOptions{
        Addrs:      cfg.Addrs,          // 多个地址 → 集群;配了 MasterName → 哨兵;单个 → 单机
        MasterName: cfg.MasterName,     // 非空则用哨兵模式
        Password:   cfg.Password,
        DB:         cfg.DB,             // 注意:集群模式下只能是 0
        PoolSize:   cfg.PoolSize,
        ClientName: cfg.ClientName,
    })
}

// UniversalClient 是接口,包含了绝大多数命令
// 但集群特有的方法(ForEachMaster 等)需要类型断言
if cluster, ok := rdb.(*redis.ClusterClient); ok {
    cluster.ForEachMaster(ctx, ...)
}

这是生产代码的推荐写法——业务代码依赖 redis.UniversalClient 接口,部署时通过配置决定用哪种模式,从单机迁移到集群不用改业务代码。


7. Pub/Sub 与 Stream

7.1 Pub/Sub

// ---------- 订阅(注意:会占用独立连接)----------
func Subscribe(ctx context.Context, rdb *redis.Client) {
    pubsub := rdb.Subscribe(ctx, "news:tech", "news:sport")
    defer pubsub.Close()          // 【必须 Close】否则连接泄漏

    // 【推荐】等待订阅确认,确保订阅真的生效了
    if _, err := pubsub.Receive(ctx); err != nil {
        log.Printf("订阅失败: %v", err)
        return
    }

    // Channel() 返回一个带缓冲的 Go channel,内部有 goroutine 在读
    ch := pubsub.Channel()
    // 也可以自定义缓冲大小和健康检查间隔
    // ch := pubsub.Channel(redis.WithChannelSize(1000),
    //                      redis.WithChannelHealthCheckInterval(3*time.Second))

    for {
        select {
        case msg, ok := <-ch:
            if !ok {
                return          // channel 关闭
            }
            log.Printf("收到 [%s]: %s", msg.Channel, msg.Payload)
        case <-ctx.Done():
            return
        }
    }
}

// 模式订阅
pubsub := rdb.PSubscribe(ctx, "news:*")

// ---------- 发布 ----------
n, err := rdb.Publish(ctx, "news:tech", "Redis 8.0 released").Result()
// n = 收到消息的订阅者数量

// 集群模式的分片 pub/sub(7.0+,第 9 篇)
rdb.SPublish(ctx, "shard-channel", "msg")
pubsub := rdb.SSubscribe(ctx, "shard-channel")

Pub/Sub 的使用要点(第 9 篇):

  1. 会丢消息(订阅者不在线就永久丢失),只能用于"允许丢失"的场景(配置更新、缓存失效通知);
  2. 必须配合 TTL 兜底(如果用于缓存失效);
  3. 占用独立连接(RESP2 下订阅连接不能执行普通命令);
  4. defer pubsub.Close() 否则连接泄漏;
  5. go-redis 的 Channel() 内部会自动重连(网络中断时),但重连期间的消息会丢

7.2 Stream 消费者(完整实现)

type StreamConsumer struct {
    rdb      redis.UniversalClient
    stream   string
    group    string
    consumer string        // 【必须稳定唯一】:用 pod 名而不是 UUID(第 9 篇)
    handler  func(ctx context.Context, msg redis.XMessage) error
}

func (c *StreamConsumer) Start(ctx context.Context) error {
    // ---------- 1. 创建消费组(MKSTREAM 让 stream 不存在时自动创建)----------
    err := c.rdb.XGroupCreateMkStream(ctx, c.stream, c.group, "0").Err()
    if err != nil && !strings.Contains(err.Error(), "BUSYGROUP") {
        return fmt.Errorf("创建消费组失败: %w", err)
    }

    // ---------- 2. 崩溃恢复:先处理自己未 ACK 的消息 ----------
    if err := c.recoverPending(ctx); err != nil {
        log.Printf("恢复 PEL 失败: %v", err)
    }

    // ---------- 3. 启动巡检协程(认领超时消息)----------
    go c.claimLoop(ctx)

    // ---------- 4. 正常消费循环 ----------
    return c.consumeLoop(ctx)
}

// 用 "0" 读自己的 PEL(">" 只给从未投递过的新消息)
func (c *StreamConsumer) recoverPending(ctx context.Context) error {
    start := "0"
    for {
        res, err := c.rdb.XReadGroup(ctx, &redis.XReadGroupArgs{
            Group:    c.group,
            Consumer: c.consumer,
            Streams:  []string{c.stream, start},
            Count:    100,
        }).Result()
        if err == redis.Nil {
            return nil
        }
        if err != nil {
            return err
        }
        if len(res) == 0 || len(res[0].Messages) == 0 {
            return nil        // PEL 已清空
        }
        for _, msg := range res[0].Messages {
            c.process(ctx, msg)
            start = msg.ID
        }
    }
}

func (c *StreamConsumer) consumeLoop(ctx context.Context) error {
    for {
        select {
        case <-ctx.Done():
            return ctx.Err()
        default:
        }

        res, err := c.rdb.XReadGroup(ctx, &redis.XReadGroupArgs{
            Group:    c.group,
            Consumer: c.consumer,
            Streams:  []string{c.stream, ">"},    // ">" = 新消息
            Count:    10,
            Block:    5 * time.Second,            // 阻塞等待,避免空轮询
        }).Result()

        if err == redis.Nil {
            continue          // Block 超时,正常
        }
        if err != nil {
            log.Printf("XReadGroup 错误: %v", err)
            time.Sleep(time.Second)
            continue
        }

        for _, stream := range res {
            for _, msg := range stream.Messages {
                c.process(ctx, msg)
            }
        }
    }
}

func (c *StreamConsumer) process(ctx context.Context, msg redis.XMessage) {
    if err := c.handler(ctx, msg); err != nil {
        log.Printf("处理消息 %s 失败: %v", msg.ID, err)
        return            // 【不 ACK】,留在 PEL 等重试
    }
    if err := c.rdb.XAck(ctx, c.stream, c.group, msg.ID).Err(); err != nil {
        log.Printf("ACK %s 失败: %v", msg.ID, err)
    }
}

// 认领其他消费者卡住的消息 + 死信处理
func (c *StreamConsumer) claimLoop(ctx context.Context) {
    ticker := time.NewTicker(30 * time.Second)
    defer ticker.Stop()

    for {
        select {
        case <-ctx.Done():
            return
        case <-ticker.C:
            start := "0-0"
            for {
                msgs, next, err := c.rdb.XAutoClaim(ctx, &redis.XAutoClaimArgs{
                    Stream:   c.stream,
                    Group:    c.group,
                    Consumer: c.consumer,
                    // 【必须 > 单条消息的最大处理时间 × 2~3】(第 9 篇)
                    MinIdle: 60 * time.Second,
                    Start:   start,
                    Count:   10,
                }).Result()
                if err != nil {
                    log.Printf("XAutoClaim 错误: %v", err)
                    break
                }
                if len(msgs) == 0 {
                    break
                }
                for _, msg := range msgs {
                    c.processWithDeadLetter(ctx, msg)
                }
                if next == "0-0" {
                    break
                }
                start = next
            }
        }
    }
}

// 死信处理:重试次数过多的消息转存并 ACK
func (c *StreamConsumer) processWithDeadLetter(ctx context.Context, msg redis.XMessage) {
    pending, err := c.rdb.XPendingExt(ctx, &redis.XPendingExtArgs{
        Stream: c.stream,
        Group:  c.group,
        Start:  msg.ID,
        End:    msg.ID,
        Count:  1,
    }).Result()

    if err == nil && len(pending) > 0 && pending[0].RetryCount > 5 {
        log.Printf("消息 %s 重试超限,转入死信", msg.ID)
        c.rdb.XAdd(ctx, &redis.XAddArgs{
            Stream: c.stream + ":dead",
            Values: msg.Values,
        })
        c.rdb.XAck(ctx, c.stream, c.group, msg.ID)
        return
    }
    c.process(ctx, msg)
}

// ---------- 生产者 ----------
func Produce(ctx context.Context, rdb redis.UniversalClient, values map[string]interface{}) error {
    return rdb.XAdd(ctx, &redis.XAddArgs{
        Stream: "stream:orders",
        MaxLen: 100000,      // 【必须设裁剪】否则内存无限增长
        Approx: true,        // 用 ~ 近似裁剪,性能好得多
        ID:     "*",
        Values: values,
    }).Err()
}

8. 常见坑与最佳实践

8.1 十个常见坑

坑一:把 redis.Nil 当错误处理

// ❌ 缓存未命中被当成系统故障,触发不必要的告警/降级
val, err := rdb.Get(ctx, key).Result()
if err != nil {
    return nil, err       // redis.Nil 也被抛出去了
}

// ✅
val, err := rdb.Get(ctx, key).Result()
if errors.Is(err, redis.Nil) {
    return nil, ErrNotFound    // 正常业务分支
}
if err != nil {
    return nil, err            // 真正的错误
}

坑二:忘了设置超时

// ❌ 网络故障时请求永久挂起,goroutine 无限堆积 → OOM
redis.NewClient(&redis.Options{Addr: addr})

// ✅
redis.NewClient(&redis.Options{
    Addr:         addr,
    DialTimeout:  1 * time.Second,
    ReadTimeout:  500 * time.Millisecond,
    WriteTimeout: 500 * time.Millisecond,
    PoolTimeout:  1 * time.Second,
})

坑三:阻塞命令用了共享的 client(连接池被耗尽)

// ❌ BRPop 会长期占用连接,10 个消费者就占了 10 个连接不放
// 如果 PoolSize 是 10,业务请求就完全拿不到连接了
rdb.BRPop(ctx, 0, "queue")

// ✅ 为阻塞消费创建【独立的 client】
blockingRdb := redis.NewClient(&redis.Options{
    Addr:        addr,
    PoolSize:    consumerCount + 2,
    ReadTimeout: 0,        // 【阻塞命令要把 ReadTimeout 设为 0 或大于 Block 时间】
})
blockingRdb.BRPop(ctx, 0, "queue")

注意 ReadTimeout 与阻塞命令的关系:如果 ReadTimeout = 500msBRPop(ctx, 0, ...) 要阻塞很久,会触发读超时错误。go-redis 对阻塞命令做了特殊处理(会用 Block 时间 + 余量作为读超时),但 Block: 0(永久阻塞)时必须把 ReadTimeout 设为 0

坑四:SET 清除了 TTL(导致内存泄漏)

// ❌ 更新缓存时忘了带 TTL,key 变成永久(第 2、15 篇的经典事故)
rdb.Set(ctx, key, newVal, 0)

// ✅ 三种正确做法
rdb.Set(ctx, key, newVal, randomTTL(time.Hour))     // 带新 TTL
rdb.Set(ctx, key, newVal, redis.KeepTTL)            // 保留原 TTL(6.0+)
rdb.Del(ctx, key)                                    // 干脆删除,让下次读重建(Cache Aside)

坑五:Pipeline 只检查整体错误

// ❌ 某条命令失败被忽略
cmds, err := rdb.Pipelined(ctx, fn)
if err != nil {
    return err
}
// 直接用结果,但某条命令可能失败了

// ✅ 遍历检查每条
cmds, err := rdb.Pipelined(ctx, fn)
if err != nil && !errors.Is(err, redis.Nil) {
    return err
}
for _, cmd := range cmds {
    if cmd.Err() != nil && !errors.Is(cmd.Err(), redis.Nil) {
        return fmt.Errorf("命令 %v 失败: %w", cmd.Args(), cmd.Err())
    }
}

坑六:集群模式下用了跨槽的事务/Lua

// ❌ CROSSSLOT 错误
clusterRdb.TxPipelined(ctx, func(pipe redis.Pipeliner) error {
    pipe.Set(ctx, "user:1001:name", "tom", 0)
    pipe.Set(ctx, "user:1001:age", "20", 0)
    return nil
})

// ✅ hash tag(且粒度要细)
pipe.Set(ctx, "{user:1001}:name", "tom", 0)
pipe.Set(ctx, "{user:1001}:age", "20", 0)

坑七:忘了 defer pubsub.Close()

// ❌ 连接泄漏
pubsub := rdb.Subscribe(ctx, "channel")
for msg := range pubsub.Channel() { ... }

// ✅
pubsub := rdb.Subscribe(ctx, "channel")
defer pubsub.Close()

坑八:写命令被自动重试导致重复执行

// go-redis 默认 MaxRetries=3,会在网络错误时重试
// 【问题】:如果是"命令已发送但响应丢失",重试会导致命令执行两次!
// 对 INCR、LPUSH 这类【非幂等】命令是危险的

// ✅ 方案一:对关键的非幂等操作禁用重试
criticalRdb := redis.NewClient(&redis.Options{
    Addr:       addr,
    MaxRetries: -1,       // -1 表示不重试
})

// ✅ 方案二:让操作本身幂等(用 SET 代替 INCR、用唯一 ID 去重)

// ✅ 方案三:用 Lua 脚本 + 业务侧的幂等键

坑九:HGetAll/SMembers/LRange 0 -1 用在大 key 上

// ❌ 大 hash 上会阻塞主线程 + 打满网络(第 14、15 篇)
all, _ := rdb.HGetAll(ctx, "bighash").Result()

// ✅ 只取需要的字段
vals, _ := rdb.HMGet(ctx, "bighash", "f1", "f2").Result()

// ✅ 或用 HScan 分批
iter := rdb.HScan(ctx, "bighash", 0, "", 100).Iterator()

坑十:连接池大小设置不当

// ❌ PoolSize 设成 1000(每个连接在 Redis 端都有缓冲区,会撑爆 maxclients 和内存)
// ❌ PoolSize 设成 5 但并发有 100(大量 PoolTimeout)

// ✅ 按"核心数 × 10"或"QPS/1000"估算,然后【看 PoolStats().Timeouts 调整】

8.2 生产级封装示例

package cache

import (
    "context"
    "encoding/json"
    "errors"
    "fmt"
    "math/rand"
    "time"

    "github.com/redis/go-redis/v9"
    "golang.org/x/sync/singleflight"
)

var ErrNotFound = errors.New("not found")

type Cache struct {
    rdb redis.UniversalClient
    sf  singleflight.Group      // 进程内合并(防击穿,第 14 篇)
}

func New(rdb redis.UniversalClient) *Cache {
    return &Cache{rdb: rdb}
}

// GetOrLoad 是缓存读取的标准封装,包含了:
// 空值缓存(防穿透)+ singleflight(防击穿)+ 随机 TTL(防雪崩)
func (c *Cache) GetOrLoad(ctx context.Context, key string, ttl time.Duration,
    dest interface{}, loader func(ctx context.Context) (interface{}, error)) error {

    // 1. 读缓存
    val, err := c.rdb.Get(ctx, key).Result()
    if err == nil {
        if val == nullPlaceholder {
            return ErrNotFound        // 【命中空值缓存】
        }
        return json.Unmarshal([]byte(val), dest)
    }
    if !errors.Is(err, redis.Nil) {
        // Redis 故障 → 【降级】直接走 loader(但要有限流保护数据源!)
        return c.loadDirect(ctx, loader, dest)
    }

    // 2. 缓存未命中:用 singleflight 合并同一个 key 的并发请求
    v, err, _ := c.sf.Do(key, func() (interface{}, error) {
        // 【Double Check】等 singleflight 期间可能已被填充
        if val, err := c.rdb.Get(ctx, key).Result(); err == nil {
            if val == nullPlaceholder {
                return nil, ErrNotFound
            }
            return []byte(val), nil
        }

        // 3. 调用数据源
        data, err := loader(ctx)
        if errors.Is(err, ErrNotFound) {
            // 【空值缓存,TTL 要短】(防穿透,第 14 篇)
            c.rdb.Set(ctx, key, nullPlaceholder, 60*time.Second)
            return nil, ErrNotFound
        }
        if err != nil {
            return nil, err
        }

        // 4. 回写缓存(【随机 TTL 抖动】防雪崩)
        b, err := json.Marshal(data)
        if err != nil {
            return nil, err
        }
        c.rdb.Set(ctx, key, b, jitter(ttl))
        return b, nil
    })
    if err != nil {
        return err
    }
    return json.Unmarshal(v.([]byte), dest)
}

// Delete 更新数据后调用(Cache Aside:先更新 DB,再删缓存)
func (c *Cache) Delete(ctx context.Context, keys ...string) error {
    if err := c.rdb.Del(ctx, keys...).Err(); err != nil {
        // 【删除失败要重试或告警】(第 14 篇)
        return fmt.Errorf("删除缓存失败: %w", err)
    }
    return nil
}

const nullPlaceholder = "\x00__NULL__"

// jitter 给 TTL 加 ±10% 的随机抖动
func jitter(d time.Duration) time.Duration {
    delta := int64(float64(d) * 0.1)
    if delta <= 0 {
        return d
    }
    return d + time.Duration(rand.Int63n(delta*2)-delta)
}

func (c *Cache) loadDirect(ctx context.Context,
    loader func(ctx context.Context) (interface{}, error), dest interface{}) error {
    data, err := loader(ctx)
    if err != nil {
        return err
    }
    b, _ := json.Marshal(data)
    return json.Unmarshal(b, dest)
}

使用

var user User
err := cache.GetOrLoad(ctx, fmt.Sprintf("user:%d", id), time.Hour, &user,
    func(ctx context.Context) (interface{}, error) {
        u, err := db.GetUser(ctx, id)
        if errors.Is(err, sql.ErrNoRows) {
            return nil, ErrNotFound      // 会被缓存为空值
        }
        return u, err
    })

8.3 可观测性

import (
    "github.com/redis/go-redis/extra/redisotel/v9"
    "github.com/redis/go-redis/extra/redisprometheus/v9"
    "github.com/prometheus/client_golang/prometheus"
)

// ---------- OpenTelemetry 链路追踪 + Metrics ----------
if err := redisotel.InstrumentTracing(rdb); err != nil {
    log.Fatal(err)
}
if err := redisotel.InstrumentMetrics(rdb); err != nil {
    log.Fatal(err)
}

// ---------- Prometheus 连接池指标 ----------
collector := redisprometheus.NewCollector("myapp", "redis", rdb)
prometheus.MustRegister(collector)

// ---------- 自定义 Hook(记录慢命令、埋点)----------
type slowLogHook struct {
    threshold time.Duration
}

func (h slowLogHook) DialHook(next redis.DialHook) redis.DialHook {
    return next
}

func (h slowLogHook) ProcessHook(next redis.ProcessHook) redis.ProcessHook {
    return func(ctx context.Context, cmd redis.Cmder) error {
        start := time.Now()
        err := next(ctx, cmd)
        cost := time.Since(start)

        // 上报每个命令的延迟(便于和服务端 commandstats 对比,第 15 篇)
        metrics.Histogram("redis.cmd.duration",
            cost.Seconds(), "cmd", cmd.Name())

        if cost > h.threshold {
            log.Warnf("慢命令 %v 耗时 %v", cmd.Args(), cost)
        }
        return err
    }
}

func (h slowLogHook) ProcessPipelineHook(next redis.ProcessPipelineHook) redis.ProcessPipelineHook {
    return func(ctx context.Context, cmds []redis.Cmder) error {
        start := time.Now()
        err := next(ctx, cmds)
        metrics.Histogram("redis.pipeline.duration",
            time.Since(start).Seconds(), "size", len(cmds))
        return err
    }
}

rdb.AddHook(slowLogHook{threshold: 100 * time.Millisecond})

8.4 生产检查清单

  • 设置了 ClientName(服务端排查时能定位到是哪个服务)
  • 设置了所有超时Dial/Read/Write/Pool
  • PoolSize 按核数/QPS 估算,并监控 PoolStats().Timeouts
  • 启动时 Ping 验证连通性
  • 正确区分 redis.Nil 和真正的错误
  • 阻塞命令用独立的 clientReadTimeout: 0
  • Pipeline 遍历检查每条命令的错误
  • 更新缓存时不会意外清除 TTL(用 KeepTTL 或删除)
  • 所有缓存 key 都有 TTL 且带随机抖动
  • Lua 脚本用 NewScript(自动处理 NOSCRIPT),key 全部通过 KEYS 传入
  • 集群模式下事务/Lua 的 key 用 hash tag 同槽,且 tag 粒度是高基数标识
  • defer pubsub.Close() / defer rdb.Close()
  • 非幂等的关键写操作考虑禁用自动重试MaxRetries: -1
  • 不用 KEYS,用 Scan 迭代器;不在大 key 上用 HGetAll/SMembers/LRange 0 -1
  • Stream 生产者设了 MaxLen + Approx: true;消费者名稳定唯一、有崩溃恢复 + XAutoClaim 巡检 + 死信处理
  • 接入了链路追踪和 metricsredisotel + redisprometheus
  • redis.UniversalClient 接口,便于在单机/哨兵/集群之间切换

9. 高频面试题

Q1:go-redis 里 redis.Nil 是什么?为什么要特殊处理?

redis.Nil 是 go-redis 定义的一个哨兵错误,表示"Redis 返回了空回复",对应 RESP 协议里的 $-1(Null Bulk String)或 *-1(Null Array)。

它不是错误,而是一个正常的业务状态。会返回它的场景:

  • GET 时 key 不存在;
  • LPOP/RPOP/SPOP 时集合为空;
  • BLPOP/BRPOP/XReadGroup 阻塞超时;
  • HGET 时 field 不存在;
  • ZSCORE 时 member 不存在。

必须特殊处理的原因:如果直接把它当错误往上抛,“缓存未命中"就会被当成"系统故障”,导致不必要的告警、错误率飙升、甚至触发熔断降级。

if errors.Is(err, redis.Nil) {
    return nil, ErrNotFound     // 正常分支
}
if err != nil {
    return nil, err             // 真正的错误
}

Q2:go-redis 的连接池应该怎么配?

PoolSize:经验值 核心数 × 10QPS / 1000

  • 太小PoolTimeout 错误、延迟升高(请求排队等连接);
  • 太大:撑爆 Redis 的 maxclientsrejected_connections 增长),且每个连接在 Redis 端都有输入输出缓冲区(输入上限 1GB、输出 16KB 固定 + 可增长链表),连接多了服务端内存开销可观。

必须配的超时(否则网络故障时 goroutine 无限堆积 → OOM):

DialTimeout:  1 * time.Second      // 建连(含 TCP 握手 + AUTH)
ReadTimeout:  500 * time.Millisecond   // 正常命令应该 < 1ms
WriteTimeout: 500 * time.Millisecond
PoolTimeout:  1 * time.Second      // 等待可用连接

MinIdleConnsPoolSize/5:预热连接,避免突发流量时的建连延迟。

调优依据:看 rdb.PoolStats()——Timeouts > 0 说明池不够用;Misses 高说明连接经常新建(调大 MinIdleConns)。

Q3:为什么阻塞命令(BRPop)要用独立的 client?

因为阻塞命令会长期独占一个连接(第 5 篇:Redis 服务端把客户端挂到阻塞队列上,主线程照常工作,但这个 TCP 连接被占死了)。

如果用共享的 client:10 个消费者 BRPop(ctx, 0, ...) 就占了 10 个连接不放。如果 PoolSize 是 10,业务请求就完全拿不到连接了(全部 PoolTimeout)。

// 为阻塞消费创建独立 client
blockingRdb := redis.NewClient(&redis.Options{
    Addr:        addr,
    PoolSize:    consumerCount + 2,
    ReadTimeout: 0,        // 【关键】永久阻塞时必须设为 0
})

ReadTimeout 的坑:如果 ReadTimeout = 500msBRPop(ctx, 0, ...) 要阻塞很久,会触发读超时错误。go-redis 对带 Block 参数的命令会自动用"Block 时间 + 余量"作为读超时,但 Block: 0(永久阻塞)时必须显式把 ReadTimeout 设为 0

Q4:go-redis 的自动重试可能带来什么问题?

go-redis 默认 MaxRetries: 3,会在网络错误时自动重试。

问题:如果失败的原因是"命令已经发送到服务端并执行了,但响应在网络上丢了",重试会导致命令执行两次

非幂等命令是危险的:

  • INCR → 计数多加了;
  • LPUSH → 消息重复入队;
  • ZINCRBY → 分数多加了;
  • SPOP → 多弹出一个元素。

三种应对

  1. 对关键的非幂等操作禁用重试MaxRetries: -1
  2. 让操作本身幂等:用 SET 代替 INCR、用唯一 ID 做业务去重;
  3. 用 Lua 脚本 + 业务侧幂等键(脚本内先检查幂等键是否存在)。

这个问题不是 go-redis 独有的——任何带自动重试的客户端都有,是分布式系统里"至少一次 vs 恰好一次"的经典问题。

Q5:集群模式下 MGET 能用吗?go-redis 是怎么处理的?

能用。虽然 Redis 服务端会对跨槽的 MGET 返回 CROSSSLOT 错误,但 go-redis 的 ClusterClient 会自动把 key 按槽分组,向各个节点并行发送多个请求,再合并结果

同样被自动处理的MSET、普通 Pipeline(按节点分组并行执行)。

但这些仍然必须同槽(客户端无法拆分):

  • TxPipeline(MULTI/EXEC)
  • Lua 脚本(所有 KEYS 必须同槽);
  • WATCH
  • RENAMESMOVEBITOPSINTERSTORE 等需要服务端跨 key 操作的命令。

解决办法是 hash tag{user:1001}:name{user:1001}:age注意 tag 粒度必须是高基数标识(用户 ID/订单 ID),用 {user} 这种粗粒度会导致所有数据挤在一个槽,而一个槽不能拆分,无法通过扩容解决(第 12 篇)。

Q6:go-redis 的集群客户端是怎么工作的?

它是"智能客户端"(smart client),实现"一跳直达"(第 12 篇):

  1. 启动时执行 CLUSTER SLOTS(或 CLUSTER SHARDS)获取完整的槽→节点映射并缓存在本地
  2. 每次请求本地计算 CRC16(key) & 16383 得到槽号,查缓存直接连到正确的节点——稳定状态下零重定向开销,这是它优于 Proxy 方案(多一跳)的核心;
  3. 收到 MOVED → 刷新槽映射表并重试(MaxRedirects 控制次数);
  4. 收到 ASK不更新映射表,向目标节点先发 ASKING 再发命令
  5. 自动拆分跨槽的 MGET/MSET/Pipeline
  6. 定期刷新拓扑,感知扩缩容和故障转移。

另外它提供 ForEachMaster/ForEachSlave/ForEachShard 来遍历节点(做 SCANFLUSHDB、收集 INFO 时需要)。

注意 PoolSize 在集群模式下是"每个节点"的,总连接数 = PoolSize × 节点数

Q7:哨兵模式下客户端怎么感知主节点切换?

redis.NewFailoverClient 内部自动完成(第 11 篇):

  1. 连接任意一个哨兵(配置里给一组地址,逐个尝试);
  2. SENTINEL GET-MASTER-ADDR-BY-NAME 查询当前主节点地址;
  3. SUBSCRIBE +switch-master 订阅切换事件;
  4. 收到切换通知时重新查询地址并重建整个连接池(旧连接全部作废);
  5. 收到 READONLY You can't write against a read only replica 错误时重新查询地址(说明连到了被降级的旧主)。

客户端侧要注意

  • 必须配多个哨兵地址(单个哨兵挂了就拿不到主节点地址);
  • 故障转移期间有 10~35 秒的写不可用窗口,要有带退避的重试和业务降级;
  • 不要自己缓存主节点地址

如果需要读写分离,用 NewFailoverClusterClient(把只读命令发到从节点),但要清楚主从延迟的风险(第 10 篇)。

Q8:redis.NewScript 相比直接 Eval 有什么好处?

redis.NewScript 自动处理了 EVALSHANOSCRIPTEVAL 的回退逻辑(第 8 篇):

1. 先尝试 EVALSHA <sha1>(只传 40 字符,省网络)
2. 如果返回 NOSCRIPT(脚本不在服务端缓存里)
   → 自动降级为 EVAL <完整脚本>(会顺便把脚本加入缓存)
3. 之后继续用 EVALSHA

为什么必须有这个回退脚本缓存不持久化,以下情况都会丢失:

  • Redis 重启
  • 执行了 SCRIPT FLUSH
  • 主从切换后新主节点没有这个脚本的缓存;
  • 客户端连到了新扩容的节点。

自己写 Eval 就要手动实现这套逻辑,容易漏。

(7.0+ 的 Function 从根本上解决了这个问题——函数库会持久化到 RDB/AOF 并复制到从节点,不需要 NOSCRIPT 回退。)

Q9:怎么在 go-redis 里做可观测性?

三个层次

(1)官方扩展包

import "github.com/redis/go-redis/extra/redisotel/v9"
redisotel.InstrumentTracing(rdb)     // OpenTelemetry 链路追踪
redisotel.InstrumentMetrics(rdb)     // Metrics

import "github.com/redis/go-redis/extra/redisprometheus/v9"
prometheus.MustRegister(redisprometheus.NewCollector("app", "redis", rdb))  // 连接池指标

(2)自定义 Hook(记录慢命令、上报延迟)

实现 redis.Hook 接口的 DialHook/ProcessHook/ProcessPipelineHook,在命令前后埋点。

(3)连接池指标

s := rdb.PoolStats()
// Hits / Misses / Timeouts / TotalConns / IdleConns / StaleConns

关键实践(第 15 篇):把客户端埋点的 P99 延迟和服务端 INFO commandstatsusec_per_call 对比——两者差值大说明问题在网络或客户端(GC、连接池等待、序列化),而不是 Redis 本身。这是区分"Redis 慢"和"客户端慢"最有效的手段。

Q10:redis.UniversalClient 有什么用?

它是一个接口NewUniversalClient 会根据配置自动返回单机/哨兵/集群客户端:

redis.NewUniversalClient(&redis.UniversalOptions{
    Addrs:      addrs,          // 多个地址 → ClusterClient
    MasterName: masterName,     // 非空 → FailoverClient(哨兵)
    // 单个地址且无 MasterName → Client(单机)
})

价值:让业务代码依赖接口而不是具体实现,部署时通过配置决定用哪种模式。这样从单机迁移到哨兵、再到集群,都不用改业务代码——这在项目演进过程中非常有用(开发环境用单机、生产用集群)。

注意集群特有的方法(ForEachMaster 等)不在接口里,需要类型断言。另外集群模式下 DB 只能是 0(第 12 篇)。


小结

  • 新项目用 github.com/redis/go-redis/v9(官方接管的事实标准);极致性能场景可关注 rueidis(原生 RESP3 + 自动 pipeline + 客户端缓存)。
  • redis.Nil 不是错误,是"key 不存在/集合为空/阻塞超时"的正常状态。把它当错误抛出是最常见的 bug(会让缓存未命中变成系统故障告警)。用 errors.Is(err, redis.Nil) 判断。
  • 必须配置的四个超时DialTimeoutReadTimeout(正常命令 < 1ms,设 300ms~1s)、WriteTimeoutPoolTimeout。不配的话网络故障时 goroutine 无限堆积。
  • PoolSize 按"核数 × 10"或"QPS/1000",太大会撑爆服务端 maxclients 和内存。PoolStats().Timeouts 调优。设 MinIdleConns 预热。
  • ClientName 务必设置——服务端的 CLIENT LISTSLOWLOG 靠它定位是哪个服务(第 15 篇的事故里就因为没设而多花了排查时间)。
  • 阻塞命令(BRPop/XReadGroup Block)必须用独立的 client,否则会耗尽连接池;Block: 0 时要把 ReadTimeout 设为 0
  • Pipeline 要遍历检查每条命令的 cmd.Err(),整体 err 为 nil 不代表全部成功;每批 100~1000 条,不要一次几万条(回复会堆爆缓冲区)。
  • Lua 用 redis.NewScript(自动处理 EVALSHANOSCRIPTEVAL 回退,因为脚本缓存不持久化);key 必须全部通过 KEYS 传入
  • 集群客户端是"智能客户端":本地缓存槽映射实现一跳直达,自动处理 MOVED/ASK(ASK 会先发 ASKING)、自动拆分跨槽的 MGET/MSET/Pipeline。但 TxPipeline/Lua/WATCH 仍必须同槽(用 hash tag,且 tag 粒度要是高基数标识)。
  • 哨兵客户端自动完成"查主节点地址 → 订阅 +switch-master → 切换时重建连接池 → 收到 READONLY 时重新查询"。必须配多个哨兵地址;要容忍 10~35 秒的切换窗口。
  • 自动重试对非幂等命令有风险(“命令已执行但响应丢失"时重试会执行两次):关键的 INCR/LPUSH 类操作考虑 MaxRetries: -1 或让操作幂等。
  • SET 会清除 TTL——更新缓存时用 redis.KeepTTL 或改成 Del(Cache Aside),否则会造成内存泄漏(第 15 篇的真实事故)。
  • Stream 消费者三件事:启动时用 STREAMS key 0 恢复自己的 PEL、正常消费用 > + Block巡检 XAutoClaimMinIdle > 最大处理时间 × 2~3)+ 死信处理。消费者名要稳定唯一(用 pod 名)。生产者必须设 MaxLen + Approx: true
  • 生产封装应该内建三大防护空值缓存(防穿透)+ singleflight + Double Check(防击穿)+ 随机 TTL 抖动(防雪崩)。
  • 可观测性redisotel(链路 + metrics)+ redisprometheus(连接池)+ 自定义 Hook(慢命令)。把客户端 P99 和服务端 commandstats 对比,差值大说明是网络/客户端问题。
  • redis.UniversalClient 接口编写业务代码,让单机/哨兵/集群的切换只是配置变更。