客户端与命令管线
Lettuce 客户端与命令管线 源码分析
Lettuce 解决的问题:让一条 TCP 连接被任意多线程安全共享,同时天然支持 pipeline。
核心思路是把每条 Redis 命令做成一个”命令对象”(RedisCommand,本质是个 CompletableFuture),写出去时压进 CommandHandler 里的一个队列,收到回包时按 FIFO 顺序逐个出队配对——因为 RESP 协议没有请求 ID,单连接上的回包顺序就是请求顺序,这个约定就是整套管线的地基。
// 基于本地 Lettuce 仓 (JavaSourceReadingLab, 6.x), io/lettuce/core/RedisClient.java
public class RedisClient extends AbstractRedisClient {
public <K, V> StatefulRedisConnection<K, V> connect(RedisCodec<K, V> codec) {
checkForRedisURI();
// 注意:同步 connect 只是 "异步建链 + 阻塞等 future",模式和后面的 sync 命令 API 完全一致
return getConnection(connectStandaloneAsync(codec, this.redisURI, getDefaultTimeout()));
}
private <K, V> ConnectionFuture<StatefulRedisConnection<K, V>> connectStandaloneAsync(RedisCodec<K, V> codec,
RedisURI redisURI, Duration timeout) {
assertNotNull(codec);
checkValidRedisURI(redisURI);
logger.debug("Trying to get a Redis connection for: {}", redisURI);
// Endpoint 是"命令写入口":channel 还没建好/断线重连期间,命令先缓冲在这里,不会直接怼到 socket
DefaultEndpoint endpoint = new DefaultEndpoint(getOptions(), getResources());
RedisChannelWriter writer = endpoint;
// 装饰器链:超时写、监听器写都是包在 endpoint 外面的 writer,职责分离
if (CommandExpiryWriter.isSupported(getOptions())) {
writer = new CommandExpiryWriter(writer, getOptions(), getResources());
}
if (CommandListenerWriter.isSupported(getCommandListeners())) {
writer = new CommandListenerWriter(writer, getCommandListeners());
}
// 先造好连接门面对象(含 sync/async/reactive 三套 API),此时 TCP 还没连
StatefulRedisConnectionImpl<K, V> connection = newStatefulRedisConnection(writer, endpoint, codec, timeout);
ConnectionFuture<StatefulRedisConnection<K, V>> future = connectStatefulAsync(connection, endpoint, redisURI,
() -> new CommandHandler(getOptions(), getResources(), endpoint)); // 每条连接一个 CommandHandler
future.whenComplete((channelHandler, throwable) -> {
if (throwable != null) {
connection.closeAsync(); // 建链失败要把已创建的门面对象关掉,防泄漏
}
});
return future;
}
@SuppressWarnings("unchecked")
private <K, V, S> ConnectionFuture<S> connectStatefulAsync(StatefulRedisConnectionImpl<K, V> connection, Endpoint endpoint,
RedisURI redisURI, Supplier<CommandHandler> commandHandlerSupplier) {
ConnectionBuilder connectionBuilder;
if (redisURI.isSsl()) {
SslConnectionBuilder sslConnectionBuilder = SslConnectionBuilder.sslConnectionBuilder();
sslConnectionBuilder.ssl(redisURI);
connectionBuilder = sslConnectionBuilder;
} else {
connectionBuilder = ConnectionBuilder.connectionBuilder();
}
ConnectionState state = connection.getConnectionState();
state.apply(redisURI);
state.setDb(redisURI.getDatabase());
connectionBuilder.connection(connection);
connectionBuilder.clientOptions(getOptions());
connectionBuilder.clientResources(getResources());
connectionBuilder.commandHandler(commandHandlerSupplier).endpoint(endpoint);
connectionBuilder(getSocketAddressSupplier(redisURI), connectionBuilder, connection.getConnectionEvents(), redisURI);
connectionBuilder.connectionInitializer(createHandshake(state)); // HELLO/AUTH/SELECT 握手也是走命令管线发的
// 真正的 netty Bootstrap.connect 在这里面:组 pipeline(CommandEncoder + CommandHandler 等)、发起 TCP 连接
ConnectionFuture<RedisChannelHandler<K, V>> future = initializeChannelAsync(connectionBuilder);
return future.thenApply(channelHandler -> (S) connection);
}
}// 基于本地 Lettuce 仓 (JavaSourceReadingLab, 6.x), io/lettuce/core/StatefulRedisConnectionImpl.java
public class StatefulRedisConnectionImpl<K, V> extends RedisChannelHandler<K, V> implements StatefulRedisConnection<K, V> {
protected final RedisCodec<K, V> codec;
protected final RedisCommands<K, V> sync; // 三套 API 是同一条连接的三个视图,不是三条连接
protected final RedisAsyncCommandsImpl<K, V> async;
protected final RedisReactiveCommandsImpl<K, V> reactive;
// ...
public StatefulRedisConnectionImpl(RedisChannelWriter writer, PushHandler pushHandler, RedisCodec<K, V> codec,
Duration timeout) {
super(writer, timeout);
this.pushHandler = pushHandler;
this.codec = codec;
this.async = newRedisAsyncCommandsImpl(); // async 是本体
this.sync = newRedisSyncCommandsImpl(); // sync 基于 async 造出来,见下
this.reactive = newRedisReactiveCommandsImpl();
}
protected RedisCommands<K, V> newRedisSyncCommandsImpl() {
// 关键一行:sync API 是 JDK 动态代理,包住 async()。每次调用 = 发异步命令 + 阻塞等 future(超时用连接的 timeout)
// 所以"同步"只是调用方线程在等,连接本身从头到尾都是异步的
return syncHandler(async(), RedisCommands.class, RedisClusterCommands.class);
}
@Override
public RedisCommands<K, V> sync() {
return sync;
}
@Override
public <T> RedisCommand<K, V, T> dispatch(RedisCommand<K, V, T> command) {
RedisCommand<K, V, T> toSend = preProcessCommand(command); // AUTH/SELECT/MULTI 等命令挂回调维护连接状态,重连后能恢复现场
potentiallyEnableMulti(command);
return super.dispatch(toSend); // 交给 RedisChannelWriter(endpoint),最终 channel.write 进 netty pipeline
}
}// 基于本地 Lettuce 仓 (JavaSourceReadingLab, 6.x), io/lettuce/core/protocol/CommandHandler.java
public class CommandHandler extends ChannelDuplexHandler implements HasQueuedCommands {
// 命名叫 stack,实际用法是 FIFO 队列:写出时 add 到尾,回包时 peek/poll 头。
// 无锁 ArrayDeque 就够了——netty 保证同一个 channel 的所有 handler 回调都在同一个 EventLoop 线程上跑
private final ArrayDeque<RedisCommand<?, ?, ?>> stack = new ArrayDeque<>();
// ...
@Override
@SuppressWarnings("unchecked")
public void write(ChannelHandlerContext ctx, Object msg, ChannelPromise promise) throws Exception {
// ...
if (msg instanceof RedisCommand) {
writeSingleCommand(ctx, (RedisCommand<?, ?, ?>) msg, promise);
return;
}
if (msg instanceof List) { // 批量分支:一次 write 多条命令,pipeline 的由来
List<RedisCommand<?, ?, ?>> batch = (List<RedisCommand<?, ?, ?>>) msg;
if (batch.size() == 1) {
writeSingleCommand(ctx, batch.get(0), promise);
return;
}
writeBatch(ctx, batch, promise);
return;
}
if (msg instanceof Collection) {
writeBatch(ctx, (Collection<RedisCommand<?, ?, ?>>) msg, promise);
}
}
private void writeSingleCommand(ChannelHandlerContext ctx, RedisCommand<?, ?, ?> command, ChannelPromise promise) {
if (!isWriteable(command)) { // 已 cancel/超时完成的命令不再写出
promise.trySuccess();
return;
}
addToStack(command, promise); // 先入队再写!否则回包可能比入队先到(都在同一线程,这里其实是保证顺序语义清晰)
attachTracing(ctx, command);
ctx.write(command, promise); // 继续往 pipeline 下游传,由 CommandEncoder 调 command.encode() 生成 RESP 字节
}
private void addToStack(RedisCommand<?, ?, ?> command, ChannelPromise promise) {
try {
if (!ActivationCommand.isActivationCommand(command)) {
validateWrite(1); // 有界队列保护:在途命令超过 requestQueueSize 直接拒绝,防止 OOM
}
if (command.getOutput() == null) {
// fire&forget commands are excluded from metrics and replies
complete(command); // 没有 output 的命令不期待回包,立即完成,不占配对队列
}
RedisCommand<?, ?, ?> redisCommand = potentiallyWrapLatencyCommand(command);
stack.add(redisCommand);
if (!promise.isVoid()) {
// 写 socket 失败时从队列移除该命令,否则后续所有回包都会错位配对——这是 FIFO 方案最危险的坑
promise.addListener(AddToStack.newInstance(stack, redisCommand));
}
} catch (Exception e) {
command.completeExceptionally(e);
throw e;
}
}
@Override
public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception {
ByteBuf input = (ByteBuf) msg;
// ...
try {
// ...
buffer.writeBytes(input); // 攒进内部聚合 buffer:TCP 是字节流,一个回包可能拆在多次 read 里到达
decode(ctx, buffer);
} finally {
input.release();
}
}
protected void decode(ChannelHandlerContext ctx, ByteBuf buffer) throws InterruptedException {
// ... (pristine 分支:无命令上下文时消费流浪回包)
while (canDecode(buffer)) { // 一次 read 可能带回多条命令的回包(pipeline),循环解码
if (isPushDecode(buffer)) {
// ... (RESP3 push 消息不占配对队列,走 pushListener)
} else {
RedisCommand<?, ?, ?> command = stack.peek(); // 只 peek 不 poll:回包可能不完整,命令要留在队头等下一波字节
// ...
pristine = false;
try {
if (!decode(ctx, buffer, command)) { // 内部是 RedisStateMachine 增量解析 RESP,字节不够返回 false
hasDecodeProgress = true;
decodeBufferPolicy.afterPartialDecode(buffer);
return; // 半包:直接退出,等 channelRead 下次续上
}
} catch (Exception e) {
ctx.close(); // 解码异常说明流已错乱,只能断连重来,不能跳过——跳过就全体错位
throw e;
}
hasDecodeProgress = false;
if (isProtectedMode(command)) {
onProtectedMode(command.getOutput().getError());
} else {
if (canComplete(command)) {
stack.poll(); // 解码完整才正式出队
try {
// ...
complete(command); // command.complete() -> CompletableFuture 完成,唤醒等待的业务线程
} catch (Exception e) {
logger.warn("{} Unexpected exception during request: {}", logPrefix, e.toString(), e);
}
}
}
afterDecode(ctx, command);
}
}
decodeBufferPolicy.afterDecoding(buffer);
}
}原理串讲
以 connection.sync().set("k", "v") 走一遍全链路。建链阶段:RedisClient.connect() 调 connectStandaloneAsync(),先 new 出 DefaultEndpoint(命令写入口兼断线缓冲)和 StatefulRedisConnectionImpl(门面,构造时把 async/sync/reactive 三套 API 都造好),再由 connectStatefulAsync() 组装 ConnectionBuilder——socket 地址、SSL、握手命令、CommandHandler 工厂全塞进去,最后 initializeChannelAsync() 里才真正执行 netty Bootstrap.connect,把 CommandEncoder、CommandHandler 等装进 channel pipeline。
注意顺序:连接对象先于 TCP 存在,channel 只是它的可替换部件,这就是断线重连后”连接对象不变、pending 命令重发”的结构基础。
命令阶段:sync API 是 syncHandler() 造的 JDK 动态代理,set() 一进来就被转成对 RedisAsyncCommandsImpl.set() 的调用——后者把命令名、参数、输出解析器打包成 AsyncCommand(继承 CompletableFuture),经 StatefulRedisConnectionImpl.dispatch() 预处理(AUTH/SELECT 挂状态回调、MULTI 包事务)后交给 endpoint,endpoint channel.write() 写入 pipeline。
CommandHandler.write() 在 EventLoop 线程上执行:addToStack() 把命令对象加进 ArrayDeque 队尾,然后 ctx.write() 传给下游 CommandEncoder 编码成 RESP 字节流出网。
此时业务线程在动态代理里阻塞等这个 future;如果用的是 async(),则拿到 future 直接返回,同一条连接、同一套管线,唯一区别是谁在等。
代码块收起展开
回包阶段:`channelRead()` 把字节攒进聚合 buffer 后进 `decode()` 循环——`stack.peek()` 取队头命令,`RedisStateMachine` 按该命令的 `CommandOutput` 增量解析 RESP;半包就返回等下一波,完整了才 `stack.poll()` + `complete(command)`,future 完成,阻塞中的业务线程醒来拿到结果。
为什么敢用 FIFO 盲配?因为 Redis 单连接上严格按请求顺序回包,协议层没有请求 ID,顺序就是唯一的关联依据;也因此写失败必须靠 `AddToStack` 监听器把命令从队列里摘掉,解码异常只能 `ctx.close()`,任何一次错位都会让后续所有配对全错。
为什么多线程能共享一条连接?因为线程间根本不竞争 socket——大家只是并发地把命令对象丢进管线,真正碰 socket 和队列的只有 EventLoop 一个线程,`ArrayDeque` 连锁都不用加;每个线程等的是自己那个命令的 future,互不干扰。对比 Jedis:Jedis 在调用线程上同步读写 socket,连接内部有请求-响应的时序状态,所以必须每线程一连接、靠连接池隔离,而 Lettuce 一条连接就能喂饱整个线程池(阻塞命令和事务除外)。
设计取舍
stack名为栈实为 FIFO 队列(add 尾 / poll 头),读源码别被字段名带偏。- 用”回包顺序 = 请求顺序”的协议约定省掉请求 ID,代价是队列绝不能错位:写失败要摘除、解码异常直接断连,没有中间态。
- sync 只是 async 加阻塞等待的动态代理,一套实现三种 API;代价是 sync 调用也有 future 分配开销,且阻塞的是业务线程不是 IO 线程。
- 单 EventLoop 处理一条连接的全部 IO,队列无锁、天然 pipeline;代价是解码回调里不能干重活,
CommandOutput解析慢会拖住这个 EventLoop 上的所有连接。 - 共享连接对阻塞命令(BLPOP)和 MULTI/EXEC 失效——前者堵住所有人回包,后者的状态属于连接不属于线程,这类场景仍要独占连接。