Go + Redis 实战:go-redis/v9 连接池、Pipeline 与事务

使用 go-redis/v9 在 Go 中操作 Redis,覆盖连接池、Pipeline 批量操作、事务(WATCH/MULTI/EXEC)、Lua 脚本与完整缓存封装

Go 语言凭借简洁的语法、出色的并发模型与高效的运行时,在服务端开发领域占据重要位置。Redis 作为内存型键值数据库,以亚毫秒级延迟与丰富的数据结构成为缓存与实时数据存储的首选。将两者结合,能够构建出高性能、高可用的后端服务。本文基于 go-redis/v9 这一官方推荐的 Go Redis 客户端,系统讲解连接管理、数据操作、Pipeline 批量处理、乐观锁事务、Lua 脚本、发布订阅,并最终落地为带 TTL 和防穿透的 Cache-Aside 缓存封装,以及计数器/限流器实战项目。

go-redis/v9 安装与连接

安装依赖

通过 Go Modules 引入 go-redis/v9:

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

基础单机连接

最简单的使用方式是创建 redis.Client 实例并连接到单机 Redis。go-redis/v9 自动管理连接池,默认行为开箱即用,但生产环境建议显式配置。

package main

import (
	"context"
	"fmt"
	"time"
	"github.com/redis/go-redis/v9"
)

func main() {
	ctx := context.Background()
	rdb := redis.NewClient(&redis.Options{
		Addr:         "localhost:6379",
		Password:     "",
		DB:           0,
		PoolSize:     10,
		MinIdleConns: 3,
		MaxRetries:   3,
		DialTimeout:  5 * time.Second,
		ReadTimeout:  3 * time.Second,
		WriteTimeout: 3 * time.Second,
		PoolTimeout:  4 * time.Second,
	})
	if err := rdb.Ping(ctx).Err(); err != nil {
		panic(fmt.Sprintf("redis ping failed: %v", err))
	}
	fmt.Println("redis connected")
	defer rdb.Close()
}

PoolSize 是连接池最大连接数。go-redis/v9 对每个连接维护独立读取通道,因此并发安全。PoolTimeout 控制从连接池获取连接的最长等待时间;如果池满且超时,会返回 redis: connection pool timeout 错误,生产环境需要对此做降级处理。

集群连接 ClusterClient

当数据量或请求量超过单机承载能力时,需要使用 Redis Cluster。go-redis/v9 提供了 redis.ClusterClient,能够自动处理 MOVED/ASK 重定向和槽位映射。

package main

import (
	"context"
	"fmt"
	"time"
	"github.com/redis/go-redis/v9"
)

func main() {
	ctx := context.Background()
	rdb := redis.NewClusterClient(&redis.ClusterOptions{
		Addrs: []string{
			"192.168.1.10:6379",
			"192.168.1.11:6379",
			"192.168.1.12:6379",
			"192.168.1.13:6379",
			"192.168.1.14:6379",
			"192.168.1.15:6379",
		},
		Password:       "cluster-password",
		PoolSize:       20,
		MinIdleConns:   5,
		ReadOnly:       true,
		RouteRandomly:  true,
		MaxRetries:     3,
		DialTimeout:    5 * time.Second,
		PoolTimeout:    4 * time.Second,
	})
	if err := rdb.Ping(ctx).Err(); err != nil {
		panic(err)
	}
	err := rdb.Set(ctx, "user:1001", "alice", 10*time.Minute).Err()
	if err != nil {
		fmt.Println("set error:", err)
	}
	defer rdb.Close()
}

集群模式下有两个关键限制:事务与 Pipeline 的键必须位于同一槽位。可通过 Hash Tag(如 {user}:1001{user}:profile)强制同槽。多键命令如 MGET、MSET 也要求同槽。

哨兵连接 SentinelClient

Redis Sentinel 为高可用架构提供故障自动转移。go-redis/v9 的 redis.NewFailoverClient 会连接 Sentinel 并自动发现主节点地址。

func sentinelDemo(ctx context.Context) {
	rdb := redis.NewFailoverClient(&redis.FailoverOptions{
		MasterName: "mymaster",
		SentinelAddrs: []string{
			"192.168.1.20:26379",
			"192.168.1.21:26379",
			"192.168.1.22:26379",
		},
		Password:         "redis-password",
		SentinelPassword: "sentinel-password",
		DB:               0,
		PoolSize:         10,
		MinIdleConns:     3,
		MaxRetries:       3,
		DialTimeout:      5 * time.Second,
		PoolTimeout:      4 * time.Second,
	})
	if err := rdb.Ping(ctx).Err(); err != nil {
		panic(err)
	}
	defer rdb.Close()
}

