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 应远大于 Misses;Timeouts 非零意味着需要扩容连接池或降低并发;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)
}
BLPop 的 time.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 float64 和 Member 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 注意事项:
- 不保证原子性:命令顺序执行,但中间失败不会阻止后续命令。
- 集群模式下键必须同槽:未使用 Hash Tag 会收到
CROSSSLOT错误。 - 命令过多时分批:每 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 脚本注意事项
- 脚本执行期间阻塞 Redis,应避免长循环处理大量数据。
- KEYS 必须在调用时确定,不能动态拼接键名,否则集群模式下无法路由。
- 脚本默认无超时,
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
}
实现原理:
- Sorted Set 的 score 存储毫秒时间戳,天然有序,支持按范围裁剪。
ZRemRangeByScore删除窗口起始时间之前的记录。ZAdd记录当前请求,ZCard统计窗口内数量。- 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
}
关键特性:
- 纯 Lua 原子执行,读取令牌、计算新令牌、扣减、写回一步完成,无竞态条件。
- 时间由客户端传入,避免服务端与客户端时钟不同步。
- 使用 Hash 存储桶状态,内存占用极低。
requested参数支持按请求权重限流。
实战总结
| 组件 | 数据结构 | 适用场景 | 精度 |
|---|---|---|---|
| 全局计数器 | String (INCR) | PV/UV 统计、点赞 | 精确 |
| 滑动窗口限流器 | Sorted Set | API 限流、防刷 | 精确 |
| 令牌桶限流器 | 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 的组合将成为构建高性能服务端应用的坚实基础。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。