WR Blog 加载中...
返回文章
JavaJava 25 LTSRedis

Java Redis 工程:Lettuce、连接、序列化、缓存、锁和 Streams

Java Redis 工程:Lettuce、连接、序列化、缓存、锁和 Streams封面

Java 基础体系 · 第 87/100 篇。示例统一以 Java 25 LTS 为语言和 JVM 基线;框架示例使用与其兼容的现代稳定版本。

Java Redis 工程:Lettuce、连接、序列化、缓存、锁和 Streams

Redis 工程问题通常不是“会不会调用 GETSET”,而是要同时处理六件事:

  1. Lettuce 客户端如何建立连接、复用连接并关闭资源;
  2. Redis 字节如何映射为 Java 类型;
  3. 缓存数据如何写入、过期、重建和失效;
  4. 多线程、多实例下的锁是否真的具有互斥和安全释放能力;
  5. Redis Streams 如何投递、确认、重试和恢复消息;
  6. 网络断开、超时、主从切换、进程崩溃时,业务状态如何变化。

本文以 Java 25 LTS 为运行范围,以 Lettuce 的同步、异步和响应式 API 为客户端基础。示例使用 Redis 7.x 常见命令语义;具体 Lettuce 版本应根据项目发布周期选择兼容版本,而不应把 Java 版本直接等同于 Redis 或 Lettuce 版本。


1. Redis 与 Lettuce 的职责边界

1.1 Redis 是服务端数据结构与协议

Redis 是一个通过 RESP 协议提供命令服务的内存数据存储。客户端发送命令:

SET user:42:name Alice EX 60

Redis 返回结果:

OK

服务端负责:

  • 解析命令;
  • 修改内存中的数据结构;
  • 执行过期判断;
  • 按连接返回响应;
  • 根据配置进行持久化或复制。

Redis 的命令通常具有单线程顺序执行的可见性特征:对于单个 Redis 实例,两个命令不会在同一时刻修改同一执行上下文。但是,这不意味着多个命令天然组成事务。例如:

GET stock:sku-1
# Java 计算 stock - 1
SET stock:sku-1 9

两个客户端可能都读到 10,最后都写入 9。单条命令的原子性不能自动扩展为多条命令的业务原子性。

1.2 Lettuce 是 Java 客户端

Lettuce 负责:

  • 建立 TCP 或 TLS 连接;
  • 编码 Java 参数为 Redis 协议;
  • 解码 Redis 响应;
  • 提供同步、异步和响应式接口;
  • 管理连接重连、拓扑刷新等客户端能力;
  • 将 Redis 协议错误、连接错误转换为 Java 异常或异步失败。

Lettuce 不会自动保证:

  • 你的对象序列化格式永远兼容;
  • 缓存一定不会击穿;
  • 分布式锁一定安全;
  • Streams 消费一定不重复;
  • Redis 故障时业务一定不丢数据。

这些是应用协议和故障模型问题。


2. 建立一个可运行的 Lettuce 工程

2.1 Maven 依赖

下面是一个普通 Maven 依赖示例。版本应在项目中统一管理,并在升级时查看对应 Lettuce 的兼容说明:

<dependency>
    <groupId>io.lettuce</groupId>
    <artifactId>lettuce-core</artifactId>
    <version>6.5.0.RELEASE</version>
</dependency>

Java 25 项目可以使用如下编译配置:

<properties>
    <maven.compiler.release>25</maven.compiler.release>
    <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
</properties>

启动本地 Redis:

docker run --rm --name redis-demo -p 6379:6379 redis:7

验证服务端:

redis-cli -h 127.0.0.1 -p 6379 PING

预期输出:

PONG

2.2 最小端到端程序

import io.lettuce.core.RedisClient;
import io.lettuce.core.RedisURI;
import io.lettuce.core.api.StatefulRedisConnection;
import io.lettuce.core.api.sync.RedisCommands;

import java.time.Duration;

public class LettuceDemo {
    public static void main(String[] args) {
        RedisURI uri = RedisURI.builder()
                .withHost("127.0.0.1")
                .withPort(6379)
                .withTimeout(Duration.ofSeconds(2))
                .build();

        RedisClient client = RedisClient.create(uri);

        try (StatefulRedisConnection<String, String> connection =
                     client.connect()) {
            RedisCommands<String, String> commands = connection.sync();

            commands.set("demo:name", "Alice");
            commands.expire("demo:name", 60);

            String value = commands.get("demo:name");
            System.out.println(value);       // Alice
            System.out.println(commands.ttl("demo:name")); // 约 60
        } finally {
            client.shutdown();
        }
    }
}

这里有三个不同生命周期:

  1. RedisClient:应用级客户端,通常创建一次;
  2. StatefulRedisConnection<K,V>:一个有状态的 Redis 连接;
  3. RedisCommands<K,V>:连接上的同步命令门面。

try 只负责关闭连接,finally 负责关闭客户端。Web 应用中通常不应每个请求都创建和销毁 RedisClient,否则会反复建立 TCP 连接、增加握手和线程管理成本。

2.3 连接 URI 与认证

生产环境通常至少配置地址、数据库、密码、TLS 和超时:

RedisURI uri = RedisURI.builder()
        .withHost("redis.example.internal")
        .withPort(6379)
        .withDatabase(0)
        .withAuthentication("app-user", "strong-password".toCharArray())
        .withSsl(true)
        .withTimeout(Duration.ofSeconds(2))
        .build();

密码不应硬编码在源代码中,应由环境变量、密钥管理系统或容器 Secret 注入。TLS 保护传输过程,但不会自动解决 Redis 命令权限;Redis ACL 仍需限制账号可执行的命令和键空间。


3. Lettuce 连接模型:同步、异步、响应式与并发

3.1 同步 API

同步调用会阻塞当前 Java 线程:

String value = connection.sync().get("key");

网络往返完成前,调用线程不能继续执行。同步 API 适合:

  • 业务代码本身是同步模型;
  • Redis 调用数量可控;
  • 已正确配置超时。

同步 API 不适合在事件循环线程中执行,因为一次网络阻塞会阻塞整个事件循环。

3.2 异步 API

var async = connection.async();

async.set("key", "value")
     .thenAccept(result -> System.out.println(result))
     .exceptionally(error -> {
         error.printStackTrace();
         return null;
     });

Lettuce 异步 API 返回 RedisFuture<T>,它与 CompletionStage 兼容。错误不会在调用处同步抛出,而通常会通过 Future 完成异常传播:

async.get("missing-key")
     .whenComplete((value, error) -> {
         if (error != null) {
             // 连接错误、Redis 错误等
             error.printStackTrace();
         } else {
             System.out.println(value);
         }
     });

如果业务需要严格超时,不能只依赖网络连接超时,还应对 Future 设置业务级超时:

async.get("key")
     .orTimeout(500, java.util.concurrent.TimeUnit.MILLISECONDS);

业务级超时只会让调用方停止等待,不一定能撤销服务端已经执行的命令。因此超时后的重试必须考虑命令是否已经成功执行。

3.3 响应式 API

Lettuce 提供 Reactor 适配的响应式命令接口。典型形式如下:

var reactive = connection.reactive();

reactive.set("key", "value")
        .then(reactive.get("key"))
        .doOnNext(System.out::println)
        .subscribe();

响应式调用通常是惰性的:只有订阅后才开始执行。若应用使用 Spring WebFlux 或其他响应式栈,应避免在响应式链中调用同步 API,否则会把阻塞操作带入非阻塞线程模型。

3.4 连接是否可以并发使用

Lettuce 的 StatefulConnection 通常设计为线程安全,可被多个线程共享。它内部会把请求写入连接,并按照协议关联响应。

但是,线程安全不等于业务操作隔离。以下场景不应依赖一个共享连接完成“只属于当前业务流程”的状态操作:

  • MULTI / EXEC 事务;
  • WATCH 乐观锁;
  • SUBSCRIBE 或阻塞式 XREADGROUP
  • 需要独占连接的交互流程;
  • 依赖连接级状态的命令。

例如,事务命令在连接上具有状态:

MULTI
SET a 1
SET b 2
EXEC

如果多个业务线程共享同一连接,命令可能交错,事务边界就会被破坏。此时应为该事务获取独立连接,或者使用 Lettuce 的连接池方案。连接池不是所有场景都必需;普通无状态命令共享一个连接通常足够,但事务、订阅和阻塞消费需要单独考虑。

3.5 Java 25 虚拟线程的边界

Java 25 的虚拟线程降低了阻塞线程的创建成本,但不会消除:

  • Redis 网络往返延迟;
  • Redis 服务端处理时间;
  • 连接数量限制;
  • 连接上的命令顺序和业务竞态。

可以让虚拟线程执行同步 Lettuce 调用,但仍需限制并发量,避免瞬间产生过多 Redis 请求。虚拟线程解决的是 Java 线程等待成本,不是 Redis 容量规划问题。


4. 序列化:Redis 存的是字节,不是 Java 对象

4.1 编解码器决定 Redis 看到的内容

Redis 的值本质上是字节数组。Lettuce 通过 RedisCodec<K,V> 将 Java 的键和值编码为字节,并把响应解码回来。

默认的 StringRedisCodec.UTF8 适合字符串:

var connection = client.connect(
        io.lettuce.core.codec.StringCodec.UTF8
);

此时:

commands.set("user:42", "{\"id\":42,\"name\":\"Alice\"}");

Redis CLI 可以直接观察值:

redis-cli GET user:42

输出:

{"id":42,"name":"Alice"}

如果使用 Java 原生序列化、Kryo、JSON、Protobuf 或自定义二进制协议,Redis CLI 看到的内容可能不可读。可读性并不等于正确性,但可读格式通常更便于诊断。

4.2 字符串 JSON 的显式序列化

一种稳妥的应用边界是:Redis 层只接受 String,对象序列化在业务代码中明确完成。以 Jackson 为例:

record User(long id, String name) {}

ObjectMapper mapper = new ObjectMapper();

User user = new User(42, "Alice");
String json = mapper.writeValueAsString(user);

commands.set("user:42", json);

String stored = commands.get("user:42");
User restored = mapper.readValue(stored, User.class);

这里必须明确处理 JsonProcessingException,因为序列化失败会导致缓存写入失败,反序列化失败则可能造成读取路径异常。

更重要的是,JSON 格式需要版本策略。例如旧值:

{"id":42,"name":"Alice"}

新代码增加字段:

record User(long id, String name, String level) {}

如果新字段没有默认兼容处理,旧缓存反序列化可能失败。常见做法是:

  • 新字段允许缺省;
  • 在 JSON 中增加显式版本字段;
  • 变更结构时修改键前缀,例如 v2:user:42
  • 不能解析的缓存直接删除并回源,而不是无限重试。

4.3 不要默认使用不安全的多态反序列化

某些序列化框架允许数据携带 Java 类型信息,然后按类型实例化对象。若数据来源可被外部写入,宽泛的多态反序列化可能造成反序列化攻击。

缓存数据也不能因为“只是缓存”就跳过安全边界。应优先使用固定目标类型:

User user = mapper.readValue(json, User.class);

而不是接受任意类名并动态实例化。

4.4 Hash、JSON 字符串与结构化数据

Redis Hash 可以把对象字段拆开:

commands.hset("user:42", java.util.Map.of(
        "id", "42",
        "name", "Alice"
));

读取:

java.util.Map<String, String> fields =
        commands.hgetall("user:42");

Hash 的优点是可以只更新一个字段,缺点是字段类型通常被转换成字符串,结构演进和原子更新逻辑需要应用自己维护。

JSON 字符串的优点是对象边界清晰,读写一次完成;缺点是修改单个字段通常需要读、修改、写,或者依赖 RedisJSON 等额外模块。选择哪一种取决于访问模式,不是“对象一定 JSON、字段一定 Hash”的固定规则。

4.5 键名也是协议

键名应包含稳定的命名空间和版本信息:

app:v1:user:42
app:v1:product:sku-100
app:v1:lock:order:9001
app:v1:stream:order-events

冒号只是约定,不是 Redis 的层级目录。键名协议需要定义:

  • 业务域;
  • 数据类型;
  • 版本;
  • 标识符编码;
  • 是否允许跨服务共享。

若不同服务对同一个键使用不同序列化格式,即使键名相同,也只是制造了隐蔽的数据兼容故障。


5. Redis 数据结构与原子性基础

5.1 常用类型

Redis 类型 典型用途 关键边界
String 缓存对象、计数器、锁令牌 复杂更新需 Lua 或条件命令
Hash 对象字段、属性集合 字段类型和版本由应用维护
List 简单队列、双端队列 消费确认能力弱于 Streams
Set 去重、集合关系 无序
Sorted Set 排行榜、延迟任务索引 分数相同需定义排序规则
Stream 追加事件、消费组 需要 ACK、重试和积压治理

Redis 的原子性一般指一条命令在服务端执行期间不会被另一条命令插入。INCR 是原子递增:

INCR page:views

但“读取后计算再写入”不是原子的。可以用服务端命令、Lua 脚本或事务重新建立原子边界。

5.2 TTL 不是永久保证

SET key value EX 60

表示设置值并附带 60 秒过期时间。相比先 SETEXPIRE,单条 SET ... EX 避免了两条命令之间进程崩溃导致“有值但无 TTL”。

检查:

TTL key

结果含义通常是:

  • >= 0:剩余秒数;
  • -1:键存在但没有过期时间;
  • -2:键不存在。

过期删除可能包含惰性删除和主动扫描,因此 TTL 到期不应被理解为精确到某一毫秒的定时器。若业务要求严格的时间调度,Redis TTL 不能单独承担该职责。


6. 缓存:从读路径到故障路径

6.1 Cache-aside 模式

Cache-aside,也称旁路缓存,由应用显式管理缓存:

读请求
  │
  ├─ GET cache
  │    ├─ 命中:返回缓存
  │    └─ 未命中
  │          │
  │          ├─ 查询数据库
  │          ├─ SET cache EX ttl
  │          └─ 返回结果
  │
写请求
  ├─ 更新数据库
  └─ 删除缓存

读取伪代码:

String key = "app:v1:user:" + userId;
String cached = commands.get(key);

if (cached != null) {
    return mapper.readValue(cached, User.class);
}

User user = userRepository.findById(userId); // 回源数据库
if (user != null) {
    String json = mapper.writeValueAsString(user);
    commands.setex(key, 300, json);
}
return user;

这里的 setex 是常见同步 API 形式;也可以使用 set(key, value, SetArgs.Builder.ex(300)),后者更明确地表达命令参数。

6.2 为什么写数据库后删除缓存

假设采用“更新数据库后删除缓存”:

T1:更新数据库为 B
T1:删除缓存

后续读取会从数据库加载 B,再写回缓存。若改成“更新数据库后更新缓存”,多个并发写请求可能按不同顺序完成,最终缓存值不一定对应最后一次数据库提交。

但删除缓存也不是绝对完美。存在如下竞态:

T1:读取数据库,得到旧值 A
T2:更新数据库为 B
T2:删除缓存
T1:把旧值 A 写回缓存

最终缓存重新出现旧值 A。常见缓解方式包括:

  • 缓存写入时携带版本号,只接受更高版本;
  • 延迟双删;
  • 通过消息或 CDC 进行失效传播;
  • 对关键读路径增加短暂一致性校验。

这些方法不能消除所有分布式故障,只能改变旧值重新出现的概率和可检测性。若业务要求强一致,应直接读数据库或采用带版本条件的更新协议。

6.3 缓存穿透、击穿、雪崩

缓存穿透:请求查询一个数据库中不存在的键,缓存也不存,恶意或高频请求持续打到数据库。

缓解方法:

commands.setex("app:v1:user:not-found:" + id, 30, "1");

这叫负缓存。它必须有较短 TTL,否则新创建的数据可能长期被“未找到”遮蔽。布隆过滤器可以降低无效查询,但存在误判,不能把“布隆过滤器不存在”之外的所有逻辑写错。

缓存击穿:某个热点键过期,大量线程同时回源。可以使用互斥重建锁,但锁失败的线程应等待、返回旧值或降级,而不是全部继续查询数据库。

缓存雪崩:大量键在相近时间过期,造成集中回源。可以给 TTL 增加随机扰动:

long ttl = 300 + ThreadLocalRandom.current().nextLong(60);
commands.set(key, json, SetArgs.Builder.ex(ttl));

随机 TTL 只是分散过期时间,不能替代容量、限流和回源保护。

6.4 逻辑过期与物理过期

物理过期依赖 Redis TTL,过期后键不可读。逻辑过期则把过期时间写入值:

{
  "data": {"id": 42, "name": "Alice"},
  "expireAt": 1730000000000
}

读取到逻辑过期数据后,可以先返回旧值,再异步刷新。它牺牲了一段时间的新鲜度,换取热点请求不同时阻塞回源。

两者可以组合:Redis 设置较长物理 TTL,值内部设置较短逻辑 TTL。物理 TTL 防止刷新任务永久失效,逻辑 TTL 控制业务新鲜度。

6.5 缓存错误如何处理

缓存读取失败时通常有两种策略:

  • 旁路降级:记录错误,直接回源数据库;
  • 快速失败:对核心链路不允许继续施压数据库,返回降级结果或错误。

不能无条件地把所有 Redis 异常转成缓存未命中。在 Redis 大面积故障时,所有请求回源可能造成数据库雪崩。应通过超时、并发舱壁、限流和降级控制回源规模。


7. 分布式锁:互斥、令牌与安全释放

7.1 一个锁至少需要什么性质

设锁键为 lock:order:42。一个基本分布式锁需要:

  • 互斥:同一时刻最多一个持有者;
  • 可释放:持有者完成后能删除;
  • 租约:持有者崩溃后锁最终能自动消失;
  • 所有权校验:旧持有者不能删除新持有者的锁。

只执行:

SET lock:order:42 owner EX 30

并不能保证互斥,因为没有 NX。正确的获取条件是:

SET lock:order:42 random-token NX EX 30

Java 示例:

String token = UUID.randomUUID().toString();
String result = commands.set(
        "app:v1:lock:order:42",
        token,
        SetArgs.Builder.nx().ex(30)
);

boolean acquired = "OK".equals(result);

NX 表示仅当键不存在时设置,EX 30 设置租约。两者放在同一个 SET 命令中,避免“先判断再设置”的竞态。

7.2 为什么释放不能直接 DEL

以下流程会错误释放别人的锁:

T1 获取锁,租约 30 秒
T1 执行任务超过 30 秒
锁自动过期
T2 获取同一把锁
T1 完成并执行 DEL
T2 的锁被 T1 删除

因此释放必须比较令牌,只有值仍等于自己的令牌时才删除。比较和删除必须是一个原子操作,Lua 脚本可表达:

if redis.call('get', KEYS[1]) == ARGV[1] then
    return redis.call('del', KEYS[1])
else
    return 0
end

Lettuce 调用:

String unlockScript = """
        if redis.call('get', KEYS[1]) == ARGV[1] then
            return redis.call('del', KEYS[1])
        else
            return 0
        end
        """;

Long deleted = commands.eval(
        unlockScript,
        ScriptOutputType.INTEGER,
        new String[]{"app:v1:lock:order:42"},
        token
);

if (deleted == 1L) {
    System.out.println("lock released");
}

释放结果为 1 表示删除了自己的锁,0 表示锁已不存在或已经属于其他持有者。

7.3 租约时间不是任务时间

如果任务最长执行时间为 TT,锁租约为 LL,还存在网络延迟、调度暂停和 GC 暂停 DD,要避免锁在任务完成前过期,至少需要满足:

L>T+D+网络与调度余量L > T + D + \text{网络与调度余量}

TT 往往不是严格上界,因此实际工程中还需要:

  • 自动续租;
  • 任务分段并校验所有权;
  • 让业务提交带 fencing token;
  • 接受锁只是“减少并发”,不是强一致事务。

Fencing token 是递增的所有权序号。每次成功获取锁获得更大的 token,写入下游资源时把 token 一起提交,下游拒绝比当前 token 更旧的写入。这样即使旧持有者因暂停后恢复,也不能覆盖新持有者的结果。

7.4 Redis 锁的边界

单 Redis 实例上的 SET NX EX 可以提供单实例租约锁语义,但在主从异步复制、故障切换和网络分区下,锁的全局安全性取决于部署模型和业务要求。

Redis 锁不等价于数据库事务,也不自动保护数据库写入。如果锁用于支付、库存等关键状态,通常还需要数据库唯一约束、条件更新、幂等键或 fencing token。将锁作为唯一正确性保障,是常见误解。

Redlock 等多节点算法有自己的假设、时钟和故障模型争议。使用前必须明确要解决的是:

  • 防止重复执行的性能优化;
  • 还是严格的线性一致性所有权。

两者的可靠性要求不同,不能只因为“部署了多个 Redis 节点”就认为问题已经解决。


8. Redis Streams:追加日志与消费组

8.1 Stream 的基本结构

Redis Stream 是按消息 ID 排序的追加型日志。消息形如:

ID                  field       value
1710000000000-0     orderId     9001
                    amount      100

生产消息:

String id = commands.xadd(
        "app:v1:stream:orders",
        java.util.Map.of(
                "orderId", "9001",
                "amount", "100"
        )
);

System.out.println(id); // 例如 1710000000000-0

消息 ID 通常由毫秒时间和序号组成。应用不应把 ID 当作业务唯一幂等键的唯一来源;业务事件仍应携带自己的 eventId

8.2 普通读取

var messages = commands.xread(
        io.lettuce.core.XReadArgs.Builder.count(10),
        io.lettuce.core.StreamOffset.from("app:v1:stream:orders", "0-0")
);

for (var message : messages) {
    System.out.println(message.getId());
    System.out.println(message.getBody());
}

0-0 表示从较早位置开始读取。普通 XREAD 不提供消费组的待确认列表,适合简单读取或广播式场景,但不适合需要明确 ACK 和故障接管的工作队列。

8.3 消费组、消费者和 PEL

消费组模型包含:

  • Stream:消息日志;
  • Consumer Group:一组协作消费者;
  • Consumer:组内一个逻辑消费者名称;
  • PEL(Pending Entries List):已经投递但尚未确认的消息集合。

创建消费组:

XGROUP CREATE app:v1:stream:orders orders-group 0 MKSTREAM

如果组可能已经存在,应用应把“组已存在”作为可接受结果处理,而不是把启动失败当成系统故障。Lettuce 中可以捕获 Redis 的 BUSYGROUP 错误。

读取新消息:

var records = commands.xreadgroup(
        Consumer.from("orders-group", "consumer-1"),
        XReadArgs.Builder.count(10).block(Duration.ofSeconds(5)),
        StreamOffset.lastConsumed("app:v1:stream:orders")
);

lastConsumed 表示从消费组记录的位置继续读取。对于新建组,起始 ID 是 0 还是 $ 会影响是否读取历史消息:

  • 0:从已有历史开始;
  • $:只关注创建组之后的新消息。

处理成功后确认:

for (var record : records) {
    try {
        handle(record.getBody());
        commands.xack(
                "app:v1:stream:orders",
                "orders-group",
                record.getId()
        );
    } catch (Exception e) {
        // 不 ACK,让消息保留在 PEL 中
        // 记录错误并等待重试或转入死信
    }
}

ACK 不是删除消息。 XACK 只从消费组的待确认集合中移除消息,Stream 中的消息仍然存在,直到通过保留策略或 XDEL 等方式清理。

8.4 崩溃与消息重投

考虑以下时序:

C1 读取消息 M
C1 业务处理成功
C1 进程在 XACK 前崩溃
C2 后续接管 M
C2 再次处理 M

因此 Streams 消费通常是至少一次投递,不是自动恰好一次。消费者必须幂等。例如数据库处理可以使用唯一约束:

CREATE TABLE processed_event (
    event_id VARCHAR(128) PRIMARY KEY,
    processed_at TIMESTAMP NOT NULL
);

处理事件时:

  1. 尝试插入 event_id
  2. 若唯一键冲突,说明已经处理过,直接 ACK;
  3. 插入成功后执行业务更新;
  4. 业务成功后 ACK。

但“插入幂等表”和“业务更新”如果不在同一个数据库事务中,仍可能出现中间失败。因此幂等记录与业务变更最好处于同一事务内。

8.5 接管空闲 Pending 消息

消费者崩溃后,消息仍属于旧消费者。其他消费者需要根据空闲时间接管。Redis 版本支持的命令和 Lettuce API 形式可能随版本变化,常见机制包括:

  • XAUTOCLAIM:按最小空闲时间自动扫描并认领;
  • XCLAIM:指定消息 ID 认领;
  • XPENDING:查看待确认消息和空闲时间。

运维或消费逻辑应先检查:

XPENDING app:v1:stream:orders orders-group

关注:

  • Pending 总数;
  • 最早消息 ID;
  • 消费者分布;
  • 长时间未确认的消息。

重试流程应有次数限制。超过阈值后写入死信 Stream:

XADD app:v1:stream:orders-dlq * orderId 9001 reason timeout

然后 ACK 原消息,避免坏消息永久阻塞消费。死信不是删除错误,它是把错误从主处理路径隔离出来,之后仍需人工或自动修复。

8.6 消费组并发和顺序

同一个消费组内,一条消息只会投递给一个消费者,但不同消息可以被不同消费者并行处理。因此:

  • 组内不保证所有消息的全局业务处理顺序;
  • 同一订单的事件如果要求顺序,需要按订单分区、串行处理或使用版本号校验;
  • eventId 幂等解决重复,不解决乱序;
  • 业务版本号可以拒绝旧事件,例如只接受 eventVersion > currentVersion

阻塞读取连接应与普通命令连接分离。一个 XREADGROUP BLOCK 长时间等待,不应占用事务、健康检查或普通请求使用的连接。


9. 事务、Lua 与批量操作的选择

9.1 Redis 事务不是回滚事务

Redis MULTI/EXEC 会把命令排队后依次执行,但通常没有关系数据库那样的回滚机制。若某个命令执行时报错,之前成功执行的命令不会自动撤销。

适合使用 Redis 事务的场景:

  • 需要将一组命令作为连续执行单元;
  • 通过 WATCH 实现乐观并发控制;
  • 不要求异常时自动回滚。

需要条件判断、比较值和修改值时,Lua 脚本通常更直接,因为脚本在 Redis 服务端原子执行:

local current = redis.call('get', KEYS[1])
if current == ARGV[1] then
    redis.call('set', KEYS[1], ARGV[2], 'EX', ARGV[3])
    return 1
end
return 0

脚本应保持短小,避免阻塞 Redis 主执行线程。Lua 不能把耗时数据库查询放进去;Redis 脚本只能处理 Redis 内部数据。

9.2 Pipeline 不是事务

Pipeline 把多个命令批量发送,减少网络往返:

commands.set("k1", "v1");
commands.set("k2", "v2");
commands.get("k1");

具体 Lettuce API 可通过异步批量、自动 flush 控制等方式实现。Pipeline 的核心是传输优化,不保证命令之间的事务原子性。网络断开时,客户端可能不知道批量中的哪些命令已经在服务端执行,因此重试必须基于幂等性设计。


10. 超时、重连与错误诊断

10.1 三类时间限制

Redis 客户端至少要区分:

  1. 连接建立超时:TCP、TLS、认证建立多久;
  2. 命令响应超时:发送命令后等待响应多久;
  3. 业务截止时间:整个请求允许 Redis 占用多久。

业务截止时间通常最严格。即使底层允许 2 秒,HTTP 请求只剩 100 毫秒,也不应再发一个必然超时的 Redis 操作。

10.2 错误分类

典型错误包括:

  • RedisCommandExecutionException:服务端返回错误,例如类型错误、权限错误、脚本错误;
  • RedisConnectionException:连接不可用或断开;
  • 超时异常:响应未在指定时间内到达;
  • 解码异常:响应字节无法按指定 Codec 转换。

诊断时至少记录:

  • 命令类型,而不是完整敏感值;
  • 键名的脱敏版本;
  • Redis 节点;
  • 超时时间;
  • 重试次数;
  • 请求 trace ID;
  • PING、连接状态和 Redis INFO 指标。

不要在日志中记录密码、锁令牌、用户隐私数据和完整缓存 JSON。

10.3 重试的危险

以下命令天然更适合重试:

GET
SET
DEL
EXPIRE

但“适合”不表示无条件安全。对于:

INCR
XADD
LPUSH

请求超时后,命令可能已经成功执行,只是响应丢失;直接重试可能造成重复计数或重复消息。解决方案包括:

  • 使用业务幂等键;
  • 让事件 ID 可去重;
  • 使用条件写入;
  • 查询结果确认,而不是盲目重发;
  • 明确接受重复并在下游消除。

11. 内存、持久化与高可用边界

Redis 是内存优先系统。即使配置了 RDB 或 AOF,也不能简单等同于每次写入都已同步落盘。

11.1 RDB 与 AOF

  • RDB:周期性生成快照,恢复较快,但两次快照之间的数据可能丢失;
  • AOF:记录写命令,通常能减少数据丢失窗口,但文件更大、重写和恢复成本不同;
  • appendfsync alwayseverysecno 在性能和持久性之间有不同取舍。

缓存数据通常可以接受丢失;订单事件、任务状态则需要更严格地评估持久化和恢复策略。

11.2 复制不是强一致

Redis 主从复制通常是异步的。主节点确认写入后,从节点可能尚未收到数据。主节点故障切换时,最近写入可能丢失。

因此:

  • 缓存丢失通常表现为回源;
  • Streams 消息丢失可能表现为业务事件缺口;
  • 锁状态在故障切换中可能出现旧锁和新锁并存的风险;
  • 需要关键持久化语义时,数据库或专门消息系统可能更适合。

Redis Cluster 将键分布到不同槽位。涉及多个键的事务或 Lua 脚本通常要求相关键落在同一槽位,可以使用 hash tag:

order:{9001}:state
order:{9001}:lock

{9001} 中的内容用于计算槽位,使两个键更可能位于同一槽。键设计必须在业务建模阶段完成,不能等跨槽错误出现后再补救。


12. 一个完整的订单事件流程

把缓存、锁和 Streams 放在一起,可以形成如下流程:

sequenceDiagram
    participant API as Java API
    participant R as Redis
    participant DB as Database
    participant C as Stream Consumer

    API->>R: GET user:v1:42
    alt cache hit
        R-->>API: JSON
    else cache miss
        API->>R: SET lock:user:42 token NX EX
        alt lock acquired
            API->>DB: SELECT user
            DB-->>API: user
            API->>R: SET user:v1:42 JSON EX
            API->>R: Lua compare-token DEL
        else lock not acquired
            API->>R: GET user:v1:42 or wait/backoff
        end
    end

    API->>DB: UPDATE order
    API->>R: XADD order-events eventId...
    R-->>API: stream message ID

    C->>R: XREADGROUP
    R-->>C: pending message
    C->>DB: idempotent business transaction
    C->>R: XACK

关键路径有三个不同的一致性边界:

  1. 缓存读取是性能优化边界,缓存丢失后可以回源;
  2. 锁是并发协调边界,需要令牌、租约和安全释放;
  3. Stream 是事件投递边界,需要 ACK、幂等、重试和积压监控。