Sentinel 客户端内部维护 Sentinel 连接池,定期查询主节点地址。故障转移后客户端会在下一次命令执行时检测到连接错误,然后重新获取新主节点。配合 MaxRetries 与合理的超时时间,可实现平滑过渡。

连接池监控与调优

go-redis/v9 内置连接池统计接口:

func logPoolStats(rdb *redis.Client) {
	stats := rdb.PoolStats()
	fmt.Printf("hits=%d misses=%d timeouts=%d total_conns=%d idle_conns=%d stale_conns=%d\n",
		stats.Hits, stats.Misses, stats.Timeouts,
		stats.TotalConns, stats.IdleConns, stats.StaleConns)
}

关键指标:Hits 应远大于 MissesTimeouts 非零意味着需要扩容连接池或降低并发;StaleConns 持续增加说明服务端空闲超时小于客户端保活周期。高并发短连接场景应增大 PoolSize 并设置 MinIdleConns 预热连接。

基本操作:GET/SET/EXPIRE、Hash/List/Set/ZSet

go-redis/v9 对所有 Redis 数据类型提供完整 API 封装,返回值以 *redis.Cmd 派生类型承载,每个命令都返回 error,必须显式检查。

字符串操作

func stringOperations(ctx context.Context, rdb *redis.Client) {
	err := rdb.Set(ctx, "key:string", "hello redis", 10*time.Minute).Err()
	if err != nil {
		panic(err)
	}
	ok, err := rdb.SetNX(ctx, "key:lock", "locked", 30*time.Second).Result()
	if err != nil {
		panic(err)
	}
	fmt.Println("setnx result:", ok)

	val, err := rdb.Get(ctx, "key:string").Result()
	if err == redis.Nil {
		fmt.Println("key does not exist")
	} else if err != nil {
		panic(err)
	} else {
		fmt.Println("value:", val)
	}

	err = rdb.MSet(ctx, "k1", "v1", "k2", "v2", "k3", "v3").Err()
	if err != nil {
		panic(err)
	}
	vals, err := rdb.MGet(ctx, "k1", "k2", "k3").Result()
	if err != nil {
		panic(err)
	}
	for i, v := range vals {
		fmt.Printf("k%d = %v\n", i+1, v)
	}

	rdb.Set(ctx, "counter:views", 0, 0)
	newVal, err := rdb.Incr(ctx, "counter:views").Result()
	if err != nil {
		panic(err)
	}
	fmt.Println("counter after incr:", newVal)

	rdb.Expire(ctx, "key:string", 5*time.Minute)
	ttl, err := rdb.TTL(ctx, "key:string").Result()
	if err != nil {
		panic(err)
	}
	fmt.Println("ttl:", ttl)
}

特别注意 redis.Nil 的处理。当键不存在时,GET 返回的不是 Go 的 nil,而是 redis.Nil 哨兵错误,方便调用方区分"键不存在"与"网络/协议错误"。

Hash 操作

Hash 适合存储对象属性,每个 Hash 最多可容纳约 40 亿个字段。

func hashOperations(ctx context.Context, rdb *redis.Client) {
	err := rdb.HSet(ctx, "user:1001", "name", "Alice").Err()
	if err != nil {
		panic(err)
	}
	err = rdb.HSet(ctx, "user:1001", map[string]interface{}{
		"age": 28, "email": "alice@example.com", "country": "CN",
	}).Err()
	if err != nil {
		panic(err)
	}
	name, err := rdb.HGet(ctx, "user:1001", "name").Result()
	if err != nil {
		panic(err)
	}
	fmt.Println("name:", name)

	fields, err := rdb.HGetAll(ctx, "user:1001").Result()
	if err != nil {
		panic(err)
	}
	for k, v := range fields {
		fmt.Printf("  %s = %s\n", k, v)
	}

	vals, err := rdb.HMGet(ctx, "user:1001", "name", "age", "gender").Result()
	if err != nil {
		panic(err)
	}
	keys := []string{"name", "age", "gender"}
	for i, v := range vals {
		fmt.Printf("  %s = %v\n", keys[i], v)
	}

	newAge, err := rdb.HIncrBy(ctx, "user:1001", "age", 1).Result()
	if err != nil {
		panic(err)
	}
	fmt.Println("new age:", newAge)

	rdb.HDel(ctx, "user:1001", "temporary_field")
	exists, err := rdb.HExists(ctx, "user:1001", "email").Result()
	if err != nil {
		panic(err)
	}
	fmt.Println("email exists:", exists)
}

