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 |
核心数 × 10 或 QPS/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 篇):
- 启动时执行
CLUSTER SLOTS获取完整槽映射并缓存 → 本地计算 CRC16 直接连正确节点,实现"一跳直达"; - 收到
MOVED时刷新映射表并重试; - 收到
ASK时不更新映射,向目标节点先发ASKING再发命令; - 自动把跨槽的
MGET/MSET拆分成按节点分组的多个并行请求(所以业务代码里MGET跨槽是能用的!); - 定期刷新拓扑,感知扩缩容和故障转移。
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 篇):
- 会丢消息(订阅者不在线就永久丢失),只能用于"允许丢失"的场景(配置更新、缓存失效通知);
- 必须配合 TTL 兜底(如果用于缓存失效);
- 占用独立连接(RESP2 下订阅连接不能执行普通命令);
defer pubsub.Close()否则连接泄漏;- 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 = 500ms 而 BRPop(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和真正的错误 - 阻塞命令用独立的 client 且
ReadTimeout: 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 巡检 + 死信处理 - 接入了链路追踪和 metrics(
redisotel+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:经验值 核心数 × 10 或 QPS / 1000。
- 太小:
PoolTimeout错误、延迟升高(请求排队等连接); - 太大:撑爆 Redis 的
maxclients(rejected_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 // 等待可用连接
MinIdleConns 设 PoolSize/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 = 500ms 而 BRPop(ctx, 0, ...) 要阻塞很久,会触发读超时错误。go-redis 对带 Block 参数的命令会自动用"Block 时间 + 余量"作为读超时,但 Block: 0(永久阻塞)时必须显式把 ReadTimeout 设为 0。
Q4:go-redis 的自动重试可能带来什么问题?
go-redis 默认 MaxRetries: 3,会在网络错误时自动重试。
问题:如果失败的原因是"命令已经发送到服务端并执行了,但响应在网络上丢了",重试会导致命令执行两次。
对非幂等命令是危险的:
INCR→ 计数多加了;LPUSH→ 消息重复入队;ZINCRBY→ 分数多加了;SPOP→ 多弹出一个元素。
三种应对:
- 对关键的非幂等操作禁用重试:
MaxRetries: -1; - 让操作本身幂等:用
SET代替INCR、用唯一 ID 做业务去重; - 用 Lua 脚本 + 业务侧幂等键(脚本内先检查幂等键是否存在)。
这个问题不是 go-redis 独有的——任何带自动重试的客户端都有,是分布式系统里"至少一次 vs 恰好一次"的经典问题。
Q5:集群模式下 MGET 能用吗?go-redis 是怎么处理的?
能用。虽然 Redis 服务端会对跨槽的 MGET 返回 CROSSSLOT 错误,但 go-redis 的 ClusterClient 会自动把 key 按槽分组,向各个节点并行发送多个请求,再合并结果。
同样被自动处理的:MSET、普通 Pipeline(按节点分组并行执行)。
但这些仍然必须同槽(客户端无法拆分):
TxPipeline(MULTI/EXEC);- Lua 脚本(所有
KEYS必须同槽); WATCH;RENAME、SMOVE、BITOP、SINTERSTORE等需要服务端跨 key 操作的命令。
解决办法是 hash tag:{user:1001}:name、{user:1001}:age。注意 tag 粒度必须是高基数标识(用户 ID/订单 ID),用 {user} 这种粗粒度会导致所有数据挤在一个槽,而一个槽不能拆分,无法通过扩容解决(第 12 篇)。
Q6:go-redis 的集群客户端是怎么工作的?
它是"智能客户端"(smart client),实现"一跳直达"(第 12 篇):
- 启动时执行
CLUSTER SLOTS(或CLUSTER SHARDS)获取完整的槽→节点映射并缓存在本地; - 每次请求本地计算
CRC16(key) & 16383得到槽号,查缓存直接连到正确的节点——稳定状态下零重定向开销,这是它优于 Proxy 方案(多一跳)的核心; - 收到
MOVED→ 刷新槽映射表并重试(MaxRedirects控制次数); - 收到
ASK→ 不更新映射表,向目标节点先发ASKING再发命令; - 自动拆分跨槽的
MGET/MSET/Pipeline; - 定期刷新拓扑,感知扩缩容和故障转移。
另外它提供 ForEachMaster/ForEachSlave/ForEachShard 来遍历节点(做 SCAN、FLUSHDB、收集 INFO 时需要)。
注意 PoolSize 在集群模式下是"每个节点"的,总连接数 = PoolSize × 节点数。
Q7:哨兵模式下客户端怎么感知主节点切换?
redis.NewFailoverClient 内部自动完成(第 11 篇):
- 连接任意一个哨兵(配置里给一组地址,逐个尝试);
SENTINEL GET-MASTER-ADDR-BY-NAME查询当前主节点地址;SUBSCRIBE +switch-master订阅切换事件;- 收到切换通知时重新查询地址并重建整个连接池(旧连接全部作废);
- 收到
READONLY You can't write against a read only replica错误时重新查询地址(说明连到了被降级的旧主)。
客户端侧要注意:
- 必须配多个哨兵地址(单个哨兵挂了就拿不到主节点地址);
- 故障转移期间有 10~35 秒的写不可用窗口,要有带退避的重试和业务降级;
- 不要自己缓存主节点地址。
如果需要读写分离,用 NewFailoverClusterClient(把只读命令发到从节点),但要清楚主从延迟的风险(第 10 篇)。
Q8:redis.NewScript 相比直接 Eval 有什么好处?
redis.NewScript 自动处理了 EVALSHA → NOSCRIPT → EVAL 的回退逻辑(第 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 commandstats 的 usec_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)判断。- 必须配置的四个超时:
DialTimeout、ReadTimeout(正常命令 < 1ms,设 300ms~1s)、WriteTimeout、PoolTimeout。不配的话网络故障时 goroutine 无限堆积。 PoolSize按"核数 × 10"或"QPS/1000",太大会撑爆服务端maxclients和内存。用PoolStats().Timeouts调优。设MinIdleConns预热。ClientName务必设置——服务端的CLIENT LIST和SLOWLOG靠它定位是哪个服务(第 15 篇的事故里就因为没设而多花了排查时间)。- 阻塞命令(
BRPop/XReadGroup Block)必须用独立的 client,否则会耗尽连接池;Block: 0时要把ReadTimeout设为 0。 - Pipeline 要遍历检查每条命令的
cmd.Err(),整体 err 为 nil 不代表全部成功;每批 100~1000 条,不要一次几万条(回复会堆爆缓冲区)。 - Lua 用
redis.NewScript(自动处理EVALSHA→NOSCRIPT→EVAL回退,因为脚本缓存不持久化);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、巡检XAutoClaim(MinIdle> 最大处理时间 × 2~3)+ 死信处理。消费者名要稳定唯一(用 pod 名)。生产者必须设MaxLen+Approx: true。 - 生产封装应该内建三大防护:空值缓存(防穿透)+
singleflight+ Double Check(防击穿)+ 随机 TTL 抖动(防雪崩)。 - 可观测性:
redisotel(链路 + metrics)+redisprometheus(连接池)+ 自定义 Hook(慢命令)。把客户端 P99 和服务端commandstats对比,差值大说明是网络/客户端问题。 - 用
redis.UniversalClient接口编写业务代码,让单机/哨兵/集群的切换只是配置变更。
xingliuhua