客户端与命令管线

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,把 CommandEncoderCommandHandler 等装进 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 直接返回,同一条连接、同一套管线,唯一区别是谁在等。

代码块JAVA · 3 行收起展开
回包阶段:`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 失效——前者堵住所有人回包,后者的状态属于连接不属于线程,这类场景仍要独占连接。

延伸阅读