HGetAll 返回 map[string]string,便于直接映射到 Go 结构体。如果需要更复杂的类型转换,可结合 encoding/json 将对象序列化为字符串后存入 String 类型,但这会丧失 Hash 的字段级操作能力。

List 操作

List 是有序字符串集合,支持从两端 push/pop,适合实现队列和栈。

func listOperations(ctx context.Context, rdb *redis.Client) {
	key := "list:messages"
	rdb.LPush(ctx, key, "msg3", "msg2", "msg1")
	rdb.RPush(ctx, key, "msg4", "msg5")

	length, err := rdb.LLen(ctx, key).Result()
	if err != nil {
		panic(err)
	}
	fmt.Println("list length:", length)

	items, err := rdb.LRange(ctx, key, 0, -1).Result()
	if err != nil {
		panic(err)
	}
	fmt.Println("all items:", items)

	left, err := rdb.LPop(ctx, key).Result()
	if err != nil {
		panic(err)
	}
	fmt.Println("popped from left:", left)

	result, err := rdb.BLPop(ctx, 10*time.Second, key).Result()
	if err == redis.Nil {
		fmt.Println("blpop timeout")
	} else if err != nil {
		panic(err)
	} else {
		fmt.Println("blpop result:", result)
	}

	rdb.LTrim(ctx, key, 0, 99)
	item, err := rdb.LIndex(ctx, key, 0).Result()
	if err != nil {
		panic(err)
	}
	fmt.Println("index 0:", item)
}

BLPoptime.Duration 参数由 go-redis/v9 内部自动转换为秒。0 表示无限阻塞,设计系统时需明确设置合理上限避免 goroutine 永久挂起。

Set 操作

Set 是无序且唯一的字符串集合,支持交并差运算。

func setOperations(ctx context.Context, rdb *redis.Client) {
	key1 := "set:tags:article1"
	key2 := "set:tags:article2"
	rdb.SAdd(ctx, key1, "golang", "redis", "tutorial")
	rdb.SAdd(ctx, key2, "redis", "database", "performance")

	members, err := rdb.SMembers(ctx, key1).Result()
	if err != nil {
		panic(err)
	}
	fmt.Println("members:", members)

	isMember, err := rdb.SIsMember(ctx, key1, "redis").Result()
	if err != nil {
		panic(err)
	}
	fmt.Println("is member:", isMember)

	card, err := rdb.SCard(ctx, key1).Result()
	if err != nil {
		panic(err)
	}
	fmt.Println("cardinality:", card)

	inter, err := rdb.SInter(ctx, key1, key2).Result()
	if err != nil {
		panic(err)
	}
	fmt.Println("intersection:", inter)

	union, err := rdb.SUnion(ctx, key1, key2).Result()
	if err != nil {
		panic(err)
	}
	fmt.Println("union:", union)

	diff, err := rdb.SDiff(ctx, key1, key2).Result()
	if err != nil {
		panic(err)
	}
	fmt.Println("diff:", diff)

	popped, err := rdb.SPopN(ctx, key1, 1).Result()
	if err != nil {
		panic(err)
	}
	fmt.Println("popped:", popped)
}

SINTER、SUNION、SDIFF 运算在集群模式下要求所有键位于同一槽位。大集合运算在 Redis 端以 O(N) 执行,仍需谨慎。

Sorted Set (ZSet) 操作

Sorted Set 每个成员关联一个 score,按 score 排序,适合排行榜、范围查询等场景。

func zsetOperations(ctx context.Context, rdb *redis.Client) {
	key := "zset:leaderboard"

	members := []redis.Z{
		{Score: 100, Member: "player:alice"},
		{Score: 85, Member: "player:bob"},
		{Score: 120, Member: "player:charlie"},
		{Score: 95, Member: "player:david"},
	}
	rdb.ZAdd(ctx, key, members...)

	results, err := rdb.ZRangeWithScores(ctx, key, 0, 2).Result()
	if err != nil {
		panic(err)
	}
	fmt.Println("top 3:")
	for _, z := range results {
		fmt.Printf("  %s: %.0f\n", z.Member, z.Score)
	}

	topResults, err := rdb.ZRevRangeWithScores(ctx, key, 0, 0).Result()
	if err != nil {
		panic(err)
	}
	fmt.Println("#1:", topResults)

	rank, err := rdb.ZRank(ctx, key, "player:alice").Result()
	if err != nil {
		panic(err)
	}
	fmt.Println("alice rank:", rank)

	score, err := rdb.ZScore(ctx, key, "player:alice").Result()
	if err != nil {
		panic(err)
	}
	fmt.Println("alice score:", score)

	newScore, err := rdb.ZIncrBy(ctx, key, 10, "player:alice").Result()
	if err != nil {
		panic(err)
	}
	fmt.Println("alice new score:", newScore)

	rdb.ZRem(ctx, key, "player:bob")
	count, err := rdb.ZCount(ctx, key, "90", "110").Result()
	if err != nil {
		panic(err)
	}
	fmt.Println("score in [90, 110]:", count)
}