如果“更新数据库”和“发布 Stream 事件”必须严格同时成功,直接先后执行仍可能发生:

数据库提交成功
Java 进程崩溃
XADD 尚未执行

更可靠的做法是数据库 Outbox:

  1. 在同一个数据库事务中更新业务表并写入 Outbox 表;
  2. 独立发布器读取 Outbox;
  3. 成功 XADD 后标记 Outbox 已发布;
  4. 发布器重复运行也必须幂等。

这体现了一个重要边界:Redis Streams 可以作为事件传输目标,但不能自动提供数据库事务的跨系统原子提交。


13. 生产验证与诊断命令

13.1 检查键和 TTL

redis-cli TYPE app:v1:user:42
redis-cli TTL app:v1:user:42
redis-cli MEMORY USAGE app:v1:user:42

TYPE 与客户端预期不符,例如代码执行 HGETALL 却返回 String 类型,Redis 会返回 WRONGTYPE 错误。这通常意味着键名冲突、版本迁移遗漏或序列化模型改变。

13.2 检查 Streams 积压

redis-cli XLEN app:v1:stream:orders
redis-cli XINFO GROUPS app:v1:stream:orders
redis-cli XPENDING app:v1:stream:orders orders-group

重点观察:

  • Stream 长度是否持续增长;
  • lag 是否增加;
  • Pending 是否持续增加;
  • 是否有消费者长时间无心跳;
  • 是否存在单个消费者占据大量 Pending。

清理 Stream 需要结合保留策略。盲目执行:

redis-cli DEL app:v1:stream:orders

会删除整个 Stream 及消费组,属于破坏性操作。生产清理应先确认备份、消费进度、恢复方案,再使用长度或时间窗口保留策略。

13.3 监控指标

应用和 Redis 两侧都应监控:

  • 命令延迟分位数;
  • 连接数和连接失败;
  • 超时和重试次数;
  • 缓存命中率;
  • 热点键访问;
  • 内存使用、淘汰数量;
  • Stream 长度、Pending 数、最老 Pending 空闲时间;
  • Lua 脚本耗时;
  • 主从复制延迟和故障切换。

只监控 Redis CPU 而不监控缓存回源量,无法判断缓存故障是否已经传导到数据库。


14. 常见错误及其根因

错误一:每个请求创建一个 RedisClient

根因是把客户端对象误当作短生命周期连接。结果是连接建立、线程和资源开销增加,故障时还可能产生连接风暴。

改为应用启动时创建客户端,关闭时统一释放;连接是否共享则根据事务、订阅和阻塞命令决定。

错误二:把对象直接 toString() 存入 Redis

toString() 通常不是稳定序列化协议,字段变化后无法可靠反序列化,也可能包含敏感信息。应使用明确的 JSON、二进制协议或 Hash,并定义版本兼容规则。

错误三:锁释放直接 DEL

这会让过期后的旧持有者删除新持有者的锁。必须使用随机令牌和比较令牌的 Lua 脚本。

错误四:收到 Stream 消息就立即 ACK

如果先 ACK 再处理,进程在业务处理前崩溃,消息已经从 PEL 中移除,无法自动重试。通常应先完成幂等业务处理,再 ACK。

错误五:把 Stream 当作无限日志

Stream 会持续占用内存。必须定义保留策略、归档策略、死信处理和积压告警。ACK 不会自动删除历史消息。

错误六:超时后无条件重试写命令

响应丢失不等于服务端未执行。重试 XADDINCR 等命令可能造成重复结果。写入协议必须设计幂等性或可确认性。

错误七:用 Redis 锁代替数据库约束

锁可能因租约过期、进程暂停、故障切换而失效。最终业务正确性仍应由数据库条件更新、唯一约束、版本号和幂等记录等机制共同保证。


15. 选择 Redis、数据库和 Streams 的边界

Redis 适合低延迟访问、临时状态、计数、集合运算、短期缓存、协调信息和事件流转。关系数据库更适合需要持久约束、复杂查询、多表事务和审计的核心业务数据。

可以用以下因果关系判断:

  • 数据丢失后可以重建:适合缓存;
  • 数据丢失后必须审计恢复:需要持久化系统和明确备份;
  • 需要严格唯一性:使用数据库唯一约束或等价的强约束;
  • 需要至少一次异步处理:Streams 消费组可行,但必须幂等;
  • 需要跨数据库事务发布事件:使用 Outbox 等事务消息模式;
  • 需要全局强一致锁:先定义故障模型,再评估 Redis 锁是否满足,不应直接套用代码。

Lettuce 只负责把 Java 程序接入 Redis。真正可靠的 Redis 工程,依赖的是清晰的键和值协议、正确的连接生命周期、明确的原子边界、可验证的租约和所有权、可恢复的消息消费流程,以及对网络和故障的显式建模。


系列导航与关联阅读

官方资料

本文依据 Java、Spring 与相关项目官方文档重新梳理;正文、示例与生产清单由 WR BLOG 编写。

评论

0 条讨论
0/1000
还没有评论,来聊聊你的看法
WR Blog 加载中...
返回文章
JavaJava 25 LTSRedis

Java Redis 工程:Lettuce、连接、序列化、缓存、锁和 Streams

Java Redis 工程:Lettuce、连接、序列化、缓存、锁和 Streams封面

Java 基础体系 · 第 87/100 篇。示例统一以 Java 25 LTS 为语言和 JVM 基线;框架示例使用与其兼容的现代稳定版本。

Java Redis 工程:Lettuce、连接、序列化、缓存、锁和 Streams

Redis 工程问题通常不是“会不会调用 GETSET”,而是要同时处理六件事:

  1. Lettuce 客户端如何建立连接、复用连接并关闭资源;
  2. Redis 字节如何映射为 Java 类型;
  3. 缓存数据如何写入、过期、重建和失效;
  4. 多线程、多实例下的锁是否真的具有互斥和安全释放能力;
  5. Redis Streams 如何投递、确认、重试和恢复消息;
  6. 网络断开、超时、主从切换、进程崩溃时,业务状态如何变化。

本文以 Java 25 LTS 为运行范围,以 Lettuce 的同步、异步和响应式 API 为客户端基础。示例使用 Redis 7.x 常见命令语义;具体 Lettuce 版本应根据项目发布周期选择兼容版本,而不应把 Java 版本直接等同于 Redis 或 Lettuce 版本。


1. Redis 与 Lettuce 的职责边界

1.1 Redis 是服务端数据结构与协议

Redis 是一个通过 RESP 协议提供命令服务的内存数据存储。客户端发送命令:

SET user:42:name Alice EX 60

Redis 返回结果:

OK

服务端负责:

  • 解析命令;
  • 修改内存中的数据结构;
  • 执行过期判断;
  • 按连接返回响应;
  • 根据配置进行持久化或复制。

Redis 的命令通常具有单线程顺序执行的可见性特征:对于单个 Redis 实例,两个命令不会在同一时刻修改同一执行上下文。但是,这不意味着多个命令天然组成事务。例如:

GET stock:sku-1
# Java 计算 stock - 1
SET stock:sku-1 9

两个客户端可能都读到 10,最后都写入 9。单条命令的原子性不能自动扩展为多条命令的业务原子性。

1.2 Lettuce 是 Java 客户端

Lettuce 负责:

  • 建立 TCP 或 TLS 连接;
  • 编码 Java 参数为 Redis 协议;
  • 解码 Redis 响应;
  • 提供同步、异步和响应式接口;
  • 管理连接重连、拓扑刷新等客户端能力;
  • 将 Redis 协议错误、连接错误转换为 Java 异常或异步失败。

Lettuce 不会自动保证:

  • 你的对象序列化格式永远兼容;
  • 缓存一定不会击穿;
  • 分布式锁一定安全;
  • Streams 消费一定不重复;
  • Redis 故障时业务一定不丢数据。

这些是应用协议和故障模型问题。


2. 建立一个可运行的 Lettuce 工程

2.1 Maven 依赖

下面是一个普通 Maven 依赖示例。版本应在项目中统一管理,并在升级时查看对应 Lettuce 的兼容说明:

<dependency>
    <groupId>io.lettuce</groupId>
    <artifactId>lettuce-core</artifactId>
    <version>6.5.0.RELEASE</version>
</dependency>

Java 25 项目可以使用如下编译配置:

<properties>
    <maven.compiler.release>25</maven.compiler.release>
    <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
</properties>

启动本地 Redis:

docker run --rm --name redis-demo -p 6379:6379 redis:7

验证服务端:

redis-cli -h 127.0.0.1 -p 6379 PING

预期输出:

PONG

2.2 最小端到端程序

import io.lettuce.core.RedisClient;
import io.lettuce.core.RedisURI;
import io.lettuce.core.api.StatefulRedisConnection;
import io.lettuce.core.api.sync.RedisCommands;

import java.time.Duration;

public class LettuceDemo {
    public static void main(String[] args) {
        RedisURI uri = RedisURI.builder()
                .withHost("127.0.0.1")
                .withPort(6379)
                .withTimeout(Duration.ofSeconds(2))
                .build();

        RedisClient client = RedisClient.create(uri);

        try (StatefulRedisConnection<String, String> connection =
                     client.connect()) {
            RedisCommands<String, String> commands = connection.sync();

            commands.set("demo:name", "Alice");
            commands.expire("demo:name", 60);

            String value = commands.get("demo:name");
            System.out.println(value);       // Alice
            System.out.println(commands.ttl("demo:name")); // 约 60
        } finally {
            client.shutdown();
        }
    }
}

这里有三个不同生命周期:

  1. RedisClient:应用级客户端,通常创建一次;
  2. StatefulRedisConnection<K,V>:一个有状态的 Redis 连接;
  3. RedisCommands<K,V>:连接上的同步命令门面。

try 只负责关闭连接,finally 负责关闭客户端。Web 应用中通常不应每个请求都创建和销毁 RedisClient,否则会反复建立 TCP 连接、增加握手和线程管理成本。

2.3 连接 URI 与认证

生产环境通常至少配置地址、数据库、密码、TLS 和超时:

RedisURI uri = RedisURI.builder()
        .withHost("redis.example.internal")
        .withPort(6379)
        .withDatabase(0)
        .withAuthentication("app-user", "strong-password".toCharArray())
        .withSsl(true)
        .withTimeout(Duration.ofSeconds(2))
        .build();

密码不应硬编码在源代码中,应由环境变量、密钥管理系统或容器 Secret 注入。TLS 保护传输过程,但不会自动解决 Redis 命令权限;Redis ACL 仍需限制账号可执行的命令和键空间。


3. Lettuce 连接模型:同步、异步、响应式与并发

3.1 同步 API

同步调用会阻塞当前 Java 线程:

String value = connection.sync().get("key");

网络往返完成前,调用线程不能继续执行。同步 API 适合:

  • 业务代码本身是同步模型;
  • Redis 调用数量可控;
  • 已正确配置超时。

同步 API 不适合在事件循环线程中执行,因为一次网络阻塞会阻塞整个事件循环。

3.2 异步 API

var async = connection.async();

async.set("key", "value")
     .thenAccept(result -> System.out.println(result))
     .exceptionally(error -> {
         error.printStackTrace();
         return null;
     });

Lettuce 异步 API 返回 RedisFuture<T>,它与 CompletionStage 兼容。错误不会在调用处同步抛出,而通常会通过 Future 完成异常传播:

async.get("missing-key")
     .whenComplete((value, error) -> {
         if (error != null) {
             // 连接错误、Redis 错误等
             error.printStackTrace();
         } else {
             System.out.println(value);
         }
     });

如果业务需要严格超时,不能只依赖网络连接超时,还应对 Future 设置业务级超时:

async.get("key")
     .orTimeout(500, java.util.concurrent.TimeUnit.MILLISECONDS);

业务级超时只会让调用方停止等待,不一定能撤销服务端已经执行的命令。因此超时后的重试必须考虑命令是否已经成功执行。

3.3 响应式 API

Lettuce 提供 Reactor 适配的响应式命令接口。典型形式如下:

var reactive = connection.reactive();

reactive.set("key", "value")
        .then(reactive.get("key"))
        .doOnNext(System.out::println)
        .subscribe();

响应式调用通常是惰性的:只有订阅后才开始执行。若应用使用 Spring WebFlux 或其他响应式栈,应避免在响应式链中调用同步 API,否则会把阻塞操作带入非阻塞线程模型。

3.4 连接是否可以并发使用

Lettuce 的 StatefulConnection 通常设计为线程安全,可被多个线程共享。它内部会把请求写入连接,并按照协议关联响应。

但是,线程安全不等于业务操作隔离。以下场景不应依赖一个共享连接完成“只属于当前业务流程”的状态操作:

  • MULTI / EXEC 事务;
  • WATCH 乐观锁;
  • SUBSCRIBE 或阻塞式 XREADGROUP
  • 需要独占连接的交互流程;
  • 依赖连接级状态的命令。

例如,事务命令在连接上具有状态:

MULTI
SET a 1
SET b 2
EXEC

如果多个业务线程共享同一连接,命令可能交错,事务边界就会被破坏。此时应为该事务获取独立连接,或者使用 Lettuce 的连接池方案。连接池不是所有场景都必需;普通无状态命令共享一个连接通常足够,但事务、订阅和阻塞消费需要单独考虑。

3.5 Java 25 虚拟线程的边界

Java 25 的虚拟线程降低了阻塞线程的创建成本,但不会消除:

  • Redis 网络往返延迟;
  • Redis 服务端处理时间;
  • 连接数量限制;
  • 连接上的命令顺序和业务竞态。

可以让虚拟线程执行同步 Lettuce 调用,但仍需限制并发量,避免瞬间产生过多 Redis 请求。虚拟线程解决的是 Java 线程等待成本,不是 Redis 容量规划问题。


4. 序列化:Redis 存的是字节,不是 Java 对象

4.1 编解码器决定 Redis 看到的内容

Redis 的值本质上是字节数组。Lettuce 通过 RedisCodec<K,V> 将 Java 的键和值编码为字节,并把响应解码回来。

默认的 StringRedisCodec.UTF8 适合字符串:

var connection = client.connect(
        io.lettuce.core.codec.StringCodec.UTF8
);

此时:

commands.set("user:42", "{\"id\":42,\"name\":\"Alice\"}");

Redis CLI 可以直接观察值:

redis-cli GET user:42

输出:

{"id":42,"name":"Alice"}

如果使用 Java 原生序列化、Kryo、JSON、Protobuf 或自定义二进制协议,Redis CLI 看到的内容可能不可读。可读性并不等于正确性,但可读格式通常更便于诊断。

4.2 字符串 JSON 的显式序列化

一种稳妥的应用边界是:Redis 层只接受 String,对象序列化在业务代码中明确完成。以 Jackson 为例:

record User(long id, String name) {}

ObjectMapper mapper = new ObjectMapper();

User user = new User(42, "Alice");
String json = mapper.writeValueAsString(user);

commands.set("user:42", json);

String stored = commands.get("user:42");
User restored = mapper.readValue(stored, User.class);

这里必须明确处理 JsonProcessingException,因为序列化失败会导致缓存写入失败,反序列化失败则可能造成读取路径异常。

更重要的是,JSON 格式需要版本策略。例如旧值:

{"id":42,"name":"Alice"}

新代码增加字段:

record User(long id, String name, String level) {}

如果新字段没有默认兼容处理,旧缓存反序列化可能失败。常见做法是:

  • 新字段允许缺省;
  • 在 JSON 中增加显式版本字段;
  • 变更结构时修改键前缀,例如 v2:user:42
  • 不能解析的缓存直接删除并回源,而不是无限重试。

