Zer0e's Blog

为什么热key能拖垮整个Redis集群?Lettuce Pipeline踩坑记

字数统计: 4.3k阅读时长: 16 min
2026/08/30 Share

前言

最近也是遇到了一个线上问题,比较严重,因此和AI对话了几轮,让他帮我总结下问题。有了AI之后,好像很多系统的工程能力都在减弱。我还是维持我的观点,AI是好,但是他只能帮助你50%的工作,架构师、开发的工作是理解业务、设计整个方案、避免踩坑,然后再把具体的实现抛给AI去做。来看遇到的问题吧。

Lettuce 是 Java 生态事实上的标准 Redis 客户端,Pipeline 是批量场景下常用的优化——两者联手,吞吐提升一到两个数量级并不夸张。

可当这对黄金搭档走进云 Redis(LB + Proxy + 分片集群)时,翻车了:一个分片上的热 key,竟然把整条 pipeline 里的所有命令全部拖慢——命令明明分散在不同的分片上,健康分片上的命令凭什么跟着陪葬?

这篇文章把这个坑完整复盘一遍:

  • Lettuce 的使用方式与异步模型;
  • Pipeline 的本质、适用场景与常见误解;
  • 云 Redis 架构下,一个分片的异常如何扩散成全局异常;
  • 多连接攒批的破局方案,附一套可直接落地的代码。

Lettuce 客户端基础

Lettuce 是什么

Java 生态里有两大主流 Redis 客户端:Jedis 和 Lettuce。

维度 Jedis Lettuce
IO 模型 同步阻塞(BIO) 异步非阻塞(Netty/NIO)
线程安全 连接非线程安全,必须用连接池 单连接线程安全(命令按序发送)
编程模型 只有同步 同步 / 异步(Future)/ 响应式(Reactive)
集群支持 依赖连接池 原生支持,自动路由
默认选择 老项目常见 Spring Boot 2.x 起默认客户端

Lettuce 的核心特点是单连接多路复用:一条 TCP 连接上可以连续发送多条命令而不必等待响应,响应按发送顺序返回。记住这个设计——它既是 Lettuce 性能优势的来源,也是后面所有坑的根源。

三种编程模型

1
2
3
4
5
6
7
8
9
10
11
12
13
RedisClient client = RedisClient.create("redis://localhost:6379");
StatefulRedisConnection<String, String> connection = client.connect();

// 1. 同步:发一条等一条
RedisCommands<String, String> sync = connection.sync();
String value = sync.get("key");

// 2. 异步:返回 Future,不阻塞
RedisAsyncCommands<String, String> async = connection.async();
RedisFuture<String> future = async.get("key");

// 3. 响应式
RedisReactiveCommands<String, String> reactive = connection.reactive();

一个关键开关:autoFlush

Lettuce 默认开启 autoFlushCommands(true):每条命令写入 Netty 出站缓冲后立即冲刷到网络。如果关掉它:

1
2
3
connection.setAutoFlushCommands(false);  // 命令只写缓冲,不发送
// ... 发送任意多条命令 ...
connection.flushCommands(); // 手动一次性冲刷

命令会先攒在缓冲区,直到你手动调用 flushCommands() 才一次性发到网络。这正是 Pipeline 的底层机制。

Pipeline 是什么

普通模式的瓶颈:RTT

先看没有 pipeline 时的交互过程:

1
2
3
4
5
6
7
客户端                          服务端
|--- GET key1 --------------->|
|<-------------- 响应1 --------|
|--- GET key2 --------------->|
|<-------------- 响应2 --------|
|--- GET key3 --------------->|
|<-------------- 响应3 --------|

N 条命令要 N 次网络往返。同机房 RTT 约 0.5ms、跨机房约 2ms,1000 条命令光等网络就要 500ms ~ 2s,而 Redis 执行这 1000 条简单命令可能只要几毫秒——时间全花在路上,CPU 全在干等。

Pipeline 模式:批量发送

1
2
3
4
5
客户端                          服务端
|--- GET key1 --------------->|
|--- GET key2 --------------->|
|--- GET key3 --------------->|
|<-- 响应1 + 响应2 + 响应3 ----|