redis.Z 结构体包含 Score float64Member interface{}。Redis 的 score 是 64 位浮点数,使用 float64 对应非常自然。需要注意浮点精度问题,若业务要求精确数值,可将 score 放大 100 倍以整数形式存储。

Pipeline 批量操作:减少 RTT

Pipeline 允许客户端将多个命令打包一次性发送,服务端按顺序执行后一次性返回所有结果,从而将多次 RTT 压缩为 1 次。go-redis/v9 提供了 Pipeline()Pipelined() 两种方式。

func pipelineDemo(ctx context.Context, rdb *redis.Client) {
	pipe := rdb.Pipeline()
	incr := pipe.Incr(ctx, "pipeline:counter")
	pipe.Expire(ctx, "pipeline:counter", 10*time.Minute)
	pipe.Set(ctx, "pipeline:key1", "value1", 0)
	get := pipe.Get(ctx, "pipeline:key1")

	_, err := pipe.Exec(ctx)
	if err != nil {
		panic(err)
	}
	fmt.Println("incr result:", incr.Val())
	fmt.Println("get result:", get.Val())
	pipe.Close()
}

更推荐使用 Pipelined,它能确保 Pipeline 正确关闭:

func pipelinedDemo(ctx context.Context, rdb *redis.Client) {
	cmds, err := rdb.Pipelined(ctx, func(pipe redis.Pipeliner) error {
		pipe.Set(ctx, "key:a", "value-a", 0)
		pipe.Set(ctx, "key:b", "value-b", 0)
		pipe.Set(ctx, "key:c", "value-c", 0)
		pipe.Get(ctx, "key:a")
		pipe.Get(ctx, "key:b")
		pipe.Get(ctx, "key:c")
		return nil
	})
	if err != nil {
		panic(err)
	}
	for i, cmd := range cmds {
		fmt.Printf("cmd[%d] result: %v\n", i, cmd.String())
	}
	if getCmd, ok := cmds[3].(*redis.StringCmd); ok {
		val, _ := getCmd.Result()
		fmt.Println("key:a =", val)
	}
}

本地实测对比:逐条执行 10000 次 SET 耗时 13 秒,Pipeline 执行仅需 1050 毫秒,性能提升约 30~100 倍。公网环境下差距更加明显。

Pipeline 注意事项:

  1. 不保证原子性:命令顺序执行,但中间失败不会阻止后续命令。
  2. 集群模式下键必须同槽:未使用 Hash Tag 会收到 CROSSSLOT 错误。
  3. 命令过多时分批:每 1000~5000 条执行一次 Pipeline,避免内存占用过高。

事务:WATCH/MULTI/EXEC/UNWATCH、乐观锁

Redis 事务通过 MULTI/EXEC 包裹一组命令顺序执行。但不支持回滚——语法错误会被拒绝,运行时错误其余命令仍继续执行。WATCH 提供乐观锁:监控键在 MULTI 和 EXEC 之间被修改时,EXEC 返回空结果。

go-redis/v9 提供了 TxPipeline()Watch() 两种事务 API。

func txPipelineDemo(ctx context.Context, rdb *redis.Client) {
	rdb.Set(ctx, "account:alice", 1000, 0)
	rdb.Set(ctx, "account:bob", 500, 0)
	pipe := rdb.TxPipeline()
	pipe.DecrBy(ctx, "account:alice", 100)
	pipe.IncrBy(ctx, "account:bob", 100)
	_, err := pipe.Exec(ctx)
	if err != nil {
		panic(err)
	}
	aliceBalance, _ := rdb.Get(ctx, "account:alice").Int64()
	bobBalance, _ := rdb.Get(ctx, "account:bob").Int64()
	fmt.Printf("alice=%d bob=%d\n", aliceBalance, bobBalance)
}

使用 Watch 实现 CAS 乐观锁:

func watchTransfer(ctx context.Context, rdb *redis.Client, from, to string, amount int64) error {
	return rdb.Watch(ctx, func(tx *redis.Tx) error {
		fromBalance, err := tx.Get(ctx, from).Int64()
		if err != nil && err != redis.Nil {
			return err
		}
		if fromBalance < amount {
			return fmt.Errorf("insufficient balance")
		}
		toBalance, err := tx.Get(ctx, to).Int64()
		if err != nil && err != redis.Nil {
			return err
		}
		_, err = tx.TxPipelined(ctx, func(pipe redis.Pipeliner) error {
			pipe.Set(ctx, from, fromBalance-amount, 0)
			pipe.Set(ctx, to, toBalance+amount, 0)
			return nil
		})
		return err
	}, from, to)
}