4.3 不要默认使用不安全的多态反序列化

某些序列化框架允许数据携带 Java 类型信息,然后按类型实例化对象。若数据来源可被外部写入,宽泛的多态反序列化可能造成反序列化攻击。

缓存数据也不能因为“只是缓存”就跳过安全边界。应优先使用固定目标类型:

User user = mapper.readValue(json, User.class);

而不是接受任意类名并动态实例化。

4.4 Hash、JSON 字符串与结构化数据

Redis Hash 可以把对象字段拆开:

commands.hset("user:42", java.util.Map.of(
        "id", "42",
        "name", "Alice"
));

读取:

java.util.Map<String, String> fields =
        commands.hgetall("user:42");

Hash 的优点是可以只更新一个字段,缺点是字段类型通常被转换成字符串,结构演进和原子更新逻辑需要应用自己维护。

JSON 字符串的优点是对象边界清晰,读写一次完成;缺点是修改单个字段通常需要读、修改、写,或者依赖 RedisJSON 等额外模块。选择哪一种取决于访问模式,不是“对象一定 JSON、字段一定 Hash”的固定规则。

4.5 键名也是协议

键名应包含稳定的命名空间和版本信息:

app:v1:user:42
app:v1:product:sku-100
app:v1:lock:order:9001
app:v1:stream:order-events

冒号只是约定,不是 Redis 的层级目录。键名协议需要定义:

  • 业务域;
  • 数据类型;
  • 版本;
  • 标识符编码;
  • 是否允许跨服务共享。

若不同服务对同一个键使用不同序列化格式,即使键名相同,也只是制造了隐蔽的数据兼容故障。


5. Redis 数据结构与原子性基础

5.1 常用类型

Redis 类型 典型用途 关键边界
String 缓存对象、计数器、锁令牌 复杂更新需 Lua 或条件命令
Hash 对象字段、属性集合 字段类型和版本由应用维护
List 简单队列、双端队列 消费确认能力弱于 Streams
Set 去重、集合关系 无序
Sorted Set 排行榜、延迟任务索引 分数相同需定义排序规则
Stream 追加事件、消费组 需要 ACK、重试和积压治理

Redis 的原子性一般指一条命令在服务端执行期间不会被另一条命令插入。INCR 是原子递增:

INCR page:views

但“读取后计算再写入”不是原子的。可以用服务端命令、Lua 脚本或事务重新建立原子边界。

5.2 TTL 不是永久保证

SET key value EX 60

表示设置值并附带 60 秒过期时间。相比先 SETEXPIRE,单条 SET ... EX 避免了两条命令之间进程崩溃导致“有值但无 TTL”。

检查:

TTL key

结果含义通常是:

  • >= 0:剩余秒数;
  • -1:键存在但没有过期时间;
  • -2:键不存在。

过期删除可能包含惰性删除和主动扫描,因此 TTL 到期不应被理解为精确到某一毫秒的定时器。若业务要求严格的时间调度,Redis TTL 不能单独承担该职责。


6. 缓存:从读路径到故障路径

6.1 Cache-aside 模式

Cache-aside,也称旁路缓存,由应用显式管理缓存:

读请求
  │
  ├─ GET cache
  │    ├─ 命中:返回缓存
  │    └─ 未命中
  │          │
  │          ├─ 查询数据库
  │          ├─ SET cache EX ttl
  │          └─ 返回结果
  │
写请求
  ├─ 更新数据库
  └─ 删除缓存

读取伪代码:

String key = "app:v1:user:" + userId;
String cached = commands.get(key);

if (cached != null) {
    return mapper.readValue(cached, User.class);
}

User user = userRepository.findById(userId); // 回源数据库
if (user != null) {
    String json = mapper.writeValueAsString(user);
    commands.setex(key, 300, json);
}
return user;

这里的 setex 是常见同步 API 形式;也可以使用 set(key, value, SetArgs.Builder.ex(300)),后者更明确地表达命令参数。

6.2 为什么写数据库后删除缓存

假设采用“更新数据库后删除缓存”:

T1:更新数据库为 B
T1:删除缓存

后续读取会从数据库加载 B,再写回缓存。若改成“更新数据库后更新缓存”,多个并发写请求可能按不同顺序完成,最终缓存值不一定对应最后一次数据库提交。

但删除缓存也不是绝对完美。存在如下竞态:

T1:读取数据库,得到旧值 A
T2:更新数据库为 B
T2:删除缓存
T1:把旧值 A 写回缓存

最终缓存重新出现旧值 A。常见缓解方式包括:

  • 缓存写入时携带版本号,只接受更高版本;
  • 延迟双删;
  • 通过消息或 CDC 进行失效传播;
  • 对关键读路径增加短暂一致性校验。

这些方法不能消除所有分布式故障,只能改变旧值重新出现的概率和可检测性。若业务要求强一致,应直接读数据库或采用带版本条件的更新协议。

6.3 缓存穿透、击穿、雪崩

缓存穿透:请求查询一个数据库中不存在的键,缓存也不存,恶意或高频请求持续打到数据库。

缓解方法:

commands.setex("app:v1:user:not-found:" + id, 30, "1");

这叫负缓存。它必须有较短 TTL,否则新创建的数据可能长期被“未找到”遮蔽。布隆过滤器可以降低无效查询,但存在误判,不能把“布隆过滤器不存在”之外的所有逻辑写错。

缓存击穿:某个热点键过期,大量线程同时回源。可以使用互斥重建锁,但锁失败的线程应等待、返回旧值或降级,而不是全部继续查询数据库。

缓存雪崩:大量键在相近时间过期,造成集中回源。可以给 TTL 增加随机扰动:

long ttl = 300 + ThreadLocalRandom.current().nextLong(60);
commands.set(key, json, SetArgs.Builder.ex(ttl));

随机 TTL 只是分散过期时间,不能替代容量、限流和回源保护。

6.4 逻辑过期与物理过期

物理过期依赖 Redis TTL,过期后键不可读。逻辑过期则把过期时间写入值:

{
  "data": {"id": 42, "name": "Alice"},
  "expireAt": 1730000000000
}

读取到逻辑过期数据后,可以先返回旧值,再异步刷新。它牺牲了一段时间的新鲜度,换取热点请求不同时阻塞回源。

两者可以组合:Redis 设置较长物理 TTL,值内部设置较短逻辑 TTL。物理 TTL 防止刷新任务永久失效,逻辑 TTL 控制业务新鲜度。

6.5 缓存错误如何处理

缓存读取失败时通常有两种策略:

  • 旁路降级:记录错误,直接回源数据库;
  • 快速失败:对核心链路不允许继续施压数据库,返回降级结果或错误。

不能无条件地把所有 Redis 异常转成缓存未命中。在 Redis 大面积故障时,所有请求回源可能造成数据库雪崩。应通过超时、并发舱壁、限流和降级控制回源规模。


7. 分布式锁:互斥、令牌与安全释放

7.1 一个锁至少需要什么性质

设锁键为 lock:order:42。一个基本分布式锁需要:

  • 互斥:同一时刻最多一个持有者;
  • 可释放:持有者完成后能删除;
  • 租约:持有者崩溃后锁最终能自动消失;
  • 所有权校验:旧持有者不能删除新持有者的锁。

只执行:

SET lock:order:42 owner EX 30

并不能保证互斥,因为没有 NX。正确的获取条件是:

SET lock:order:42 random-token NX EX 30

Java 示例:

String token = UUID.randomUUID().toString();
String result = commands.set(
        "app:v1:lock:order:42",
        token,
        SetArgs.Builder.nx().ex(30)
);

boolean acquired = "OK".equals(result);

NX 表示仅当键不存在时设置,EX 30 设置租约。两者放在同一个 SET 命令中,避免“先判断再设置”的竞态。

7.2 为什么释放不能直接 DEL

以下流程会错误释放别人的锁:

T1 获取锁,租约 30 秒
T1 执行任务超过 30 秒
锁自动过期
T2 获取同一把锁
T1 完成并执行 DEL
T2 的锁被 T1 删除

因此释放必须比较令牌,只有值仍等于自己的令牌时才删除。比较和删除必须是一个原子操作,Lua 脚本可表达:

if redis.call('get', KEYS[1]) == ARGV[1] then
    return redis.call('del', KEYS[1])
else
    return 0
end

Lettuce 调用:

String unlockScript = """
        if redis.call('get', KEYS[1]) == ARGV[1] then
            return redis.call('del', KEYS[1])
        else
            return 0
        end
        """;

Long deleted = commands.eval(
        unlockScript,
        ScriptOutputType.INTEGER,
        new String[]{"app:v1:lock:order:42"},
        token
);

if (deleted == 1L) {
    System.out.println("lock released");
}

释放结果为 1 表示删除了自己的锁,0 表示锁已不存在或已经属于其他持有者。

7.3 租约时间不是任务时间

如果任务最长执行时间为 TT,锁租约为 LL,还存在网络延迟、调度暂停和 GC 暂停 DD,要避免锁在任务完成前过期,至少需要满足:

L>T+D+网络与调度余量L > T + D + \text{网络与调度余量}

TT 往往不是严格上界,因此实际工程中还需要:

  • 自动续租;
  • 任务分段并校验所有权;
  • 让业务提交带 fencing token;
  • 接受锁只是“减少并发”,不是强一致事务。

Fencing token 是递增的所有权序号。每次成功获取锁获得更大的 token,写入下游资源时把 token 一起提交,下游拒绝比当前 token 更旧的写入。这样即使旧持有者因暂停后恢复,也不能覆盖新持有者的结果。

7.4 Redis 锁的边界

单 Redis 实例上的 SET NX EX 可以提供单实例租约锁语义,但在主从异步复制、故障切换和网络分区下,锁的全局安全性取决于部署模型和业务要求。

Redis 锁不等价于数据库事务,也不自动保护数据库写入。如果锁用于支付、库存等关键状态,通常还需要数据库唯一约束、条件更新、幂等键或 fencing token。将锁作为唯一正确性保障,是常见误解。

Redlock 等多节点算法有自己的假设、时钟和故障模型争议。使用前必须明确要解决的是:

  • 防止重复执行的性能优化;
  • 还是严格的线性一致性所有权。

两者的可靠性要求不同,不能只因为“部署了多个 Redis 节点”就认为问题已经解决。


8. Redis Streams:追加日志与消费组

8.1 Stream 的基本结构

Redis Stream 是按消息 ID 排序的追加型日志。消息形如:

ID                  field       value
1710000000000-0     orderId     9001
                    amount      100

生产消息:

String id = commands.xadd(
        "app:v1:stream:orders",
        java.util.Map.of(
                "orderId", "9001",
                "amount", "100"
        )
);

System.out.println(id); // 例如 1710000000000-0

消息 ID 通常由毫秒时间和序号组成。应用不应把 ID 当作业务唯一幂等键的唯一来源;业务事件仍应携带自己的 eventId

8.2 普通读取

var messages = commands.xread(
        io.lettuce.core.XReadArgs.Builder.count(10),
        io.lettuce.core.StreamOffset.from("app:v1:stream:orders", "0-0")
);

for (var message : messages) {
    System.out.println(message.getId());
    System.out.println(message.getBody());
}

0-0 表示从较早位置开始读取。普通 XREAD 不提供消费组的待确认列表,适合简单读取或广播式场景,但不适合需要明确 ACK 和故障接管的工作队列。

8.3 消费组、消费者和 PEL

消费组模型包含:

  • Stream:消息日志;
  • Consumer Group:一组协作消费者;
  • Consumer:组内一个逻辑消费者名称;
  • PEL(Pending Entries List):已经投递但尚未确认的消息集合。

创建消费组:

XGROUP CREATE app:v1:stream:orders orders-group 0 MKSTREAM

如果组可能已经存在,应用应把“组已存在”作为可接受结果处理,而不是把启动失败当成系统故障。Lettuce 中可以捕获 Redis 的 BUSYGROUP 错误。

读取新消息:

var records = commands.xreadgroup(
        Consumer.from("orders-group", "consumer-1"),
        XReadArgs.Builder.count(10).block(Duration.ofSeconds(5)),
        StreamOffset.lastConsumed("app:v1:stream:orders")
);

lastConsumed 表示从消费组记录的位置继续读取。对于新建组,起始 ID 是 0 还是 $ 会影响是否读取历史消息:

  • 0:从已有历史开始;
  • $:只关注创建组之后的新消息。

处理成功后确认:

for (var record : records) {
    try {
        handle(record.getBody());
        commands.xack(
                "app:v1:stream:orders",
                "orders-group",
                record.getId()
        );
    } catch (Exception e) {
        // 不 ACK,让消息保留在 PEL 中
        // 记录错误并等待重试或转入死信
    }
}

ACK 不是删除消息。 XACK 只从消费组的待确认集合中移除消息,Stream 中的消息仍然存在,直到通过保留策略或 XDEL 等方式清理。

8.4 崩溃与消息重投

考虑以下时序:

C1 读取消息 M
C1 业务处理成功
C1 进程在 XACK 前崩溃
C2 后续接管 M
C2 再次处理 M

因此 Streams 消费通常是至少一次投递,不是自动恰好一次。消费者必须幂等。例如数据库处理可以使用唯一约束:

CREATE TABLE processed_event (
    event_id VARCHAR(128) PRIMARY KEY,
    processed_at TIMESTAMP NOT NULL
);

处理事件时:

  1. 尝试插入 event_id
  2. 若唯一键冲突,说明已经处理过,直接 ACK;
  3. 插入成功后执行业务更新;
  4. 业务成功后 ACK。

但“插入幂等表”和“业务更新”如果不在同一个数据库事务中,仍可能出现中间失败。因此幂等记录与业务变更最好处于同一事务内。

8.5 接管空闲 Pending 消息

消费者崩溃后,消息仍属于旧消费者。其他消费者需要根据空闲时间接管。Redis 版本支持的命令和 Lettuce API 形式可能随版本变化,常见机制包括:

  • XAUTOCLAIM:按最小空闲时间自动扫描并认领;
  • XCLAIM:指定消息 ID 认领;
  • XPENDING:查看待确认消息和空闲时间。

运维或消费逻辑应先检查:

XPENDING app:v1:stream:orders orders-group

关注:

  • Pending 总数;
  • 最早消息 ID;
  • 消费者分布;
  • 长时间未确认的消息。

重试流程应有次数限制。超过阈值后写入死信 Stream:

XADD app:v1:stream:orders-dlq * orderId 9001 reason timeout

然后 ACK 原消息,避免坏消息永久阻塞消费。死信不是删除错误,它是把错误从主处理路径隔离出来,之后仍需人工或自动修复。

8.6 消费组并发和顺序

同一个消费组内,一条消息只会投递给一个消费者,但不同消息可以被不同消费者并行处理。因此:

  • 组内不保证所有消息的全局业务处理顺序;
  • 同一订单的事件如果要求顺序,需要按订单分区、串行处理或使用版本号校验;
  • eventId 幂等解决重复,不解决乱序;
  • 业务版本号可以拒绝旧事件,例如只接受 eventVersion > currentVersion

阻塞读取连接应与普通命令连接分离。一个 XREADGROUP BLOCK 长时间等待,不应占用事务、健康检查或普通请求使用的连接。


9. 事务、Lua 与批量操作的选择

9.1 Redis 事务不是回滚事务

Redis MULTI/EXEC 会把命令排队后依次执行,但通常没有关系数据库那样的回滚机制。若某个命令执行时报错,之前成功执行的命令不会自动撤销。

适合使用 Redis 事务的场景:

  • 需要将一组命令作为连续执行单元;
  • 通过 WATCH 实现乐观并发控制;
  • 不要求异常时自动回滚。