N 条命令一次性发出,服务端依次执行后把 N 个响应打包返回。网络开销从 N × RTT 降到约 1 × RTT。

flowchart LR
    A["普通模式
N 条命令 = N 次 RTT"] --> C["吞吐被网络延迟锁死"] B["Pipeline 模式
N 条命令 ≈ 1 次 RTT"] --> D["吞吐提升 10~50 倍"] style A fill:#ff8a65,color:#fff style B fill:#66bb6a,color:#fff style C fill:#ff8a65,color:#fff style D fill:#66bb6a,color:#fff

两个必须澄清的误解

1. Pipeline 不是事务

它只是”连续发送”,命令之间可能插入其他客户端的命令。要原子性请用 MULTI/EXEC 或 Lua 脚本。

2. Pipeline 会改变延迟特征

pipeline 用”批次的尾延迟”换”整体吞吐”。单条命令的响应时间反而可能变长(要等整批发出、等前面的命令执行完),所以它不适合对单条延迟极度敏感的场景。

Lettuce 中的两种 pipeline 姿势

姿势一:关闭自动 flush,手动批量发送(官方推荐的标准姿势)

1
2
3
4
5
6
7
8
9
10
11
12
StatefulRedisConnection<String, String> conn = client.connect();
conn.setAutoFlushCommands(false);

RedisAsyncCommands<String, String> async = conn.async();
List<RedisFuture<?>> futures = new ArrayList<>();
for (String key : keys) {
futures.add(async.get(key)); // 只是写入 Netty 出站缓冲
}
conn.flushCommands(); // 一次性冲刷到网络
LettuceFutures.awaitAll(timeout, futures.toArray(new RedisFuture[0]));

conn.setAutoFlushCommands(true); // 用完记得恢复

姿势二:直接用 async API 连续发命令

Lettuce 天然允许”不等响应就发下一条”,连续调用 async 命令收集 Future,本质上就是 pipeline 行为,只是少了”先攒后发”的批量控制。

Pipeline 的适用场景

核心判断标准一句话:一次业务动作需要执行多条彼此无依赖的 Redis 命令,且网络延迟占大头。

场景 说明
批量读写 一次页面渲染查几百个 key;数据迁移、缓存预热批量写入
批量计数器 限流/配额扣减,一批请求的 INCRBY 攒起来一次发
混合命令批处理 MGET/MSET 只能做单一操作,pipeline 可以混着 GET、SET、HINCRBY、EXPIRE 一起发
高延迟链路 跨机房/跨可用区 RTT 达到毫秒级时,收益最大

本次我遇到的问题就是批量计数器场景。

不适合的场景:

场景 原因
命令间有依赖 下一条命令的参数依赖上一条的结果,没法攒批
需要原子性 pipeline 不是事务
含阻塞命令(BLPOP 等) 会卡住整条连接上的后续命令,必须用独立连接单独发
单条命令低频调用 没有批可攒,没意义

经验法则:命令数几十到几千、无依赖、可容忍批次级延迟——必上 pipeline;命令数小于 10 且同机房,逐条发通常就够了。

一条命令慢,会卡住后面所有命令

在聊云 Redis 之前,先搞懂 pipeline 在单实例上的行为——这一节是后面所有推理的地基。

Redis 命令执行是单线程、严格按到达顺序串行的,响应也严格按请求顺序返回。所以只要一条命令”卡住”,它后面的所有命令都被压住——这就是队头阻塞(head-of-line blocking)。

“卡住”有三种典型情况:

情况 例子 后果
命令执行慢 KEYS *、大集合 SORT、几百万 field 的 HGETALL 同连接后面的命令排队等待,甚至其他客户端也被拖慢
阻塞型命令 BLPOP key 0、XREAD BLOCK 连接挂起等数据,后续命令一直留在队列里
网络/服务端假死 TCP 拥塞、服务端停止读 socket 该连接上所有已发未回命令全部卡住

关键结论只有一句话:队头阻塞的作用单位是单条连接。

proxy+分片架构下,一个分片异常扩散为全局异常

看似合理的生产架构

我们的生产环境是典型的云 Redis 代理架构:

1
Lettuce ──单连接──→ LB ──→ Redis Proxy × 2 ──→ 所有分片
  • Lettuce 用一条连接连 LB(对客户端来说就像连一台单机);
  • LB 把连接分发到后端两个 Redis Proxy 之一;
  • Proxy 按 key 路由到具体分片,把 N 个分片伪装成一个整体。

分层清晰、职责明确,看起来毫无破绽。问题就藏在”单连接”这三个字里。

Lettuce 直连 Cluster:分片隔离是有效的

先做个对照。如果用 Lettuce 直连 Redis Cluster(不经过 proxy),Lettuce 会为每个节点维护一条独立连接,flush 时按 slot 路由:

1
2
3
4
这批命令
├── 落分片A的命令 ──→ 连接A ──→ 实例A ✅ 正常返回
├── 落分片B(热)的命令 ─→ 连接B ─→ 实例B ❌ 队头阻塞只发生在这里
└── 落分片C的命令 ──→ 连接C ──→ 实例C ✅ 正常返回

热分片 B 上的命令互相排队,但分片 A、C 的命令走各自的连接,完全不受影响。队头阻塞被牢牢限制在”热分片那条连接”内,分片隔离性是有效的。

代理架构:隔离性在最外层被丢掉了

换成 单连接 → LB → Proxy 的架构,故事走向完全不同。数据流要分两个阶段看,魔鬼藏在第二个阶段里:

阶段一:Proxy 向后端分发——执行是并行的

Proxy 收到一批命令后,按 key 路由并行转发到各分片的后端连接。所以健康分片上的命令照常快速执行,这一步没有被热分片阻塞。

阶段二:Proxy 向客户端返回——响应必须按请求顺序

这是致命的一环。RESP 协议规定:单条连接上,响应顺序必须和请求顺序一致。于是 Proxy 只能这样做:

1
2
3
4
5
命令1 → 分片B(热,500ms 才返回)
命令2 → 分片A(1ms 完成) ← 响应已拿到,但不能回
命令3 → 分片C(1ms 完成) ← 同上,压着
...
必须等命令1的响应先返回,2、3 才能依次返回

健康分片的命令执行得再快也没用——它们的响应被热分片的慢响应压在 Proxy 的缓冲区里。这就是响应层的队头阻塞,和 HTTP/1.1 pipelining 的老问题一模一样。

flowchart TD
    A["Lettuce 单连接 pipeline"] --> LB["LB(连接粘性,只命中一个 proxy)"]
    LB --> P["Redis Proxy"]
    P --> S1["分片A:健康"]
    P --> S2["分片B:热 key,慢"]
    P --> S3["分片C:健康"]

    S1 --> R1["响应快,但必须排队等 S2"]
    S2 --> R2["慢响应排在最前"]
    S3 --> R3["响应快,但必须排队等 S2"]

    R2 --> OUT["按请求顺序返回 → 整批被拖慢"]
    R1 --> OUT
    R3 --> OUT

    style A fill:#42a5f5,color:#fff
    style S2 fill:#ff8a65,color:#fff
    style OUT fill:#ef5350,color:#fff

三种架构的伤害范围对比

架构 慢命令的影响范围
单实例 + pipeline 全部卡(一条连接、一个执行线程)
Lettuce 直连 Cluster + pipeline 只卡热分片那条连接上的命令
单连接 → LB → Proxy → 分片 整条连接的所有命令都被拖住(执行并行,返回串行排队)

别漏掉第二个隐藏问题:L4 LB 是按连接分发的。Lettuce 只有一条连接,就会永远粘在同一个 Proxy 上——你以为挂了两台 Proxy 做负载均衡,实际上第二台只是待命的备胎,半点流量都分不到。

连锁反应

热分片一旦出现,这套架构就会开始循环播放下面的剧本:

  1. 热 key 导致分片 B 的 RT 升高;
  2. 落在分片 B 的命令响应变慢;
  3. Proxy 为了保序,压住所有后续命令的响应;
  4. 整条客户端连接的批次完成时间被拉高;
  5. 如果串行发批(发一批等一批),慢批次连环拖累后续批次;
  6. 期间 Lettuce 还在往这条连接攒新命令,Proxy 响应缓冲持续增长;
  7. 监控表现:所有命令的 RT 一起变高,看起来像”Redis 整体异常”,实际上只有一个分片出了问题。

这就是这个坑最迷惑人的地方:监控上全局异常,根因却只是一个热 key。不知道这套机制的人,排查方向很容易跑偏到网络、LB、Proxy 本身——明明所有指标都在喊”我慢了”,真凶却只是集群角落里一个不起眼的 key。

破局

思路

问题的本质是”顺序收敛点”太集中:所有命令的响应都要挤在一条连接上排队。破局方向只有一个——把单连接拆成多连接,把一个全局顺序点打散成多个局部顺序点:

  • 多条连接经 LB 自然分散到两个 Proxy,两台机器真正同时干活;
  • 每条连接各自保序,队头阻塞被限制在”单条连接的那一批”内;
  • 热分片再慢,也只能拖累恰好押中它的那条连接,其他批次照跑。

设计要点

  1. 预建 N 条连接做成池,每次攒批任务独占借出一条连接(setAutoFlushCommands(false) 是连接级状态,不能共享);
  2. 用 BlockingQueue 做简易池,借不到就等待——天然形成背压;
  3. 归还前必须恢复 setAutoFlushCommands(true);
  4. 设置命令超时兜底,避免热分片把 Future 无限压住。

完整代码

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
public class RedisPipelineManager implements Closeable {

private final RedisClient client;
private final List<StatefulRedisConnection<String, String>> connections;
private final BlockingQueue<StatefulRedisConnection<String, String>> pool;
private final Duration batchTimeout;

public RedisPipelineManager(String redisUri, int poolSize, Duration batchTimeout) {
this.client = RedisClient.create(redisUri);
this.batchTimeout = batchTimeout;

// 开启命令级超时,避免慢分片把 Future 无限压住
client.setOptions(ClientOptions.builder()
.timeoutOptions(TimeoutOptions.enabled(batchTimeout))
.build());

this.connections = new ArrayList<>(poolSize);
this.pool = new LinkedBlockingQueue<>(poolSize);
for (int i = 0; i < poolSize; i++) {
StatefulRedisConnection<String, String> conn = client.connect();
connections.add(conn);
pool.offer(conn);
}
}

/**
* 执行一批命令。commandBuilder 只负责往 futures 里添加命令,
* 不要自己 flush 或等待。
*/
public <T> List<T> executePipelined(
PipelineCommandBuilder<T> commandBuilder) throws InterruptedException {

// 借连接:借不到说明并发批次已达上限,自然形成背压
StatefulRedisConnection<String, String> conn =
pool.poll(batchTimeout.toMillis(), TimeUnit.MILLISECONDS);
if (conn == null) {
throw new IllegalStateException("no available connection for pipeline");
}

try {
conn.setAutoFlushCommands(false);
RedisAsyncCommands<String, String> async = conn.async();

List<RedisFuture<T>> futures = new ArrayList<>();
commandBuilder.build(async, futures);
if (futures.isEmpty()) {
return Collections.emptyList();
}

conn.flushCommands();

RedisFuture<?>[] arr = futures.toArray(new RedisFuture[0]);
boolean done = LettuceFutures.awaitAll(batchTimeout, arr);
if (!done) {
throw new RuntimeException("pipeline batch timeout");
}

// 逐条收集结果,部分失败不影响其他结果
List<T> results = new ArrayList<>(futures.size());
for (RedisFuture<T> f : futures) {
try {
results.add(f.get());
} catch (ExecutionException e) {
results.add(null); // 或抛自定义异常,按业务定
}
}
return results;

} finally {
conn.setAutoFlushCommands(true); // 必须恢复
pool.offer(conn); // 归还连接
}
}

@FunctionalInterface
public interface PipelineCommandBuilder<T> {
void build(RedisAsyncCommands<String, String> async, List<RedisFuture<T>> futures);
}

@Override
public void close() {
connections.forEach(StatefulRedisConnection::close);
client.shutdown();
}
}

使用方式:

1
2
3
4
5
6
7
8
9
RedisPipelineManager pipeline =
new RedisPipelineManager("redis://lb-host:6379", 8, Duration.ofMillis(500));

List<String> keys = collectKeys();
pipeline.executePipelined((async, futures) -> {
for (String key : keys) {
futures.add(async.incr(key));
}
});

线程池绑定模式(更简单的替代方案)

如果 pipeline 调用本来就发生在固定的工作线程池里,可以不做池化,直接用 ThreadLocal 绑定:

1
2
private static final ThreadLocal<StatefulRedisConnection<String, String>> CONN =
ThreadLocal.withInitial(() -> redisClient.connect());

每个工作线程终身使用自己那条连接,攒批时关 autoFlush → 发命令 → flush → awaitAll → 恢复。没有借还开销,连接数等于线程数,经 LB 自然分散。

两种模式选一个:请求来源杂、并发不可控就用队列池;有固定工作线程池就用 ThreadLocal 绑定。

参数与注意事项

事项 建议
连接池大小 8~16,按 Proxy 数量和吞吐压测确定
单批大小 500~1000 条,太大时 Proxy 缓冲压力和客户端内存都会上来
批次超时 500ms~1s,配合 TimeoutOptions 双重兜底
Spring 项目注意 别用 RedisTemplate.executePipelined,它走的是共享单连接,等于又把多连接收敛回一条

治标更要治本:热key治理

多连接只是把爆炸半径缩小,热 key 这个火药桶还是要拆。按性价比排序:

手段 说明
本地缓存 Caffeine 扛读热点,短 TTL(如 1~5s)就能把热分片的请求量降几个数量级
热 key 拆散 把 key 拆成 key_{0..N} 子 key 散到不同分片,写时随机/全部写,读时聚合
读副本 给热分片挂只读副本,把读流量打散
客户端限流 对热 key 的访问做速率限制,避免 pipeline 把请求成倍放大

另外,代理架构下的监控值得加一层:按分片维度采集 RT,而不是只看整体均值。热分片的异常被整体均值一稀释,等你从大盘曲线上看出不对劲时,业务往往已经开始报障了。

总结

结论 要点
Pipeline 的本质 用批次的尾延迟换整体吞吐,消灭的是 RTT 开销
队头阻塞的单位 单条连接,不是整个批次
直连 Cluster 按节点拆连接,分片间隔离有效,慢只慢热分片
单连接代理架构 执行并行、返回串行排队,热分片的慢会扩散到整条连接
破局方向 多连接并发攒批,把全局顺序点打散成局部顺序点
根因治理 热 key 本地缓存 / 拆散 / 读副本,监控按分片采集

最后用一句话收尾:在云 Redis 的代理架构里,客户端连接数决定了故障的爆炸半径。一条连接,全集群的命运绑在一起;N 条连接,故障被限制在 N 个局部。

Pipeline 没有错,错的是把 Pipeline 架在一条连接上。

CATALOG
  1. 1. 前言
  2. 2. Lettuce 客户端基础
    1. 2.1. Lettuce 是什么
    2. 2.2. 三种编程模型
    3. 2.3. 一个关键开关:autoFlush
  3. 3. Pipeline 是什么
    1. 3.1. 普通模式的瓶颈:RTT
    2. 3.2. Pipeline 模式:批量发送
    3. 3.3. 两个必须澄清的误解
    4. 3.4. Lettuce 中的两种 pipeline 姿势
  4. 4. Pipeline 的适用场景
  5. 5. 一条命令慢,会卡住后面所有命令
  6. 6. proxy+分片架构下,一个分片异常扩散为全局异常
    1. 6.1. 看似合理的生产架构
    2. 6.2. Lettuce 直连 Cluster:分片隔离是有效的
    3. 6.3. 代理架构:隔离性在最外层被丢掉了
    4. 6.4. 三种架构的伤害范围对比
    5. 6.5. 连锁反应
  7. 7. 破局
    1. 7.1. 思路
    2. 7.2. 设计要点
    3. 7.3. 完整代码
    4. 7.4. 线程池绑定模式(更简单的替代方案)
    5. 7.5. 参数与注意事项
  8. 8. 治标更要治本:热key治理
  9. 9. 总结