如果监控的键被修改,tx.TxPipelined().Exec() 返回 redis.TxFailedErr,业务层需要自行实现重试:

func withOptimisticLock(ctx context.Context, rdb *redis.Client,
	keys []string, maxRetry int, fn func(tx *redis.Tx) error) error {
	for i := 0; i < maxRetry; i++ {
		err := rdb.Watch(ctx, fn, keys...)
		if err == nil {
			return nil
		}
		if err == redis.TxFailedErr {
			backoff := time.Duration(1<<i) * time.Millisecond * 10
			if backoff > 500*time.Millisecond {
				backoff = 500 * time.Millisecond
			}
			time.Sleep(backoff)
			continue
		}
		return err
	}
	return fmt.Errorf("max retry exceeded")
}

集群模式下事务也要求所有键位于同一槽位。

Lua 脚本加载与执行

Redis Lua 脚本将多个命令和逻辑封装为原子执行,天然无竞态条件。go-redis/v9 通过 Eval 直接执行,或通过 Script 预加载后用 EVALSHA 优化。

func luaEvalDemo(ctx context.Context, rdb *redis.Client) {
	script := `
		local key = KEYS[1]
		local field = KEYS[2]
		local increment = tonumber(ARGV[1])
		local current = redis.call('HGET', key, field)
		if current == false then
			current = 0
		else
			current = tonumber(current)
		end
		local newVal = current + increment
		redis.call('HSET', key, field, newVal)
		return newVal
	`
	result, err := rdb.Eval(ctx, script, []string{"user:1001", "points"}, "50").Result()
	if err != nil {
		panic(err)
	}
	fmt.Printf("new points: %v\n", result)
}

对于高频调用的脚本,使用 Script 预加载:

var incrByScript = redis.NewScript(`
	local key = KEYS[1]
	local field = KEYS[2]
	local increment = tonumber(ARGV[1])
	local current = redis.call('HGET', key, field)
	if current == false then current = 0 else current = tonumber(current) end
	local newVal = current + increment
	redis.call('HSET', key, field, newVal)
	return newVal
`)

func luaScriptDemo(ctx context.Context, rdb *redis.Client) {
	result, err := incrByScript.Run(ctx, rdb, []string{"user:1001", "points"}, "50").Result()
	if err != nil {
		panic(err)
	}
	fmt.Printf("new points: %v\n", result)
}

NewScript 首次 Run 时发送 SCRIPT LOAD,后续使用 EVALSHA。如果 Redis 重启导致缓存丢失,会自动回退 EVAL

原子性扣减库存示例:

var deductStockScript = redis.NewScript(`
	local stockKey = KEYS[1]
	local orderKey = KEYS[2]
	local userId = ARGV[1]
	local qty = tonumber(ARGV[2])
	local exists = redis.call('SISMEMBER', orderKey, userId)
	if exists == 1 then
		return {-1, "already purchased"}
	end
	local current = tonumber(redis.call('GET', stockKey) or 0)
	if current < qty then
		return {-2, "insufficient stock"}
	end
	redis.call('DECRBY', stockKey, qty)
	redis.call('SADD', orderKey, userId)
	return {current - qty, "success"}
`)

该脚本三步原子执行:幂等检查 -> 库存检查 -> 扣减并记录。返回值以数组形式传递业务状态。

Lua 脚本注意事项

  1. 脚本执行期间阻塞 Redis,应避免长循环处理大量数据。
  2. KEYS 必须在调用时确定,不能动态拼接键名,否则集群模式下无法路由。
  3. 脚本默认无超时,lua-time-limit 仅用于触发慢日志,不会终止脚本。

发布/订阅:Subscribe/PSubscribe

Redis Pub/Sub 实现消息多播,发布者向频道发送消息,所有在线订阅者实时接收。不保存历史,只投递给当前已订阅的客户端。

