恒美微站
首页
关于我们
建站服务
主题模板
案例展示
资讯中心
联系我们
Lettuce 异步 API 实战指南:掌握 RedisFuture 与 CompletionStage 的并发编程
首页
资讯中心
/
Lettuce 异步 API 实战指南:掌握 RedisFuture 与 CompletionStage 的并发编程
Lettuce 异步 API 实战指南:掌握 RedisFuture 与 CompletionStage 的并发编程
发布时间:2026/10/12 2:03:47
数据库后端缓存【免费下载链接】lettuce-coreRedis Java client项目地址https://gitcode.com/gh_mirrors/le/lettuce-core点击查看免费下载Lettuceio.lettuce:lettuce-core是基于 Netty 构建的线程安全 Redis 客户端同时提供同步、异步与响应式三种 API。本篇指南聚焦其异步 APIAsync API从异步执行模型与 Pipelining 原理出发结合仓库源码剖析RedisFutureT的创建、消费、同步与错误处理帮助你在 Standalone、Sentinel、Pub/Sub 与 Cluster 场景下编写不阻塞、高吞吐的 Redis 客户端代码。异步 API 的设计动机为什么需要异步异步方法论的核心价值在于更充分地利用系统资源与其让线程阻塞在网络上等待 I/O不如让线程持续执行其他计算任务。Lettuce 之所以天然支持异步是因为它构建在 Netty 之上——Netty 是一个多线程、事件驱动的 I/O 框架Lettuce 的全部通信都以异步方式处理。当底层基础设施已经能够并发处理命令时上层自然而然地获得了异步能力相比之下把一套阻塞式同步软件改造成并发处理系统要困难得多。理解异步执行流程异步允许在传输尚未完成、响应尚未处理之前就继续执行其他处理。在 Lettuce 与 Redis 的语境下这意味着可以连续发出多条命令而不必等待前一条命令完成这种工作模式即著名的 Pipelining。下面用一个双客户端例子说明其工作方式给定客户端A与客户端B客户端A触发命令SET AB客户端B同时触发命令SET CDRedis 收到客户端A的命令Redis 收到客户端B的命令Redis 处理SET AB并向客户端A返回OK客户端A收到响应并将响应存入结果句柄response handleRedis 处理SET CD并向客户端B返回OK客户端B收到响应并将响应存入结果句柄例子中的两个客户端既可以是同一个应用内的两个线程或两条连接也可以是两台物理隔离的客户端。客户端之间可以以独立进程、线程、事件循环event-loop、Actor、纤程fiber等任意形式并发运行。值得注意的是Redis 本身串行处理进入的命令主体上单线程运行命令按到达顺序被执行部分特性会在后文说明。再把上面的简化例子加上程序流细节给定客户端A客户端A触发命令SET AB客户端A使用异步 API因此可以继续执行其他处理Redis 收到客户端A的命令Redis 处理SET AB并向客户端A返回OK客户端A收到响应并存入结果句柄客户端A无需等待即可访问命令结果非阻塞客户端A因为不等待命令结果所以既能做计算工作也能继续发出下一条 Redis 命令一旦响应可用客户端就能立刻使用命令结果。异步对同步 API 的影响理解异步 API 的同时有必要弄清它对同步 API 的影响。同步 API 与异步 API 的总体思路并无不同两者使用完全相同的设施来发起命令并把命令传输给 Redis 服务器唯一区别在于调用方的阻塞行为。阻塞发生在命令级别且只影响命令完成completion那一段——也就是说多个使用同步 API 的客户端可以在同一条连接上同时发起命令彼此不会互相阻塞一旦命令响应被处理完毕同步 API 上的调用即被解除阻塞。给定客户端A与客户端B客户端A在同步 API 上触发命令SET AB并等待结果客户端B同时也在同步 API 上触发命令SET CD并等待结果Redis 收到客户端A的命令Redis 收到客户端B的命令Redis 处理SET AB并向客户端A返回OK客户端A收到响应解除自身程序流阻塞Redis 处理SET CD并向客户端B返回OK客户端B收到响应解除自身程序流阻塞不过为了避免副作用有些场景不应在线程之间共享连接禁用命令后立即冲刷flush-after-command以提升性能时使用阻塞式操作如BLPOP时阻塞命令会在 Redis 端排队等待执行。一条连接被阻塞期间其他连接仍可向 Redis 发命令一旦有命令解除阻塞例如LPUSH或RPUSH命中了该列表被阻塞的连接随即解除并继续事务Transactions场景使用多个数据库multiple databases时。结果句柄RedisFuture 与 CompletionStage异步 API 上的每一次命令调用都会创建一个RedisFutureT它可以被取消cancel、被等待await和被订阅listener。一个CompletableFutureT或RedisFutureT是指向结果的指针——其值在计算完成前是未知的。RedisFutureT提供了同步synchronization与链式chaining两类操作。先看标准 Java 的CompletableFuture基础用法CompletableFutureString future new CompletableFuture(); System.out.println(Current state: future.isDone()); future.complete(my value); System.out.println(Current state: future.isDone()); System.out.println(Got value: future.get());示例输出Current state: false Current state: true Got value: my value给 future 附加 listener 即可实现链式回调。Promise 与 future 常被混用但并非每个 future 都是 promise——promise 保证会回调/通知完成结果这正是其名称由来。一个在 future 完成时被调用的简单 listenerfinal CompletableFutureString future new CompletableFuture(); future.thenRun(new Runnable() { Override public void run() { try { System.out.println(Got value: future.get()); } catch (Exception e) { e.printStackTrace(); } } }); System.out.println(Current state: future.isDone()); future.complete(my value); System.out.println(Current state: future.isDone());值处理逻辑从调用方移入 listener由完成 future 的一方来触发调用。示例输出Current state: false Got value: my value Current state: true上述代码必须处理异常因为调用get()可能抛出异常future 计算过程中产生的异常会被封装进ExecutionException传递另一类可能抛出的异常是InterruptedException——因为get()是阻塞调用阻塞中的线程随时可能被中断例如系统关机时。自 Java 8 起CompletionStageT类型提供了更精细的 future 处理方式它可以消费、转换并构建一条值处理链。上面的代码用 Java 8 风格改写CompletableFutureString future new CompletableFuture(); future.thenAccept(new ConsumerString() { Override public void accept(String value) { System.out.println(Got value: value); } }); System.out.println(Current state: future.isDone()); future.complete(my value); System.out.println(Current state: future.isDone());示例输出Current state: false Got value: my value Current state: trueCompletionStageT的完整方法参考可查阅 Java 8 API 文档中java.util.concurrent.CompletionStage的规范说明本文不再外链仓库内的响应式指南 reactive-api 也大量使用了该类型。Lettuce 的 RedisFuture 接口仓库中RedisFuture的定义位于 RedisFuture.javapublic interface RedisFutureV extends CompletionStageV, FutureV { String getError(); boolean await(long timeout, TimeUnit unit) throws InterruptedException; }从源码结构看RedisFutureV同时继承了CompletionStageV与FutureV因此CompletionStage的全部方法thenApply、thenAccept、handle、exceptionally等在 Lettuce future 上均可用此外它还额外提供getError()返回错误文本若有错误发生await(timeout, unit)最多等待指定时长直到命令输出可用返回boolean仅在等待期间线程被中断时抛出InterruptedException。使用 Lettuce 创建 FutureLettuce 的 future 既可用于初始操作也可用于链式操作。使用时会明显体会到非阻塞行为——因为所有 I/O 与命令处理都由 Netty EventLoop 异步完成。Lettuce 在 Standalone、Sentinel、Publish/Subscribe 与 Cluster 四类 API 上都暴露了 future。连接 Redis 非常简单RedisClient client RedisClient.create(redis://localhost); RedisAsyncCommandsString, String commands client.connect().async();connect().async()返回的是RedisAsyncCommandsK, V接口。从 RedisAsyncCommands.java 的声明可以看到它聚合了字符串、Hash、List、Set、SortedSet、Stream、Geo、Scripting、事务以及 JSON、Vector Set、Search 等数十个命令子接口是包含 400 方法的完整异步、线程安全 Redis API其实现类 RedisAsyncCommandsImpl.java 同时实现了RedisAsyncCommands与RedisClusterAsyncCommands。接下来获取某个 key 的值需要GET操作RedisFutureString future commands.get(key);从 AbstractRedisAsyncCommands.java 的get与set实现如public RedisFutureV get(K key)、public RedisFutureString set(K key, V value)可以看到异步命令一律通过统一的dispatch(...)通道发出先由命令构建器生成协议命令再交给 channel writer 写出最终返回RedisFuture。这也印证了同步 API 与异步 API 共用同一套命令传输设施的论断。消费 Future使用 future 的第一件事就是消费它——即取得其值。下面是阻塞调用线程并打印值的例子RedisFutureString future commands.get(key); String value future.get(); System.out.println(value);get()拉取式pull-style调用会阻塞调用线程至少阻塞到值计算完成最坏情况下可能无限期阻塞。总是使用超时是避免耗尽线程的好习惯try { RedisFutureString future commands.get(key); String value future.get(1, TimeUnit.MINUTES); System.out.println(value); } catch (Exception e) { e.printStackTrace(); }该示例最多等待 1 分钟若超时将抛出TimeoutException以提示超时。future 也可以采用**推式push-style**消费当RedisFutureT完成时自动触发后续动作RedisFutureString future commands.get(key); future.thenAccept(new ConsumerString() { Override public void accept(String value) { System.out.println(value); } });用 Java 8 lambda 简写RedisFutureString future commands.get(key); future.thenAccept(System.out::println);规则永远不要阻塞 EventLoopLettuce 的 future 是在 **Netty EventLoop 上完成complete**的。默认在线程上消费与链式处理 future 总是个好主意但有一个例外阻塞或长时间运行的操作用例。经验法则绝不阻塞事件循环。如果需要在链式中使用阻塞调用请用thenAcceptAsync()/thenRunAsync()把处理分叉到其他线程Executor sharedExecutor ... RedisFutureString future commands.get(key); future.thenAcceptAsync(new ConsumerString() { Override public void accept(String value) { System.out.println(value); } }, sharedExecutor);…async()系列方法需要一个线程基础设施来执行默认使用ForkJoinPool.commonPool()。ForkJoinPool是静态构建的不会随负载增长因此使用自定义的Executor几乎总是更好的选择。Future 的同步使用 future 的关键点是同步。future 通常用于四种目的触发多次调用而不必等待前驱Batching批处理发出命令后完全不等结果Fire Forget发出命令的同时执行其他计算Decoupling解耦为特定计算任务增加并发Concurrency。等待 future 完成或获得完成通知有多种方式不同的同步技术适用于不同的使用动机。阻塞式同步阻塞式同步适合做批处理或为系统局部增加并发。批处理示例批量写入/读取多个值并在处理过程中的某个点之前等待所有结果ListRedisFutureString futures new ArrayListRedisFutureString(); for (int i 0; i 10; i) { futures.add(commands.set(key- i, value- i)); } LettuceFutures.awaitAll(1, TimeUnit.MINUTES, futures.toArray(new RedisFuture[futures.size()]));上述代码不会等某条命令完成后再发下一条同步发生在全部命令发出之后。把对LettuceFutures.awaitAll()的调用省略掉这段代码就轻松变成了 Fire Forget 模式。单个 future 也可以被等待——即选择等待一段特定时间但不抛出异常RedisFutureString future commands.get(key); if(!future.await(1, TimeUnit.MINUTES)) { System.out.println(Could not complete within the timeout); }await()对调用方更友好它只在阻塞线程被中断时抛出InterruptedException。get()方法在前面已介绍过这里不再赘述。最后一种阻塞式同步方式是轮询isDone()。其主要陷阱是你必须自行处理线程中断。如果忽略这一点系统在运行状态下将无法被正常关闭RedisFutureString future commands.get(key); while (!future.isDone()) { // do something ... }isDone()的主要用途并不是同步但在命令执行期间顺便做点其他计算却很实用。LettuceFutures 的底层实现仓库中 LettuceFutures.java 提供了两个实用方法awaitAll(long timeout, TimeUnit unit, Future?... futures)等待所有 future 完成或超时到达。超时后不会取消命令区别于awaitOrCancel返回true表示全部在时限内完成否则返回false该方法还提供了接收Duration的重载since 5.0。awaitOrCancel(RedisFutureT cmd, long timeout, TimeUnit unit)等待命令完成或超时若超时仍未完成则取消该命令。其核心实现在内部工具类 internal/Futures.javaawaitAll逐个对 future 执行带剩余时间预算的get(nanos, TimeUnit.NANOSECONDS)捕获TimeoutException后返回false其余异常经Exceptions.fromSynchronization重新抛出awaitOrCancel在超时后执行cmd.cancel(true)并抛出超时异常。源码中的行为验证了文档的说法awaitAll只等待不取消而awaitOrCancel会在超时时取消未完成命令。链式同步future 可以以非阻塞方式同步/链式组合以提升线程利用率。链式在依赖事件驱动特性的系统中尤其好用future 链构建了一条由多个 future 组成的串行执行链每个链成员负责计算的一部分。CompletionStageTAPI 提供了多种链式与转换方法。用thenApply()做简单的值转换future.thenApply(new FunctionString, Integer() { Override public Integer apply(String value) { return value.length(); } }).thenAccept(new ConsumerInteger() { Override public void accept(Integer integer) { System.out.println(Got value: integer); } });用 Java 8 lambda 简写future.thenApply(String::length) .thenAccept(integer - System.out.println(Got value: integer));thenApply()接收一个把值转换成另一个值的函数最后的thenAccept()消费值做最终处理。前面示例中已见过的thenRun()用于在数据对你的流程不关键时处理 future 完成事件future.thenRun(new Runnable() { Override public void run() { System.out.println(Finished the future.); } });注意若Runnable中要做阻塞调用务必让它在自定义Executor上执行。另一个值得关注的是either-or 链式。CompletionStageT提供了若干…Either()方法完整参考见 Java 8 API 文档要么模式either-or消费最先完成的那个 future的值。典型场景是两个服务返回相同数据例如 Master-Replica 主从场景希望尽可能快地拿到数据RedisStringAsyncCommandsString, String master masterClient.connect().async(); RedisStringAsyncCommandsString, String replica replicaClient.connect().async(); RedisFutureString future master.get(key); future.acceptEither(replica.get(key), new ConsumerString() { Override public void accept(String value) { System.out.println(Got value: value); } });错误处理错误处理是每个真实世界应用不可或缺的组成部分应当从一开始就纳入考量。future 提供了一些错误处理机制通常你会希望以如下方式响应异常返回默认值Return a default value instead使用备用 futureUse a backup future重试该 futureRetry the future。RedisFutureT会携带发生的异常。调用get()时发生的异常会被包装进ExecutionException抛出这与 Lettuce 3.x 不同注意这一版本差异。更多细节见CompletionStage的 Javadoc。下面的代码使用handle()方法在遇到异常后回退到默认值future.handle(new BiFunctionString, Throwable, String() { Override public Integer apply(String value, Throwable throwable) { if(throwable ! null) { return default value; } return value; } }).thenAccept(new ConsumerString() { Override public void accept(String value) { System.out.println(Got value: value); } });更复杂的代码可以根据 throwable 的具体类型决定返回什么值例如下面的exceptionally()快捷写法future.exceptionally(new FunctionThrowable, String() { Override public String apply(Throwable throwable) { if (throwable instanceof IllegalStateException) { return default value; } return other default value; } });需要特别说明重试 future 以及基于 future 的恢复recovery并不属于 Java 8CompletableFutureT的能力范围。如需更舒适地处理异常请参阅仓库中的 Reactive API 指南——响应式 API 对异常组合与重试提供了更顺手的支持。实战示例汇总以下是文档中给出的完整可运行示例覆盖了阻塞拉取、超时等待与推式监听三种消费方式示例 1get()阻塞拉取RedisAsyncCommandsString, String async client.connect().async(); RedisFutureString set async.set(key, value); RedisFutureString get async.get(key); set.get() OK get.get() value示例 2await()与带超时的get()RedisAsyncCommandsString, String async client.connect().async(); RedisFutureString set async.set(key, value); RedisFutureString get async.get(key); set.await(1, SECONDS) true set.get() OK get.get(1, TimeUnit.MINUTES) value示例 3thenRun()推式监听RedisStringAsyncCommandsString, String async client.connect().async(); RedisFutureString set async.set(key, value); Runnable listener new Runnable() { Override public void run() { ...; } }; set.thenRun(listener);深入原理一条命令如何变成 Future结合仓库源码可以看清异步命令的完整生命周期这有助于你更好地预判行为、排查超时与异常命令构建与派发以get为例AbstractRedisAsyncCommands.java 中的get(K key)调用dispatch(commandBuilder.get(key))set同理RedisFutureString set(K key, V value)。底层派发通道定义在 BaseRedisAsyncCommands.java 的dispatch(ProtocolKeyword type, CommandOutputK, V, T output, CommandArgsK, V args)。写出与排队RedisChannelHandler.java 的dispatch(RedisCommand cmd)把命令交给channelWriter.write(cmd)。在 CommandHandler.java 的addToStack(...)中若命令的 output 为null即 Fire Forget 命令会立即 complete 而不进入响应处理这从源码层面印证了 Fire Forget 模式不等待响应的行为否则命令被压入栈stack等待写出并受ClientOptions.getRequestQueueSize()有界队列约束。响应到达与完成当 Netty 的channelRead读到 Redis 响应字节流时CommandHandler解析响应并通过complete(command)见 CommandHandler.java 第 772 行附近的RedisCommand#complete()完成对应命令最终完成其RedisFuture——这也解释了为什么 future 总是在 EventLoop 上完成。Cluster 场景的聚合 future在 Cluster 多节点执行跨槽位命令时Lettuce 使用 PipelinedRedisFuture.javaextends CompletableFutureV implements RedisFutureV把多个节点的执行结果用CompletableFuture.allOf(...)聚合为单个复合结果并通过CountDownLatch保证完成语义。小结Lettuce 异步 API 的核心是RedisFutureT它既是FutureV也是CompletionStageV天然支持拉取get/await/isDone与推式thenAccept/thenApply/thenRun/handle/exceptionally两类消费方式。掌握三个关键原则即可写出高质量代码消费 future 时始终携带超时绝不在 Netty EventLoop 上做阻塞或长耗时操作必要时用带自定义Executor的…async()方法分叉线程把异常处理纳入设计默认值、备用 future或直接转向响应式 API 获得更丰富的重试能力。同步 API 与异步 API 共享同一套传输设施二者可以放心地在同一连接上共存只需避开阻塞命令、事务与多数据库等需要独占连接的场景。赞分享数据库后端缓存【免费下载链接】lettuce-coreRedis Java client项目地址https://gitcode.com/gh_mirrors/le/lettuce-core点击查看免费下载相关推荐轻松掌握Iced异步编程Future集成与并发处理实战指南轻松掌握Iced异步编程Future集成与并发处理实战指南 你还在为Rust GUI应用中的异步操作头疼吗当用户点击按钮后界面卡顿、网络请求阻塞主线程、多个前端跨平台UI组件桌面应用FastAPI并发编程终极指南掌握异步处理的完整实战教程FastAPI并发编程终极指南掌握异步处理的完整实战教程 FastAPI是一个现代、快速高性能的Web框架用于构建API它基于标准Python类型提示后端前端认证鉴权Lettuce-core异步API深度解析与实战指南Lettuce core异步API深度解析与实战指南 异步API概述 Lettuce core作为高性能Redis客户端其异步API设计基于Netty框架充数据库后端缓存上一篇终极指南如何用免费开源工具深度优化AMD Ryzen处理器性能下一篇SMUDebugTool完整指南免费AMD Ryzen处理器调试工具终极教程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考