需要条件判断、比较值和修改值时,Lua 脚本通常更直接,因为脚本在 Redis 服务端原子执行:

local current = redis.call('get', KEYS[1])
if current == ARGV[1] then
    redis.call('set', KEYS[1], ARGV[2], 'EX', ARGV[3])
    return 1
end
return 0

脚本应保持短小,避免阻塞 Redis 主执行线程。Lua 不能把耗时数据库查询放进去;Redis 脚本只能处理 Redis 内部数据。

9.2 Pipeline 不是事务

Pipeline 把多个命令批量发送,减少网络往返:

commands.set("k1", "v1");
commands.set("k2", "v2");
commands.get("k1");

具体 Lettuce API 可通过异步批量、自动 flush 控制等方式实现。Pipeline 的核心是传输优化,不保证命令之间的事务原子性。网络断开时,客户端可能不知道批量中的哪些命令已经在服务端执行,因此重试必须基于幂等性设计。


10. 超时、重连与错误诊断

10.1 三类时间限制

Redis 客户端至少要区分:

  1. 连接建立超时:TCP、TLS、认证建立多久;
  2. 命令响应超时:发送命令后等待响应多久;
  3. 业务截止时间:整个请求允许 Redis 占用多久。

业务截止时间通常最严格。即使底层允许 2 秒,HTTP 请求只剩 100 毫秒,也不应再发一个必然超时的 Redis 操作。

10.2 错误分类

典型错误包括:

  • RedisCommandExecutionException:服务端返回错误,例如类型错误、权限错误、脚本错误;
  • RedisConnectionException:连接不可用或断开;
  • 超时异常:响应未在指定时间内到达;
  • 解码异常:响应字节无法按指定 Codec 转换。

诊断时至少记录:

  • 命令类型,而不是完整敏感值;
  • 键名的脱敏版本;
  • Redis 节点;
  • 超时时间;
  • 重试次数;
  • 请求 trace ID;
  • PING、连接状态和 Redis INFO 指标。

不要在日志中记录密码、锁令牌、用户隐私数据和完整缓存 JSON。

10.3 重试的危险

以下命令天然更适合重试:

GET
SET
DEL
EXPIRE

但“适合”不表示无条件安全。对于:

INCR
XADD
LPUSH

请求超时后,命令可能已经成功执行,只是响应丢失;直接重试可能造成重复计数或重复消息。解决方案包括:

  • 使用业务幂等键;
  • 让事件 ID 可去重;
  • 使用条件写入;
  • 查询结果确认,而不是盲目重发;
  • 明确接受重复并在下游消除。

11. 内存、持久化与高可用边界

Redis 是内存优先系统。即使配置了 RDB 或 AOF,也不能简单等同于每次写入都已同步落盘。

11.1 RDB 与 AOF

  • RDB:周期性生成快照,恢复较快,但两次快照之间的数据可能丢失;
  • AOF:记录写命令,通常能减少数据丢失窗口,但文件更大、重写和恢复成本不同;
  • appendfsync alwayseverysecno 在性能和持久性之间有不同取舍。

缓存数据通常可以接受丢失;订单事件、任务状态则需要更严格地评估持久化和恢复策略。

11.2 复制不是强一致

Redis 主从复制通常是异步的。主节点确认写入后,从节点可能尚未收到数据。主节点故障切换时,最近写入可能丢失。

因此:

  • 缓存丢失通常表现为回源;
  • Streams 消息丢失可能表现为业务事件缺口;
  • 锁状态在故障切换中可能出现旧锁和新锁并存的风险;
  • 需要关键持久化语义时,数据库或专门消息系统可能更适合。

Redis Cluster 将键分布到不同槽位。涉及多个键的事务或 Lua 脚本通常要求相关键落在同一槽位,可以使用 hash tag:

order:{9001}:state
order:{9001}:lock

{9001} 中的内容用于计算槽位,使两个键更可能位于同一槽。键设计必须在业务建模阶段完成,不能等跨槽错误出现后再补救。


12. 一个完整的订单事件流程

把缓存、锁和 Streams 放在一起,可以形成如下流程:

sequenceDiagram
    participant API as Java API
    participant R as Redis
    participant DB as Database
    participant C as Stream Consumer

    API->>R: GET user:v1:42
    alt cache hit
        R-->>API: JSON
    else cache miss
        API->>R: SET lock:user:42 token NX EX
        alt lock acquired
            API->>DB: SELECT user
            DB-->>API: user
            API->>R: SET user:v1:42 JSON EX
            API->>R: Lua compare-token DEL
        else lock not acquired
            API->>R: GET user:v1:42 or wait/backoff
        end
    end

    API->>DB: UPDATE order
    API->>R: XADD order-events eventId...
    R-->>API: stream message ID

    C->>R: XREADGROUP
    R-->>C: pending message
    C->>DB: idempotent business transaction
    C->>R: XACK

关键路径有三个不同的一致性边界:

  1. 缓存读取是性能优化边界,缓存丢失后可以回源;
  2. 锁是并发协调边界,需要令牌、租约和安全释放;
  3. Stream 是事件投递边界,需要 ACK、幂等、重试和积压监控。

如果“更新数据库”和“发布 Stream 事件”必须严格同时成功,直接先后执行仍可能发生:

数据库提交成功
Java 进程崩溃
XADD 尚未执行

更可靠的做法是数据库 Outbox:

  1. 在同一个数据库事务中更新业务表并写入 Outbox 表;
  2. 独立发布器读取 Outbox;
  3. 成功 XADD 后标记 Outbox 已发布;
  4. 发布器重复运行也必须幂等。

这体现了一个重要边界:Redis Streams 可以作为事件传输目标,但不能自动提供数据库事务的跨系统原子提交。


13. 生产验证与诊断命令

13.1 检查键和 TTL

redis-cli TYPE app:v1:user:42
redis-cli TTL app:v1:user:42
redis-cli MEMORY USAGE app:v1:user:42

TYPE 与客户端预期不符,例如代码执行 HGETALL 却返回 String 类型,Redis 会返回 WRONGTYPE 错误。这通常意味着键名冲突、版本迁移遗漏或序列化模型改变。

13.2 检查 Streams 积压

redis-cli XLEN app:v1:stream:orders
redis-cli XINFO GROUPS app:v1:stream:orders
redis-cli XPENDING app:v1:stream:orders orders-group

重点观察:

  • Stream 长度是否持续增长;
  • lag 是否增加;
  • Pending 是否持续增加;
  • 是否有消费者长时间无心跳;
  • 是否存在单个消费者占据大量 Pending。

清理 Stream 需要结合保留策略。盲目执行:

redis-cli DEL app:v1:stream:orders

会删除整个 Stream 及消费组,属于破坏性操作。生产清理应先确认备份、消费进度、恢复方案,再使用长度或时间窗口保留策略。

13.3 监控指标

应用和 Redis 两侧都应监控:

  • 命令延迟分位数;
  • 连接数和连接失败;
  • 超时和重试次数;
  • 缓存命中率;
  • 热点键访问;
  • 内存使用、淘汰数量;
  • Stream 长度、Pending 数、最老 Pending 空闲时间;
  • Lua 脚本耗时;
  • 主从复制延迟和故障切换。

只监控 Redis CPU 而不监控缓存回源量,无法判断缓存故障是否已经传导到数据库。


14. 常见错误及其根因

错误一:每个请求创建一个 RedisClient

根因是把客户端对象误当作短生命周期连接。结果是连接建立、线程和资源开销增加,故障时还可能产生连接风暴。

改为应用启动时创建客户端,关闭时统一释放;连接是否共享则根据事务、订阅和阻塞命令决定。

错误二:把对象直接 toString() 存入 Redis

toString() 通常不是稳定序列化协议,字段变化后无法可靠反序列化,也可能包含敏感信息。应使用明确的 JSON、二进制协议或 Hash,并定义版本兼容规则。

错误三:锁释放直接 DEL

这会让过期后的旧持有者删除新持有者的锁。必须使用随机令牌和比较令牌的 Lua 脚本。

错误四:收到 Stream 消息就立即 ACK

如果先 ACK 再处理,进程在业务处理前崩溃,消息已经从 PEL 中移除,无法自动重试。通常应先完成幂等业务处理,再 ACK。

错误五:把 Stream 当作无限日志

Stream 会持续占用内存。必须定义保留策略、归档策略、死信处理和积压告警。ACK 不会自动删除历史消息。

错误六:超时后无条件重试写命令

响应丢失不等于服务端未执行。重试 XADDINCR 等命令可能造成重复结果。写入协议必须设计幂等性或可确认性。

错误七:用 Redis 锁代替数据库约束

锁可能因租约过期、进程暂停、故障切换而失效。最终业务正确性仍应由数据库条件更新、唯一约束、版本号和幂等记录等机制共同保证。


15. 选择 Redis、数据库和 Streams 的边界

Redis 适合低延迟访问、临时状态、计数、集合运算、短期缓存、协调信息和事件流转。关系数据库更适合需要持久约束、复杂查询、多表事务和审计的核心业务数据。

可以用以下因果关系判断:

  • 数据丢失后可以重建:适合缓存;
  • 数据丢失后必须审计恢复:需要持久化系统和明确备份;
  • 需要严格唯一性:使用数据库唯一约束或等价的强约束;
  • 需要至少一次异步处理:Streams 消费组可行,但必须幂等;
  • 需要跨数据库事务发布事件:使用 Outbox 等事务消息模式;
  • 需要全局强一致锁:先定义故障模型,再评估 Redis 锁是否满足,不应直接套用代码。

Lettuce 只负责把 Java 程序接入 Redis。真正可靠的 Redis 工程,依赖的是清晰的键和值协议、正确的连接生命周期、明确的原子边界、可验证的租约和所有权、可恢复的消息消费流程,以及对网络和故障的显式建模。


系列导航与关联阅读

官方资料

本文依据 Java、Spring 与相关项目官方文档重新梳理;正文、示例与生产清单由 WR BLOG 编写。

评论