func pubSubDemo(rdb *redis.Client) {
	ctx := context.Background()
	var wg sync.WaitGroup

	wg.Add(1)
	go func() {
		defer wg.Done()
		pubsub := rdb.Subscribe(ctx, "channel:news")
		defer pubsub.Close()
		_, err := pubsub.Receive(ctx)
		if err != nil {
			fmt.Println("subscribe error:", err)
			return
		}
		ch := pubsub.Channel()
		for msg := range ch {
			fmt.Printf("received from %s: %s\n", msg.Channel, msg.Payload)
		}
	}()

	time.Sleep(100 * time.Millisecond)
	for i := 0; i < 5; i++ {
		rdb.Publish(ctx, "channel:news", fmt.Sprintf("news-%d", i))
		time.Sleep(50 * time.Millisecond)
	}
	time.Sleep(200 * time.Millisecond)
	wg.Wait()
}

模式订阅支持 glob 风格匹配:

func pSubscribeDemo(rdb *redis.Client) {
	ctx := context.Background()
	pubsub := rdb.PSubscribe(ctx, "channel:*")
	defer pubsub.Close()
	_, err := pubsub.Receive(ctx)
	if err != nil {
		panic(err)
	}
	ch := pubsub.Channel()
	go func() {
		for msg := range ch {
			fmt.Printf("pattern=%s channel=%s payload=%s\n",
				msg.Pattern, msg.Channel, msg.Payload)
		}
	}()
	rdb.Publish(ctx, "channel:sports", "goal!")
	rdb.Publish(ctx, "channel:tech", "new release")
}

生产环境中订阅者需支持优雅退出和断线重连:

func runSubscriber(ctx context.Context, rdb *redis.Client, channels ...string) {
	for {
		select {
		case <-ctx.Done():
			return
		default:
		}
		pubsub := rdb.Subscribe(ctx, channels...)
		_, err := pubsub.Receive(ctx)
		if err != nil {
			pubsub.Close()
			time.Sleep(1 * time.Second)
			continue
		}
		ch := pubsub.Channel(redis.WithChannelHealthCheckInterval(30 * time.Second))
		for msg := range ch {
			select {
			case <-ctx.Done():
				pubsub.Close()
				return
			default:
			}
			fmt.Printf("[%s] %s\n", msg.Channel, msg.Payload)
		}
		pubsub.Close()
		time.Sleep(1 * time.Second)
	}
}

Pub/Sub 适用实时通知与配置热更新,不适合要求可靠投递和消费确认的场景(应使用 Redis Streams 或专业消息队列)。

Cache-Aside 封装:带 TTL 的缓存层、缓存穿透防御

Cache-Aside 是最常用的缓存策略:读时先查缓存,命中则返回;未命中则查数据库并回写缓存。写时先更新数据库,再使缓存失效。

Cache 接口定义

type Cache interface {
	Get(ctx context.Context, key string, dest interface{}) (bool, error)
	Set(ctx context.Context, key string, value interface{}, ttl time.Duration) error
	Delete(ctx context.Context, keys ...string) error
	GetOrSet(ctx context.Context, key string, dest interface{}, ttl time.Duration, fn func() (interface{}, error)) error
}

RedisCache 实现

type RedisCache struct {
	client    *redis.Client
	baseTTL   time.Duration
	ttlJitter time.Duration
	nullTTL   time.Duration
	group     singleflight.Group
}

func NewRedisCache(client *redis.Client, baseTTL time.Duration) *RedisCache {
	return &RedisCache{
		client:    client,
		baseTTL:   baseTTL,
		ttlJitter: 30 * time.Second,
		nullTTL:   1 * time.Minute,
	}
}

func (c *RedisCache) effectiveTTL(base time.Duration) time.Duration {
	if c.ttlJitter <= 0 {
		return base
	}
	jitter := time.Duration(rand.Int63n(int64(c.ttlJitter)))
	return base + jitter
}

func (c *RedisCache) Get(ctx context.Context, key string, dest interface{}) (bool, error) {
	data, err := c.client.Get(ctx, key).Result()
	if err == redis.Nil {
		return false, nil
	}
	if err != nil {
		return false, err
	}
	if data == "__NULL__" {
		return true, nil
	}
	err = json.Unmarshal([]byte(data), dest)
	if err != nil {
		return true, fmt.Errorf("unmarshal failed: %w", err)
	}
	return true, nil
}

func (c *RedisCache) Set(ctx context.Context, key string, value interface{}, ttl time.Duration) error {
	var data string
	if value == nil {
		data = "__NULL__"
		ttl = c.nullTTL
	} else {
		b, err := json.Marshal(value)
		if err != nil {
			return fmt.Errorf("marshal failed: %w", err)
		}
		data = string(b)
	}
	return c.client.Set(ctx, key, data, c.effectiveTTL(ttl)).Err()
}

func (c *RedisCache) Delete(ctx context.Context, keys ...string) error {
	return c.client.Del(ctx, keys...).Err()
}

GetOrSet 防穿透与并发安全

func (c *RedisCache) GetOrSet(ctx context.Context, key string, dest interface{}, ttl time.Duration, fn func() (interface{}, error)) error {
	found, err := c.Get(ctx, key, dest)
	if err != nil {
		return err
	}
	if found {
		return nil
	}
	val, err, _ := c.group.Do(key, func() (interface{}, error) {
		innerFound, innerErr := c.Get(ctx, key, dest)
		if innerErr != nil {
			return nil, innerErr
		}
		if innerFound {
			return dest, nil
		}
		data, fnErr := fn()
		if fnErr != nil {
			return nil, fnErr
		}
		c.Set(ctx, key, data, ttl)
		return data, nil
	})
	if err != nil {
		return err
	}
	if val == dest {
		return nil
	}
	b, _ := json.Marshal(val)
	return json.Unmarshal(b, dest)
}

singleflight.Group 确保同一 key 的并发请求只有一个执行 fn。__NULL__ 空值缓存可防御缓存穿透,随机抖动 TTL 避免大量缓存同时过期导致雪崩。

实战项目:Go + Redis 计数器/限流器实现

本节综合前面所学,实现两个生产级组件:全局原子计数器滑动窗口限流器

全局原子计数器

基于 Redis INCR 和 EXPIRE 实现支持多维度统计的计数器,可用于接口调用量、在线人数、点赞计数等。

package counter

import (
	"context"
	"fmt"
	"time"
	"github.com/redis/go-redis/v9"
)

type Counter struct {
	client *redis.Client
	prefix string
}

func NewCounter(client *redis.Client, prefix string) *Counter {
	return &Counter{client: client, prefix: prefix}
}

func (c *Counter) Key(dimension string, t time.Time) string {
	return fmt.Sprintf("%s:%s:%s", c.prefix, dimension, t.Format("20060102"))
}

// Incr 原子增加计数,Pipeline 打包 INCR 和 EXPIRE
func (c *Counter) Incr(ctx context.Context, dimension string, t time.Time) (int64, error) {
	key := c.Key(dimension, t)
	pipe := c.client.Pipeline()
	incr := pipe.Incr(ctx, key)
	pipe.Expire(ctx, key, 48*time.Hour)
	_, err := pipe.Exec(ctx)
	if err != nil {
		return 0, err
	}
	return incr.Val(), nil
}

func (c *Counter) Get(ctx context.Context, dimension string, t time.Time) (int64, error) {
	key := c.Key(dimension, t)
	val, err := c.client.Get(ctx, key).Int64()
	if err == redis.Nil {
		return 0, nil
	}
	return val, err
}

// GetRange 获取一段时间内的汇总
func (c *Counter) GetRange(ctx context.Context, dimension string, start, end time.Time) (int64, error) {
	var total int64
	for t := start; !t.After(end); t = t.Add(24 * time.Hour) {
		val, err := c.Get(ctx, dimension, t)
		if err != nil {
			return 0, err
		}
		total += val
	}
	return total, nil
}

Pipeline 将 INCR 和 EXPIRE 压缩为一次 RTT。TTL 设置为 48 小时,满足次日统计需求又不长期占用内存。维度化设计允许一个 Counter 实例同时管理多个指标。

滑动窗口限流器

固定窗口在边界处会出现突发流量翻倍问题。滑动窗口通过 Sorted Set 记录每个请求的时间戳,精确控制窗口内总量。

package ratelimiter

import (
	"context"
	"fmt"
	"time"
	"github.com/redis/go-redis/v9"
)

type SlidingWindowLimiter struct {
	client *redis.Client
	prefix string
}

func NewSlidingWindowLimiter(client *redis.Client, prefix string) *SlidingWindowLimiter {
	return &SlidingWindowLimiter{client: client, prefix: prefix}
}

func (l *SlidingWindowLimiter) Allow(ctx context.Context, key string, limit int, window time.Duration) (bool, error) {
	redisKey := fmt.Sprintf("%s:%s", l.prefix, key)
	now := time.Now().UnixMilli()
	windowStart := now - window.Milliseconds()

	pipe := l.client.Pipeline()
	pipe.ZRemRangeByScore(ctx, redisKey, "0", fmt.Sprintf("%d", windowStart))
	pipe.ZAdd(ctx, redisKey, redis.Z{Score: float64(now), Member: now})
	countCmd := pipe.ZCard(ctx, redisKey)
	pipe.Expire(ctx, redisKey, window*2)

	_, err := pipe.Exec(ctx)
	if err != nil {
		return false, err
	}
	return countCmd.Val() <= int64(limit), nil
}