0 条讨论
0/1000
还没有评论,来聊聊你的看法
会影响是否读取历史消息:\n\n- `0`:从已有历史开始;\n- ` Java Redis 工程:Lettuce、连接、序列化、缓存、锁和 Streams - WR Blog
WR Blog 加载中...
返回文章
JavaJava 25 LTSRedis

Java Redis 工程:Lettuce、连接、序列化、缓存、锁和 Streams

Java Redis 工程:Lettuce、连接、序列化、缓存、锁和 Streams封面

Java 基础体系 · 第 87/100 篇。示例统一以 Java 25 LTS 为语言和 JVM 基线;框架示例使用与其兼容的现代稳定版本。

Java Redis 工程:Lettuce、连接、序列化、缓存、锁和 Streams

Redis 工程问题通常不是“会不会调用 GETSET”,而是要同时处理六件事:

  1. Lettuce 客户端如何建立连接、复用连接并关闭资源;
  2. Redis 字节如何映射为 Java 类型;
  3. 缓存数据如何写入、过期、重建和失效;
  4. 多线程、多实例下的锁是否真的具有互斥和安全释放能力;
  5. Redis Streams 如何投递、确认、重试和恢复消息;
  6. 网络断开、超时、主从切换、进程崩溃时,业务状态如何变化。

本文以 Java 25 LTS 为运行范围,以 Lettuce 的同步、异步和响应式 API 为客户端基础。示例使用 Redis 7.x 常见命令语义;具体 Lettuce 版本应根据项目发布周期选择兼容版本,而不应把 Java 版本直接等同于 Redis 或 Lettuce 版本。


1. Redis 与 Lettuce 的职责边界

1.1 Redis 是服务端数据结构与协议

Redis 是一个通过 RESP 协议提供命令服务的内存数据存储。客户端发送命令:

SET user:42:name Alice EX 60

Redis 返回结果:

OK

服务端负责:

  • 解析命令;
  • 修改内存中的数据结构;
  • 执行过期判断;
  • 按连接返回响应;
  • 根据配置进行持久化或复制。

Redis 的命令通常具有单线程顺序执行的可见性特征:对于单个 Redis 实例,两个命令不会在同一时刻修改同一执行上下文。但是,这不意味着多个命令天然组成事务。例如:

GET stock:sku-1
# Java 计算 stock - 1
SET stock:sku-1 9

两个客户端可能都读到 10,最后都写入 9。单条命令的原子性不能自动扩展为多条命令的业务原子性。

1.2 Lettuce 是 Java 客户端

Lettuce 负责:

  • 建立 TCP 或 TLS 连接;
  • 编码 Java 参数为 Redis 协议;
  • 解码 Redis 响应;
  • 提供同步、异步和响应式接口;
  • 管理连接重连、拓扑刷新等客户端能力;
  • 将 Redis 协议错误、连接错误转换为 Java 异常或异步失败。

Lettuce 不会自动保证:

  • 你的对象序列化格式永远兼容;
  • 缓存一定不会击穿;
  • 分布式锁一定安全;
  • Streams 消费一定不重复;
  • Redis 故障时业务一定不丢数据。

这些是应用协议和故障模型问题。


2. 建立一个可运行的 Lettuce 工程

2.1 Maven 依赖

下面是一个普通 Maven 依赖示例。版本应在项目中统一管理,并在升级时查看对应 Lettuce 的兼容说明:

<dependency>
    <groupId>io.lettuce</groupId>
    <artifactId>lettuce-core</artifactId>
    <version>6.5.0.RELEASE</version>
</dependency>

Java 25 项目可以使用如下编译配置:

<properties>
    <maven.compiler.release>25</maven.compiler.release>
    <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
</properties>

启动本地 Redis:

docker run --rm --name redis-demo -p 6379:6379 redis:7

验证服务端:

redis-cli -h 127.0.0.1 -p 6379 PING

预期输出:

PONG

2.2 最小端到端程序

import io.lettuce.core.RedisClient;
import io.lettuce.core.RedisURI;
import io.lettuce.core.api.StatefulRedisConnection;
import io.lettuce.core.api.sync.RedisCommands;

import java.time.Duration;

public class LettuceDemo {
    public static void main(String[] args) {
        RedisURI uri = RedisURI.builder()
                .withHost("127.0.0.1")
                .withPort(6379)
                .withTimeout(Duration.ofSeconds(2))
                .build();

        RedisClient client = RedisClient.create(uri);

        try (StatefulRedisConnection<String, String> connection =
                     client.connect()) {
            RedisCommands<String, String> commands = connection.sync();

            commands.set("demo:name", "Alice");
            commands.expire("demo:name", 60);

            String value = commands.get("demo:name");
            System.out.println(value);       // Alice
            System.out.println(commands.ttl("demo:name")); // 约 60
        } finally {
            client.shutdown();
        }
    }
}

这里有三个不同生命周期:

  1. RedisClient:应用级客户端,通常创建一次;
  2. StatefulRedisConnection<K,V>:一个有状态的 Redis 连接;
  3. RedisCommands<K,V>:连接上的同步命令门面。

try 只负责关闭连接,finally 负责关闭客户端。Web 应用中通常不应每个请求都创建和销毁 RedisClient,否则会反复建立 TCP 连接、增加握手和线程管理成本。

2.3 连接 URI 与认证

生产环境通常至少配置地址、数据库、密码、TLS 和超时:

RedisURI uri = RedisURI.builder()
        .withHost("redis.example.internal")
        .withPort(6379)
        .withDatabase(0)
        .withAuthentication("app-user", "strong-password".toCharArray())
        .withSsl(true)
        .withTimeout(Duration.ofSeconds(2))
        .build();

密码不应硬编码在源代码中,应由环境变量、密钥管理系统或容器 Secret 注入。TLS 保护传输过程,但不会自动解决 Redis 命令权限;Redis ACL 仍需限制账号可执行的命令和键空间。


3. Lettuce 连接模型:同步、异步、响应式与并发

3.1 同步 API

同步调用会阻塞当前 Java 线程:

String value = connection.sync().get("key");

网络往返完成前,调用线程不能继续执行。同步 API 适合:

  • 业务代码本身是同步模型;
  • Redis 调用数量可控;
  • 已正确配置超时。

同步 API 不适合在事件循环线程中执行,因为一次网络阻塞会阻塞整个事件循环。

3.2 异步 API

var async = connection.async();

async.set("key", "value")
     .thenAccept(result -> System.out.println(result))
     .exceptionally(error -> {
         error.printStackTrace();
         return null;
     });

Lettuce 异步 API 返回 RedisFuture<T>,它与 CompletionStage 兼容。错误不会在调用处同步抛出,而通常会通过 Future 完成异常传播:

async.get("missing-key")
     .whenComplete((value, error) -> {
         if (error != null) {
             // 连接错误、Redis 错误等
             error.printStackTrace();
         } else {
             System.out.println(value);
         }
     });

如果业务需要严格超时,不能只依赖网络连接超时,还应对 Future 设置业务级超时:

async.get("key")
     .orTimeout(500, java.util.concurrent.TimeUnit.MILLISECONDS);

业务级超时只会让调用方停止等待,不一定能撤销服务端已经执行的命令。因此超时后的重试必须考虑命令是否已经成功执行。

3.3 响应式 API

Lettuce 提供 Reactor 适配的响应式命令接口。典型形式如下:

var reactive = connection.reactive();

reactive.set("key", "value")
        .then(reactive.get("key"))
        .doOnNext(System.out::println)
        .subscribe();

响应式调用通常是惰性的:只有订阅后才开始执行。若应用使用 Spring WebFlux 或其他响应式栈,应避免在响应式链中调用同步 API,否则会把阻塞操作带入非阻塞线程模型。

3.4 连接是否可以并发使用

Lettuce 的 StatefulConnection 通常设计为线程安全,可被多个线程共享。它内部会把请求写入连接,并按照协议关联响应。

但是,线程安全不等于业务操作隔离。以下场景不应依赖一个共享连接完成“只属于当前业务流程”的状态操作:

  • MULTI / EXEC 事务;
  • WATCH 乐观锁;
  • SUBSCRIBE 或阻塞式 XREADGROUP
  • 需要独占连接的交互流程;
  • 依赖连接级状态的命令。

例如,事务命令在连接上具有状态:

MULTI
SET a 1
SET b 2
EXEC

如果多个业务线程共享同一连接,命令可能交错,事务边界就会被破坏。此时应为该事务获取独立连接,或者使用 Lettuce 的连接池方案。连接池不是所有场景都必需;普通无状态命令共享一个连接通常足够,但事务、订阅和阻塞消费需要单独考虑。

3.5 Java 25 虚拟线程的边界

Java 25 的虚拟线程降低了阻塞线程的创建成本,但不会消除:

  • Redis 网络往返延迟;
  • Redis 服务端处理时间;
  • 连接数量限制;
  • 连接上的命令顺序和业务竞态。

可以让虚拟线程执行同步 Lettuce 调用,但仍需限制并发量,避免瞬间产生过多 Redis 请求。虚拟线程解决的是 Java 线程等待成本,不是 Redis 容量规划问题。


4. 序列化:Redis 存的是字节,不是 Java 对象

4.1 编解码器决定 Redis 看到的内容

Redis 的值本质上是字节数组。Lettuce 通过 RedisCodec<K,V> 将 Java 的键和值编码为字节,并把响应解码回来。

默认的 StringRedisCodec.UTF8 适合字符串:

var connection = client.connect(
        io.lettuce.core.codec.StringCodec.UTF8
);

此时:

commands.set("user:42", "{\"id\":42,\"name\":\"Alice\"}");

Redis CLI 可以直接观察值:

redis-cli GET user:42

输出:

{"id":42,"name":"Alice"}

如果使用 Java 原生序列化、Kryo、JSON、Protobuf 或自定义二进制协议,Redis CLI 看到的内容可能不可读。可读性并不等于正确性,但可读格式通常更便于诊断。

4.2 字符串 JSON 的显式序列化

一种稳妥的应用边界是:Redis 层只接受 String,对象序列化在业务代码中明确完成。以 Jackson 为例:

record User(long id, String name) {}

ObjectMapper mapper = new ObjectMapper();

User user = new User(42, "Alice");
String json = mapper.writeValueAsString(user);

commands.set("user:42", json);

String stored = commands.get("user:42");
User restored = mapper.readValue(stored, User.class);

这里必须明确处理 JsonProcessingException,因为序列化失败会导致缓存写入失败,反序列化失败则可能造成读取路径异常。

更重要的是,JSON 格式需要版本策略。例如旧值:

{"id":42,"name":"Alice"}

新代码增加字段:

record User(long id, String name, String level) {}

如果新字段没有默认兼容处理,旧缓存反序列化可能失败。常见做法是:

  • 新字段允许缺省;
  • 在 JSON 中增加显式版本字段;
  • 变更结构时修改键前缀,例如 v2:user:42
  • 不能解析的缓存直接删除并回源,而不是无限重试。

4.3 不要默认使用不安全的多态反序列化

某些序列化框架允许数据携带 Java 类型信息,然后按类型实例化对象。若数据来源可被外部写入,宽泛的多态反序列化可能造成反序列化攻击。

缓存数据也不能因为“只是缓存”就跳过安全边界。应优先使用固定目标类型:

User user = mapper.readValue(json, User.class);

而不是接受任意类名并动态实例化。

4.4 Hash、JSON 字符串与结构化数据

Redis Hash 可以把对象字段拆开:

commands.hset("user:42", java.util.Map.of(
        "id", "42",
        "name", "Alice"
));

读取:

java.util.Map<String, String> fields =
        commands.hgetall("user:42");

Hash 的优点是可以只更新一个字段,缺点是字段类型通常被转换成字符串,结构演进和原子更新逻辑需要应用自己维护。

JSON 字符串的优点是对象边界清晰,读写一次完成;缺点是修改单个字段通常需要读、修改、写,或者依赖 RedisJSON 等额外模块。选择哪一种取决于访问模式,不是“对象一定 JSON、字段一定 Hash”的固定规则。

4.5 键名也是协议

键名应包含稳定的命名空间和版本信息:

app:v1:user:42
app:v1:product:sku-100
app:v1:lock:order:9001
app:v1:stream:order-events

冒号只是约定,不是 Redis 的层级目录。键名协议需要定义:

  • 业务域;
  • 数据类型;
  • 版本;
  • 标识符编码;
  • 是否允许跨服务共享。

若不同服务对同一个键使用不同序列化格式,即使键名相同,也只是制造了隐蔽的数据兼容故障。


5. Redis 数据结构与原子性基础

5.1 常用类型

Redis 类型 典型用途 关键边界
String 缓存对象、计数器、锁令牌 复杂更新需 Lua 或条件命令
Hash 对象字段、属性集合 字段类型和版本由应用维护
List 简单队列、双端队列 消费确认能力弱于 Streams
Set 去重、集合关系 无序
Sorted Set 排行榜、延迟任务索引 分数相同需定义排序规则
Stream 追加事件、消费组 需要 ACK、重试和积压治理

Redis 的原子性一般指一条命令在服务端执行期间不会被另一条命令插入。INCR 是原子递增:

INCR page:views

但“读取后计算再写入”不是原子的。可以用服务端命令、Lua 脚本或事务重新建立原子边界。

5.2 TTL 不是永久保证

SET key value EX 60

表示设置值并附带 60 秒过期时间。相比先 SETEXPIRE,单条 SET ... EX 避免了两条命令之间进程崩溃导致“有值但无 TTL”。

检查:

TTL key

结果含义通常是:

  • >= 0:剩余秒数;
  • -1:键存在但没有过期时间;
  • -2:键不存在。

过期删除可能包含惰性删除和主动扫描,因此 TTL 到期不应被理解为精确到某一毫秒的定时器。若业务要求严格的时间调度,Redis TTL 不能单独承担该职责。


6. 缓存:从读路径到故障路径

6.1 Cache-aside 模式

Cache-aside,也称旁路缓存,由应用显式管理缓存:

读请求
  │
  ├─ GET cache
  │    ├─ 命中:返回缓存
  │    └─ 未命中
  │          │
  │          ├─ 查询数据库
  │          ├─ SET cache EX ttl
  │          └─ 返回结果
  │
写请求
  ├─ 更新数据库
  └─ 删除缓存

读取伪代码:

String key = "app:v1:user:" + userId;
String cached = commands.get(key);

if (cached != null) {
    return mapper.readValue(cached, User.class);
}

User user = userRepository.findById(userId); // 回源数据库
if (user != null) {
    String json = mapper.writeValueAsString(user);
    commands.setex(key, 300, json);
}
return user;

这里的 setex 是常见同步 API 形式;也可以使用 set(key, value, SetArgs.Builder.ex(300)),后者更明确地表达命令参数。

6.2 为什么写数据库后删除缓存

假设采用“更新数据库后删除缓存”:

T1:更新数据库为 B
T1:删除缓存

后续读取会从数据库加载 B,再写回缓存。若改成“更新数据库后更新缓存”,多个并发写请求可能按不同顺序完成,最终缓存值不一定对应最后一次数据库提交。

但删除缓存也不是绝对完美。存在如下竞态:

T1:读取数据库,得到旧值 A
T2:更新数据库为 B
T2:删除缓存
T1:把旧值 A 写回缓存

最终缓存重新出现旧值 A。常见缓解方式包括:

  • 缓存写入时携带版本号,只接受更高版本;
  • 延迟双删;
  • 通过消息或 CDC 进行失效传播;
  • 对关键读路径增加短暂一致性校验。

这些方法不能消除所有分布式故障,只能改变旧值重新出现的概率和可检测性。若业务要求强一致,应直接读数据库或采用带版本条件的更新协议。

6.3 缓存穿透、击穿、雪崩

缓存穿透:请求查询一个数据库中不存在的键,缓存也不存,恶意或高频请求持续打到数据库。

缓解方法:

commands.setex("app:v1:user:not-found:" + id, 30, "1");

这叫负缓存。它必须有较短 TTL,否则新创建的数据可能长期被“未找到”遮蔽。布隆过滤器可以降低无效查询,但存在误判,不能把“布隆过滤器不存在”之外的所有逻辑写错。

缓存击穿:某个热点键过期,大量线程同时回源。可以使用互斥重建锁,但锁失败的线程应等待、返回旧值或降级,而不是全部继续查询数据库。

缓存雪崩:大量键在相近时间过期,造成集中回源。可以给 TTL 增加随机扰动:

long ttl = 300 + ThreadLocalRandom.current().nextLong(60);
commands.set(key, json, SetArgs.Builder.ex(ttl));

随机 TTL 只是分散过期时间,不能替代容量、限流和回源保护。

6.4 逻辑过期与物理过期

物理过期依赖 Redis TTL,过期后键不可读。逻辑过期则把过期时间写入值:

{
  "data": {"id": 42, "name": "Alice"},
  "expireAt": 1730000000000
}

读取到逻辑过期数据后,可以先返回旧值,再异步刷新。它牺牲了一段时间的新鲜度,换取热点请求不同时阻塞回源。

两者可以组合:Redis 设置较长物理 TTL,值内部设置较短逻辑 TTL。物理 TTL 防止刷新任务永久失效,逻辑 TTL 控制业务新鲜度。

6.5 缓存错误如何处理

缓存读取失败时通常有两种策略:

  • 旁路降级:记录错误,直接回源数据库;
  • 快速失败:对核心链路不允许继续施压数据库,返回降级结果或错误。

不能无条件地把所有 Redis 异常转成缓存未命中。在 Redis 大面积故障时,所有请求回源可能造成数据库雪崩。应通过超时、并发舱壁、限流和降级控制回源规模。


7. 分布式锁:互斥、令牌与安全释放

7.1 一个锁至少需要什么性质

设锁键为 lock:order:42。一个基本分布式锁需要:

  • 互斥:同一时刻最多一个持有者;
  • 可释放:持有者完成后能删除;
  • 租约:持有者崩溃后锁最终能自动消失;
  • 所有权校验:旧持有者不能删除新持有者的锁。

只执行:

SET lock:order:42 owner EX 30

并不能保证互斥,因为没有 NX。正确的获取条件是:

SET lock:order:42 random-token NX EX 30

Java 示例:

String token = UUID.randomUUID().toString();
String result = commands.set(
        "app:v1:lock:order:42",
        token,
        SetArgs.Builder.nx().ex(30)
);

boolean acquired = "OK".equals(result);

NX 表示仅当键不存在时设置,EX 30 设置租约。两者放在同一个 SET 命令中,避免“先判断再设置”的竞态。

7.2 为什么释放不能直接 DEL

以下流程会错误释放别人的锁:

T1 获取锁,租约 30 秒
T1 执行任务超过 30 秒
锁自动过期
T2 获取同一把锁
T1 完成并执行 DEL
T2 的锁被 T1 删除

因此释放必须比较令牌,只有值仍等于自己的令牌时才删除。比较和删除必须是一个原子操作,Lua 脚本可表达:

if redis.call('get', KEYS[1]) == ARGV[1] then
    return redis.call('del', KEYS[1])
else
    return 0
end

Lettuce 调用:

String unlockScript = """
        if redis.call('get', KEYS[1]) == ARGV[1] then
            return redis.call('del', KEYS[1])
        else
            return 0
        end
        """;

Long deleted = commands.eval(
        unlockScript,
        ScriptOutputType.INTEGER,
        new String[]{"app:v1:lock:order:42"},
        token
);

if (deleted == 1L) {
    System.out.println("lock released");
}

释放结果为 1 表示删除了自己的锁,0 表示锁已不存在或已经属于其他持有者。

7.3 租约时间不是任务时间

如果任务最长执行时间为 TT,锁租约为 LL,还存在网络延迟、调度暂停和 GC 暂停 DD,要避免锁在任务完成前过期,至少需要满足:

L>T+D+网络与调度余量L > T + D + \text{网络与调度余量}

TT 往往不是严格上界,因此实际工程中还需要:

  • 自动续租;
  • 任务分段并校验所有权;
  • 让业务提交带 fencing token;
  • 接受锁只是“减少并发”,不是强一致事务。

Fencing token 是递增的所有权序号。每次成功获取锁获得更大的 token,写入下游资源时把 token 一起提交,下游拒绝比当前 token 更旧的写入。这样即使旧持有者因暂停后恢复,也不能覆盖新持有者的结果。

7.4 Redis 锁的边界

单 Redis 实例上的 SET NX EX 可以提供单实例租约锁语义,但在主从异步复制、故障切换和网络分区下,锁的全局安全性取决于部署模型和业务要求。

Redis 锁不等价于数据库事务,也不自动保护数据库写入。如果锁用于支付、库存等关键状态,通常还需要数据库唯一约束、条件更新、幂等键或 fencing token。将锁作为唯一正确性保障,是常见误解。

Redlock 等多节点算法有自己的假设、时钟和故障模型争议。使用前必须明确要解决的是:

  • 防止重复执行的性能优化;
  • 还是严格的线性一致性所有权。

两者的可靠性要求不同,不能只因为“部署了多个 Redis 节点”就认为问题已经解决。


8. Redis Streams:追加日志与消费组

8.1 Stream 的基本结构

Redis Stream 是按消息 ID 排序的追加型日志。消息形如:

ID                  field       value
1710000000000-0     orderId     9001
                    amount      100

生产消息:

String id = commands.xadd(
        "app:v1:stream:orders",
        java.util.Map.of(
                "orderId", "9001",
                "amount", "100"
        )
);

System.out.println(id); // 例如 1710000000000-0

消息 ID 通常由毫秒时间和序号组成。应用不应把 ID 当作业务唯一幂等键的唯一来源;业务事件仍应携带自己的 eventId

8.2 普通读取

var messages = commands.xread(
        io.lettuce.core.XReadArgs.Builder.count(10),
        io.lettuce.core.StreamOffset.from("app:v1:stream:orders", "0-0")
);

for (var message : messages) {
    System.out.println(message.getId());
    System.out.println(message.getBody());
}

0-0 表示从较早位置开始读取。普通 XREAD 不提供消费组的待确认列表,适合简单读取或广播式场景,但不适合需要明确 ACK 和故障接管的工作队列。

8.3 消费组、消费者和 PEL

消费组模型包含:

  • Stream:消息日志;
  • Consumer Group:一组协作消费者;
  • Consumer:组内一个逻辑消费者名称;
  • PEL(Pending Entries List):已经投递但尚未确认的消息集合。

创建消费组:

XGROUP CREATE app:v1:stream:orders orders-group 0 MKSTREAM

如果组可能已经存在,应用应把“组已存在”作为可接受结果处理,而不是把启动失败当成系统故障。Lettuce 中可以捕获 Redis 的 BUSYGROUP 错误。

读取新消息:

var records = commands.xreadgroup(
        Consumer.from("orders-group", "consumer-1"),
        XReadArgs.Builder.count(10).block(Duration.ofSeconds(5)),
        StreamOffset.lastConsumed("app:v1:stream:orders")
);

lastConsumed 表示从消费组记录的位置继续读取。对于新建组,起始 ID 是 0 还是 $ 会影响是否读取历史消息:

  • 0:从已有历史开始;
  • $:只关注创建组之后的新消息。

处理成功后确认:

for (var record : records) {
    try {
        handle(record.getBody());
        commands.xack(
                "app:v1:stream:orders",
                "orders-group",
                record.getId()
        );
    } catch (Exception e) {
        // 不 ACK,让消息保留在 PEL 中
        // 记录错误并等待重试或转入死信
    }
}

ACK 不是删除消息。 XACK 只从消费组的待确认集合中移除消息,Stream 中的消息仍然存在,直到通过保留策略或 XDEL 等方式清理。

8.4 崩溃与消息重投

考虑以下时序:

C1 读取消息 M
C1 业务处理成功
C1 进程在 XACK 前崩溃
C2 后续接管 M
C2 再次处理 M

因此 Streams 消费通常是至少一次投递,不是自动恰好一次。消费者必须幂等。例如数据库处理可以使用唯一约束:

CREATE TABLE processed_event (
    event_id VARCHAR(128) PRIMARY KEY,
    processed_at TIMESTAMP NOT NULL
);

处理事件时:

  1. 尝试插入 event_id
  2. 若唯一键冲突,说明已经处理过,直接 ACK;
  3. 插入成功后执行业务更新;
  4. 业务成功后 ACK。

但“插入幂等表”和“业务更新”如果不在同一个数据库事务中,仍可能出现中间失败。因此幂等记录与业务变更最好处于同一事务内。

8.5 接管空闲 Pending 消息

消费者崩溃后,消息仍属于旧消费者。其他消费者需要根据空闲时间接管。Redis 版本支持的命令和 Lettuce API 形式可能随版本变化,常见机制包括:

  • XAUTOCLAIM:按最小空闲时间自动扫描并认领;
  • XCLAIM:指定消息 ID 认领;
  • XPENDING:查看待确认消息和空闲时间。

运维或消费逻辑应先检查:

XPENDING app:v1:stream:orders orders-group

关注:

  • Pending 总数;
  • 最早消息 ID;
  • 消费者分布;
  • 长时间未确认的消息。

重试流程应有次数限制。超过阈值后写入死信 Stream:

XADD app:v1:stream:orders-dlq * orderId 9001 reason timeout

然后 ACK 原消息,避免坏消息永久阻塞消费。死信不是删除错误,它是把错误从主处理路径隔离出来,之后仍需人工或自动修复。

8.6 消费组并发和顺序

同一个消费组内,一条消息只会投递给一个消费者,但不同消息可以被不同消费者并行处理。因此:

  • 组内不保证所有消息的全局业务处理顺序;
  • 同一订单的事件如果要求顺序,需要按订单分区、串行处理或使用版本号校验;
  • eventId 幂等解决重复,不解决乱序;
  • 业务版本号可以拒绝旧事件,例如只接受 eventVersion > currentVersion

阻塞读取连接应与普通命令连接分离。一个 XREADGROUP BLOCK 长时间等待,不应占用事务、健康检查或普通请求使用的连接。


9. 事务、Lua 与批量操作的选择

9.1 Redis 事务不是回滚事务

Redis MULTI/EXEC 会把命令排队后依次执行,但通常没有关系数据库那样的回滚机制。若某个命令执行时报错,之前成功执行的命令不会自动撤销。

适合使用 Redis 事务的场景:

  • 需要将一组命令作为连续执行单元;
  • 通过 WATCH 实现乐观并发控制;
  • 不要求异常时自动回滚。

需要条件判断、比较值和修改值时,Lua 脚本通常更直接,因为脚本在 Redis 服务端原子执行:

local current = redis.call('get', KEYS[1])
if current == ARGV[1] then
    redis.call('set', KEYS[1], ARGV[2], 'EX', ARGV[3])
    return 1
end
return 0

脚本应保持短小,避免阻塞 Redis 主执行线程。Lua 不能把耗时数据库查询放进去;Redis 脚本只能处理 Redis 内部数据。

9.2 Pipeline 不是事务

Pipeline 把多个命令批量发送,减少网络往返:

commands.set("k1", "v1");
commands.set("k2", "v2");
commands.get("k1");

具体 Lettuce API 可通过异步批量、自动 flush 控制等方式实现。Pipeline 的核心是传输优化,不保证命令之间的事务原子性。网络断开时,客户端可能不知道批量中的哪些命令已经在服务端执行,因此重试必须基于幂等性设计。


10. 超时、重连与错误诊断

10.1 三类时间限制

Redis 客户端至少要区分:

  1. 连接建立超时:TCP、TLS、认证建立多久;
  2. 命令响应超时:发送命令后等待响应多久;
  3. 业务截止时间:整个请求允许 Redis 占用多久。

业务截止时间通常最严格。即使底层允许 2 秒,HTTP 请求只剩 100 毫秒,也不应再发一个必然超时的 Redis 操作。

10.2 错误分类

典型错误包括:

  • RedisCommandExecutionException:服务端返回错误,例如类型错误、权限错误、脚本错误;
  • RedisConnectionException:连接不可用或断开;
  • 超时异常:响应未在指定时间内到达;
  • 解码异常:响应字节无法按指定 Codec 转换。

诊断时至少记录:

  • 命令类型,而不是完整敏感值;
  • 键名的脱敏版本;
  • Redis 节点;
  • 超时时间;
  • 重试次数;
  • 请求 trace ID;
  • PING、连接状态和 Redis INFO 指标。

不要在日志中记录密码、锁令牌、用户隐私数据和完整缓存 JSON。

10.3 重试的危险

以下命令天然更适合重试:

GET
SET
DEL
EXPIRE

但“适合”不表示无条件安全。对于:

INCR
XADD
LPUSH

请求超时后,命令可能已经成功执行,只是响应丢失;直接重试可能造成重复计数或重复消息。解决方案包括:

  • 使用业务幂等键;
  • 让事件 ID 可去重;
  • 使用条件写入;
  • 查询结果确认,而不是盲目重发;
  • 明确接受重复并在下游消除。

11. 内存、持久化与高可用边界

Redis 是内存优先系统。即使配置了 RDB 或 AOF,也不能简单等同于每次写入都已同步落盘。

11.1 RDB 与 AOF

  • RDB:周期性生成快照,恢复较快,但两次快照之间的数据可能丢失;
  • AOF:记录写命令,通常能减少数据丢失窗口,但文件更大、重写和恢复成本不同;
  • appendfsync alwayseverysecno 在性能和持久性之间有不同取舍。

缓存数据通常可以接受丢失;订单事件、任务状态则需要更严格地评估持久化和恢复策略。

11.2 复制不是强一致

Redis 主从复制通常是异步的。主节点确认写入后,从节点可能尚未收到数据。主节点故障切换时,最近写入可能丢失。

因此:

  • 缓存丢失通常表现为回源;
  • Streams 消息丢失可能表现为业务事件缺口;
  • 锁状态在故障切换中可能出现旧锁和新锁并存的风险;
  • 需要关键持久化语义时,数据库或专门消息系统可能更适合。

Redis Cluster 将键分布到不同槽位。涉及多个键的事务或 Lua 脚本通常要求相关键落在同一槽位,可以使用 hash tag:

order:{9001}:state
order:{9001}:lock

{9001} 中的内容用于计算槽位,使两个键更可能位于同一槽。键设计必须在业务建模阶段完成,不能等跨槽错误出现后再补救。


12. 一个完整的订单事件流程

把缓存、锁和 Streams 放在一起,可以形成如下流程:

sequenceDiagram
    participant API as Java API
    participant R as Redis
    participant DB as Database
    participant C as Stream Consumer

    API->>R: GET user:v1:42
    alt cache hit
        R-->>API: JSON
    else cache miss
        API->>R: SET lock:user:42 token NX EX
        alt lock acquired
            API->>DB: SELECT user
            DB-->>API: user
            API->>R: SET user:v1:42 JSON EX
            API->>R: Lua compare-token DEL
        else lock not acquired
            API->>R: GET user:v1:42 or wait/backoff
        end
    end

    API->>DB: UPDATE order
    API->>R: XADD order-events eventId...
    R-->>API: stream message ID

    C->>R: XREADGROUP
    R-->>C: pending message
    C->>DB: idempotent business transaction
    C->>R: XACK

关键路径有三个不同的一致性边界:

  1. 缓存读取是性能优化边界,缓存丢失后可以回源;
  2. 锁是并发协调边界,需要令牌、租约和安全释放;
  3. Stream 是事件投递边界,需要 ACK、幂等、重试和积压监控。

如果“更新数据库”和“发布 Stream 事件”必须严格同时成功,直接先后执行仍可能发生:

数据库提交成功
Java 进程崩溃
XADD 尚未执行

更可靠的做法是数据库 Outbox:

  1. 在同一个数据库事务中更新业务表并写入 Outbox 表;
  2. 独立发布器读取 Outbox;
  3. 成功 XADD 后标记 Outbox 已发布;
  4. 发布器重复运行也必须幂等。

这体现了一个重要边界:Redis Streams 可以作为事件传输目标,但不能自动提供数据库事务的跨系统原子提交。


13. 生产验证与诊断命令

13.1 检查键和 TTL

redis-cli TYPE app:v1:user:42
redis-cli TTL app:v1:user:42
redis-cli MEMORY USAGE app:v1:user:42

TYPE 与客户端预期不符,例如代码执行 HGETALL 却返回 String 类型,Redis 会返回 WRONGTYPE 错误。这通常意味着键名冲突、版本迁移遗漏或序列化模型改变。

13.2 检查 Streams 积压

redis-cli XLEN app:v1:stream:orders
redis-cli XINFO GROUPS app:v1:stream:orders
redis-cli XPENDING app:v1:stream:orders orders-group

重点观察:

  • Stream 长度是否持续增长;
  • lag 是否增加;
  • Pending 是否持续增加;
  • 是否有消费者长时间无心跳;
  • 是否存在单个消费者占据大量 Pending。

清理 Stream 需要结合保留策略。盲目执行:

redis-cli DEL app:v1:stream:orders

会删除整个 Stream 及消费组,属于破坏性操作。生产清理应先确认备份、消费进度、恢复方案,再使用长度或时间窗口保留策略。

13.3 监控指标

应用和 Redis 两侧都应监控:

  • 命令延迟分位数;
  • 连接数和连接失败;
  • 超时和重试次数;
  • 缓存命中率;
  • 热点键访问;
  • 内存使用、淘汰数量;
  • Stream 长度、Pending 数、最老 Pending 空闲时间;
  • Lua 脚本耗时;
  • 主从复制延迟和故障切换。

只监控 Redis CPU 而不监控缓存回源量,无法判断缓存故障是否已经传导到数据库。


14. 常见错误及其根因

错误一:每个请求创建一个 RedisClient

根因是把客户端对象误当作短生命周期连接。结果是连接建立、线程和资源开销增加,故障时还可能产生连接风暴。

改为应用启动时创建客户端,关闭时统一释放;连接是否共享则根据事务、订阅和阻塞命令决定。

错误二:把对象直接 toString() 存入 Redis

toString() 通常不是稳定序列化协议,字段变化后无法可靠反序列化,也可能包含敏感信息。应使用明确的 JSON、二进制协议或 Hash,并定义版本兼容规则。

错误三:锁释放直接 DEL

这会让过期后的旧持有者删除新持有者的锁。必须使用随机令牌和比较令牌的 Lua 脚本。

错误四:收到 Stream 消息就立即 ACK

如果先 ACK 再处理,进程在业务处理前崩溃,消息已经从 PEL 中移除,无法自动重试。通常应先完成幂等业务处理,再 ACK。

错误五:把 Stream 当作无限日志

Stream 会持续占用内存。必须定义保留策略、归档策略、死信处理和积压告警。ACK 不会自动删除历史消息。

错误六:超时后无条件重试写命令

响应丢失不等于服务端未执行。重试 XADDINCR 等命令可能造成重复结果。写入协议必须设计幂等性或可确认性。

错误七:用 Redis 锁代替数据库约束

锁可能因租约过期、进程暂停、故障切换而失效。最终业务正确性仍应由数据库条件更新、唯一约束、版本号和幂等记录等机制共同保证。


15. 选择 Redis、数据库和 Streams 的边界

Redis 适合低延迟访问、临时状态、计数、集合运算、短期缓存、协调信息和事件流转。关系数据库更适合需要持久约束、复杂查询、多表事务和审计的核心业务数据。

可以用以下因果关系判断:

  • 数据丢失后可以重建:适合缓存;
  • 数据丢失后必须审计恢复:需要持久化系统和明确备份;
  • 需要严格唯一性:使用数据库唯一约束或等价的强约束;
  • 需要至少一次异步处理:Streams 消费组可行,但必须幂等;
  • 需要跨数据库事务发布事件:使用 Outbox 等事务消息模式;
  • 需要全局强一致锁:先定义故障模型,再评估 Redis 锁是否满足,不应直接套用代码。

Lettuce 只负责把 Java 程序接入 Redis。真正可靠的 Redis 工程,依赖的是清晰的键和值协议、正确的连接生命周期、明确的原子边界、可验证的租约和所有权、可恢复的消息消费流程,以及对网络和故障的显式建模。


系列导航与关联阅读

官方资料

本文依据 Java、Spring 与相关项目官方文档重新梳理;正文、示例与生产清单由 WR BLOG 编写。

评论

0 条讨论
0/1000
还没有评论,来聊聊你的看法
:只关注创建组之后的新消息。\n\n处理成功后确认:\n\n```java\nfor (var record : records) {\n try {\n handle(record.getBody());\n commands.xack(\n \"app:v1:stream:orders\",\n \"orders-group\",\n record.getId()\n );\n } catch (Exception e) {\n // 不 ACK,让消息保留在 PEL 中\n // 记录错误并等待重试或转入死信\n }\n}\n```\n\n**ACK 不是删除消息。** `XACK` 只从消费组的待确认集合中移除消息,Stream 中的消息仍然存在,直到通过保留策略或 `XDEL` 等方式清理。\n\n### 8.4 崩溃与消息重投\n\n考虑以下时序:\n\n```text\nC1 读取消息 M\nC1 业务处理成功\nC1 进程在 XACK 前崩溃\nC2 后续接管 M\nC2 再次处理 M\n```\n\n因此 Streams 消费通常是**至少一次投递**,不是自动恰好一次。消费者必须幂等。例如数据库处理可以使用唯一约束:\n\n```sql\nCREATE TABLE processed_event (\n event_id VARCHAR(128) PRIMARY KEY,\n processed_at TIMESTAMP NOT NULL\n);\n```\n\n处理事件时:\n\n1. 尝试插入 `event_id`;\n2. 若唯一键冲突,说明已经处理过,直接 ACK;\n3. 插入成功后执行业务更新;\n4. 业务成功后 ACK。\n\n但“插入幂等表”和“业务更新”如果不在同一个数据库事务中,仍可能出现中间失败。因此幂等记录与业务变更最好处于同一事务内。\n\n### 8.5 接管空闲 Pending 消息\n\n消费者崩溃后,消息仍属于旧消费者。其他消费者需要根据空闲时间接管。Redis 版本支持的命令和 Lettuce API 形式可能随版本变化,常见机制包括:\n\n- `XAUTOCLAIM`:按最小空闲时间自动扫描并认领;\n- `XCLAIM`:指定消息 ID 认领;\n- `XPENDING`:查看待确认消息和空闲时间。\n\n运维或消费逻辑应先检查:\n\n```text\nXPENDING app:v1:stream:orders orders-group\n```\n\n关注:\n\n- Pending 总数;\n- 最早消息 ID;\n- 消费者分布;\n- 长时间未确认的消息。\n\n重试流程应有次数限制。超过阈值后写入死信 Stream:\n\n```text\nXADD app:v1:stream:orders-dlq * orderId 9001 reason timeout\n```\n\n然后 ACK 原消息,避免坏消息永久阻塞消费。死信不是删除错误,它是把错误从主处理路径隔离出来,之后仍需人工或自动修复。\n\n### 8.6 消费组并发和顺序\n\n同一个消费组内,一条消息只会投递给一个消费者,但不同消息可以被不同消费者并行处理。因此:\n\n- 组内不保证所有消息的全局业务处理顺序;\n- 同一订单的事件如果要求顺序,需要按订单分区、串行处理或使用版本号校验;\n- `eventId` 幂等解决重复,不解决乱序;\n- 业务版本号可以拒绝旧事件,例如只接受 `eventVersion > currentVersion`。\n\n阻塞读取连接应与普通命令连接分离。一个 `XREADGROUP BLOCK` 长时间等待,不应占用事务、健康检查或普通请求使用的连接。\n\n---\n\n## 9. 事务、Lua 与批量操作的选择\n\n### 9.1 Redis 事务不是回滚事务\n\nRedis `MULTI/EXEC` 会把命令排队后依次执行,但通常没有关系数据库那样的回滚机制。若某个命令执行时报错,之前成功执行的命令不会自动撤销。\n\n适合使用 Redis 事务的场景:\n\n- 需要将一组命令作为连续执行单元;\n- 通过 `WATCH` 实现乐观并发控制;\n- 不要求异常时自动回滚。\n\n需要条件判断、比较值和修改值时,Lua 脚本通常更直接,因为脚本在 Redis 服务端原子执行:\n\n```lua\nlocal current = redis.call('get', KEYS[1])\nif current == ARGV[1] then\n redis.call('set', KEYS[1], ARGV[2], 'EX', ARGV[3])\n return 1\nend\nreturn 0\n```\n\n脚本应保持短小,避免阻塞 Redis 主执行线程。Lua 不能把耗时数据库查询放进去;Redis 脚本只能处理 Redis 内部数据。\n\n### 9.2 Pipeline 不是事务\n\nPipeline 把多个命令批量发送,减少网络往返:\n\n```java\ncommands.set(\"k1\", \"v1\");\ncommands.set(\"k2\", \"v2\");\ncommands.get(\"k1\");\n```\n\n具体 Lettuce API 可通过异步批量、自动 flush 控制等方式实现。Pipeline 的核心是传输优化,不保证命令之间的事务原子性。网络断开时,客户端可能不知道批量中的哪些命令已经在服务端执行,因此重试必须基于幂等性设计。\n\n---\n\n## 10. 超时、重连与错误诊断\n\n### 10.1 三类时间限制\n\nRedis 客户端至少要区分:\n\n1. **连接建立超时**:TCP、TLS、认证建立多久;\n2. **命令响应超时**:发送命令后等待响应多久;\n3. **业务截止时间**:整个请求允许 Redis 占用多久。\n\n业务截止时间通常最严格。即使底层允许 2 秒,HTTP 请求只剩 100 毫秒,也不应再发一个必然超时的 Redis 操作。\n\n### 10.2 错误分类\n\n典型错误包括:\n\n- `RedisCommandExecutionException`:服务端返回错误,例如类型错误、权限错误、脚本错误;\n- `RedisConnectionException`:连接不可用或断开;\n- 超时异常:响应未在指定时间内到达;\n- 解码异常:响应字节无法按指定 Codec 转换。\n\n诊断时至少记录:\n\n- 命令类型,而不是完整敏感值;\n- 键名的脱敏版本;\n- Redis 节点;\n- 超时时间;\n- 重试次数;\n- 请求 trace ID;\n- PING、连接状态和 Redis `INFO` 指标。\n\n不要在日志中记录密码、锁令牌、用户隐私数据和完整缓存 JSON。\n\n### 10.3 重试的危险\n\n以下命令天然更适合重试:\n\n```text\nGET\nSET\nDEL\nEXPIRE\n```\n\n但“适合”不表示无条件安全。对于:\n\n```text\nINCR\nXADD\nLPUSH\n```\n\n请求超时后,命令可能已经成功执行,只是响应丢失;直接重试可能造成重复计数或重复消息。解决方案包括:\n\n- 使用业务幂等键;\n- 让事件 ID 可去重;\n- 使用条件写入;\n- 查询结果确认,而不是盲目重发;\n- 明确接受重复并在下游消除。\n\n---\n\n## 11. 内存、持久化与高可用边界\n\nRedis 是内存优先系统。即使配置了 RDB 或 AOF,也不能简单等同于每次写入都已同步落盘。\n\n### 11.1 RDB 与 AOF\n\n- RDB:周期性生成快照,恢复较快,但两次快照之间的数据可能丢失;\n- AOF:记录写命令,通常能减少数据丢失窗口,但文件更大、重写和恢复成本不同;\n- `appendfsync always`、`everysec`、`no` 在性能和持久性之间有不同取舍。\n\n缓存数据通常可以接受丢失;订单事件、任务状态则需要更严格地评估持久化和恢复策略。\n\n### 11.2 复制不是强一致\n\nRedis 主从复制通常是异步的。主节点确认写入后,从节点可能尚未收到数据。主节点故障切换时,最近写入可能丢失。\n\n因此:\n\n- 缓存丢失通常表现为回源;\n- Streams 消息丢失可能表现为业务事件缺口;\n- 锁状态在故障切换中可能出现旧锁和新锁并存的风险;\n- 需要关键持久化语义时,数据库或专门消息系统可能更适合。\n\nRedis Cluster 将键分布到不同槽位。涉及多个键的事务或 Lua 脚本通常要求相关键落在同一槽位,可以使用 hash tag:\n\n```text\norder:{9001}:state\norder:{9001}:lock\n```\n\n`{9001}` 中的内容用于计算槽位,使两个键更可能位于同一槽。键设计必须在业务建模阶段完成,不能等跨槽错误出现后再补救。\n\n---\n\n## 12. 一个完整的订单事件流程\n\n把缓存、锁和 Streams 放在一起,可以形成如下流程:\n\n```mermaid\nsequenceDiagram\n participant API as Java API\n participant R as Redis\n participant DB as Database\n participant C as Stream Consumer\n\n API->>R: GET user:v1:42\n alt cache hit\n R--\u003e>API: JSON\n else cache miss\n API->>R: SET lock:user:42 token NX EX\n alt lock acquired\n API->>DB: SELECT user\n DB--\u003e>API: user\n API->>R: SET user:v1:42 JSON EX\n API->>R: Lua compare-token DEL\n else lock not acquired\n API->>R: GET user:v1:42 or wait/backoff\n end\n end\n\n API->>DB: UPDATE order\n API->>R: XADD order-events eventId...\n R--\u003e>API: stream message ID\n\n C->>R: XREADGROUP\n R--\u003e>C: pending message\n C->>DB: idempotent business transaction\n C->>R: XACK\n```\n\n关键路径有三个不同的一致性边界:\n\n1. 缓存读取是性能优化边界,缓存丢失后可以回源;\n2. 锁是并发协调边界,需要令牌、租约和安全释放;\n3. Stream 是事件投递边界,需要 ACK、幂等、重试和积压监控。\n\n如果“更新数据库”和“发布 Stream 事件”必须严格同时成功,直接先后执行仍可能发生:\n\n```text\n数据库提交成功\nJava 进程崩溃\nXADD 尚未执行\n```\n\n更可靠的做法是数据库 Outbox:\n\n1. 在同一个数据库事务中更新业务表并写入 Outbox 表;\n2. 独立发布器读取 Outbox;\n3. 成功 `XADD` 后标记 Outbox 已发布;\n4. 发布器重复运行也必须幂等。\n\n这体现了一个重要边界:Redis Streams 可以作为事件传输目标,但不能自动提供数据库事务的跨系统原子提交。\n\n---\n\n## 13. 生产验证与诊断命令\n\n### 13.1 检查键和 TTL\n\n```bash\nredis-cli TYPE app:v1:user:42\nredis-cli TTL app:v1:user:42\nredis-cli MEMORY USAGE app:v1:user:42\n```\n\n若 `TYPE` 与客户端预期不符,例如代码执行 `HGETALL` 却返回 String 类型,Redis 会返回 WRONGTYPE 错误。这通常意味着键名冲突、版本迁移遗漏或序列化模型改变。\n\n### 13.2 检查 Streams 积压\n\n```bash\nredis-cli XLEN app:v1:stream:orders\nredis-cli XINFO GROUPS app:v1:stream:orders\nredis-cli XPENDING app:v1:stream:orders orders-group\n```\n\n重点观察:\n\n- Stream 长度是否持续增长;\n- `lag` 是否增加;\n- Pending 是否持续增加;\n- 是否有消费者长时间无心跳;\n- 是否存在单个消费者占据大量 Pending。\n\n清理 Stream 需要结合保留策略。盲目执行:\n\n```bash\nredis-cli DEL app:v1:stream:orders\n```\n\n会删除整个 Stream 及消费组,属于破坏性操作。生产清理应先确认备份、消费进度、恢复方案,再使用长度或时间窗口保留策略。\n\n### 13.3 监控指标\n\n应用和 Redis 两侧都应监控:\n\n- 命令延迟分位数;\n- 连接数和连接失败;\n- 超时和重试次数;\n- 缓存命中率;\n- 热点键访问;\n- 内存使用、淘汰数量;\n- Stream 长度、Pending 数、最老 Pending 空闲时间;\n- Lua 脚本耗时;\n- 主从复制延迟和故障切换。\n\n只监控 Redis CPU 而不监控缓存回源量,无法判断缓存故障是否已经传导到数据库。\n\n---\n\n## 14. 常见错误及其根因\n\n### 错误一:每个请求创建一个 RedisClient\n\n根因是把客户端对象误当作短生命周期连接。结果是连接建立、线程和资源开销增加,故障时还可能产生连接风暴。\n\n改为应用启动时创建客户端,关闭时统一释放;连接是否共享则根据事务、订阅和阻塞命令决定。\n\n### 错误二:把对象直接 `toString()` 存入 Redis\n\n`toString()` 通常不是稳定序列化协议,字段变化后无法可靠反序列化,也可能包含敏感信息。应使用明确的 JSON、二进制协议或 Hash,并定义版本兼容规则。\n\n### 错误三:锁释放直接 `DEL`\n\n这会让过期后的旧持有者删除新持有者的锁。必须使用随机令牌和比较令牌的 Lua 脚本。\n\n### 错误四:收到 Stream 消息就立即 ACK\n\n如果先 ACK 再处理,进程在业务处理前崩溃,消息已经从 PEL 中移除,无法自动重试。通常应先完成幂等业务处理,再 ACK。\n\n### 错误五:把 Stream 当作无限日志\n\nStream 会持续占用内存。必须定义保留策略、归档策略、死信处理和积压告警。ACK 不会自动删除历史消息。\n\n### 错误六:超时后无条件重试写命令\n\n响应丢失不等于服务端未执行。重试 `XADD`、`INCR` 等命令可能造成重复结果。写入协议必须设计幂等性或可确认性。\n\n### 错误七:用 Redis 锁代替数据库约束\n\n锁可能因租约过期、进程暂停、故障切换而失效。最终业务正确性仍应由数据库条件更新、唯一约束、版本号和幂等记录等机制共同保证。\n\n---\n\n## 15. 选择 Redis、数据库和 Streams 的边界\n\nRedis 适合低延迟访问、临时状态、计数、集合运算、短期缓存、协调信息和事件流转。关系数据库更适合需要持久约束、复杂查询、多表事务和审计的核心业务数据。\n\n可以用以下因果关系判断:\n\n- 数据丢失后可以重建:适合缓存;\n- 数据丢失后必须审计恢复:需要持久化系统和明确备份;\n- 需要严格唯一性:使用数据库唯一约束或等价的强约束;\n- 需要至少一次异步处理:Streams 消费组可行,但必须幂等;\n- 需要跨数据库事务发布事件:使用 Outbox 等事务消息模式;\n- 需要全局强一致锁:先定义故障模型,再评估 Redis 锁是否满足,不应直接套用代码。\n\nLettuce 只负责把 Java 程序接入 Redis。真正可靠的 Redis 工程,依赖的是清晰的键和值协议、正确的连接生命周期、明确的原子边界、可验证的租约和所有权、可恢复的消息消费流程,以及对网络和故障的显式建模。\n\n---\n\n## 系列导航与关联阅读\n\n- 系列入口:[Java 完整学习路线:从 Java 25 语言与 JVM 到 Spring、微服务和生产交付](https://wrblog.cn/articles/ee791cb5-6de0-5903-b22b-047cf231e387)\n- 上一篇:[Spring 事务深入:传播、隔离、代理、自调用和事件边界](https://wrblog.cn/articles/0eb2c40b-6262-5be3-9f6c-2330212cb5cb)\n- 下一篇:[Java Elasticsearch 工程:客户端、Mapping、查询、Bulk 和重试](https://wrblog.cn/articles/b9238834-5abb-5e2f-8eec-21ec380318d0)\n\n## 官方资料\n\n- [JDBC Basics](https://docs.oracle.com/javase/tutorial/jdbc/basics/)\n- [Jakarta Persistence Specification](https://jakarta.ee/specifications/persistence/)\n\n> 本文依据 Java、Spring 与相关项目官方文档重新梳理;正文、示例与生产清单由 WR BLOG 编写。\n","tags":["Java","Java 25 LTS","Redis"],"likeCount":0,"commentCount":0,"createdByUserId":"10000000000","createdByDisplayName":"小郝","createdByAvatar":"/public/profile/10000000000/avatar/2026/08/04/db02b81c-42f2-441b-8a80-61370cdbb581.webp","publishTime":"2026-09-01 13:12:21","updateTime":"2026-09-01 13:12:21"}},"status":200,"locale":"zh-CN","theme":"light"}