实现原理:

  1. Sorted Set 的 score 存储毫秒时间戳,天然有序,支持按范围裁剪。
  2. ZRemRangeByScore 删除窗口起始时间之前的记录。
  3. ZAdd 记录当前请求,ZCard 统计窗口内数量。
  4. Pipeline 将四条命令打包为单次 RTT。

使用示例:

func apiHandler(ctx context.Context, limiter *ratelimiter.SlidingWindowLimiter, userID string) {
	allowed, err := limiter.Allow(ctx, userID, 100, time.Minute)
	if err != nil {
		fmt.Println("limiter error:", err)
		return
	}
	if !allowed {
		fmt.Println("rate limited")
		return
	}
	fmt.Println("request allowed")
}

令牌桶限流器(Lua 原子版)

令牌桶在内存和时间复杂度上更优,且天然支持突发流量平滑。以下用 Lua 脚本实现分布式令牌桶:

var tokenBucketScript = redis.NewScript(`
	local key = KEYS[1]
	local rate = tonumber(ARGV[1])
	local capacity = tonumber(ARGV[2])
	local now = tonumber(ARGV[3])
	local requested = tonumber(ARGV[4])

	local bucket = redis.call('HMGET', key, 'tokens', 'last_time')
	local tokens = tonumber(bucket[1]) or capacity
	local last_time = tonumber(bucket[2]) or now

	local delta = math.max(0, now - last_time)
	local add_tokens = delta * rate / 1000.0
	tokens = math.min(capacity, tokens + add_tokens)

	local allowed = 0
	if tokens >= requested then
		tokens = tokens - requested
		allowed = 1
	end

	redis.call('HMSET', key, 'tokens', tokens, 'last_time', now)
	redis.call('EXPIRE', key, 60)
	return allowed
`)

type TokenBucketLimiter struct {
	client   *redis.Client
	prefix   string
	rate     float64
	capacity int64
}

func NewTokenBucketLimiter(client *redis.Client, prefix string, rate float64, capacity int64) *TokenBucketLimiter {
	return &TokenBucketLimiter{client: client, prefix: prefix, rate: rate, capacity: capacity}
}

func (l *TokenBucketLimiter) Allow(ctx context.Context, key string, tokens int64) (bool, error) {
	redisKey := fmt.Sprintf("%s:%s", l.prefix, key)
	now := time.Now().UnixMilli()
	result, err := tokenBucketScript.Run(ctx, l.client,
		[]string{redisKey},
		l.rate, l.capacity, now, tokens,
	).Int64()
	if err != nil {
		return false, err
	}
	return result == 1, nil
}

关键特性:

  1. 纯 Lua 原子执行,读取令牌、计算新令牌、扣减、写回一步完成,无竞态条件。
  2. 时间由客户端传入,避免服务端与客户端时钟不同步。
  3. 使用 Hash 存储桶状态,内存占用极低。
  4. requested 参数支持按请求权重限流。

实战总结

组件数据结构适用场景精度
全局计数器String (INCR)PV/UV 统计、点赞精确
滑动窗口限流器Sorted SetAPI 限流、防刷精确
令牌桶限流器Hash + Lua平滑流量控制近似

在实际部署时,API 网关层可做粗粒度限流,业务代码中做细粒度限流(如按用户ID限制调用功能次数)。超高并发场景可结合本地缓存做 L1 防御,仅当本地窗口用尽时才访问 Redis。

总结

本文系统讲解了 go-redis/v9 的核心用法,从单机/集群/哨兵连接,到 Hash/List/Set/ZSet 全数据类型操作,再到 Pipeline 批量优化、WATCH 乐观锁事务、Lua 脚本原子执行。随后通过发布订阅展示实时消息能力,通过 Cache-Aside 封装演示了生产级缓存设计,最后以计数器和两种限流器收束,覆盖了 Go + Redis 实战中的高频场景。

go-redis/v9 API 与 Redis 协议高度对应,充分利用了 Go 的 context 与 error 处理特性。使用中需要记住:始终检查 error 并区分 redis.Nil;高延迟网络中使用 Pipeline 减少 RTT;需"先读后写"的竞争操作优先使用 Lua 脚本;集群模式下多键命令确保同槽位;连接池参数根据 QPS 和命令耗时调优并监控 PoolStats。掌握这些要点后,Go + Redis 的组合将成为构建高性能服务端应用的坚实基础。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「database」更多文章

  1. 缓存架构演进之路:从单机 Redis 到亿级分布式多级缓存体系
  2. Redis 7.x 重大新特性与架构升级深度解析
  3. Redis 消息队列深度对比:Pub/Sub、Streams 与 Kafka/RabbitMQ 选型指南