Java 基础体系 · 第 87/100 篇。示例统一以 Java 25 LTS 为语言和 JVM 基线;框架示例使用与其兼容的现代稳定版本。
Java Redis 工程:Lettuce、连接、序列化、缓存、锁和 Streams
Redis 工程问题通常不是“会不会调用 GET 和 SET”,而是要同时处理六件事:
- Lettuce 客户端如何建立连接、复用连接并关闭资源;
- Redis 字节如何映射为 Java 类型;
- 缓存数据如何写入、过期、重建和失效;
- 多线程、多实例下的锁是否真的具有互斥和安全释放能力;
- Redis Streams 如何投递、确认、重试和恢复消息;
- 网络断开、超时、主从切换、进程崩溃时,业务状态如何变化。
本文以 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();
}
}
}
这里有三个不同生命周期:
RedisClient:应用级客户端,通常创建一次;StatefulRedisConnection<K,V>:一个有状态的 Redis 连接;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 秒过期时间。相比先 SET 再 EXPIRE,单条 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 租约时间不是任务时间
如果任务最长执行时间为 ,锁租约为 ,还存在网络延迟、调度暂停和 GC 暂停 ,要避免锁在任务完成前过期,至少需要满足:
但 往往不是严格上界,因此实际工程中还需要:
- 自动续租;
- 任务分段并校验所有权;
- 让业务提交带 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
);
处理事件时:
- 尝试插入
event_id; - 若唯一键冲突,说明已经处理过,直接 ACK;
- 插入成功后执行业务更新;
- 业务成功后 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 客户端至少要区分:
- 连接建立超时:TCP、TLS、认证建立多久;
- 命令响应超时:发送命令后等待响应多久;
- 业务截止时间:整个请求允许 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 always、everysec、no在性能和持久性之间有不同取舍。
缓存数据通常可以接受丢失;订单事件、任务状态则需要更严格地评估持久化和恢复策略。
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
关键路径有三个不同的一致性边界:
- 缓存读取是性能优化边界,缓存丢失后可以回源;
- 锁是并发协调边界,需要令牌、租约和安全释放;
- Stream 是事件投递边界,需要 ACK、幂等、重试和积压监控。
如果“更新数据库”和“发布 Stream 事件”必须严格同时成功,直接先后执行仍可能发生:
数据库提交成功
Java 进程崩溃
XADD 尚未执行
更可靠的做法是数据库 Outbox:
- 在同一个数据库事务中更新业务表并写入 Outbox 表;
- 独立发布器读取 Outbox;
- 成功
XADD后标记 Outbox 已发布; - 发布器重复运行也必须幂等。
这体现了一个重要边界: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 不会自动删除历史消息。
错误六:超时后无条件重试写命令
响应丢失不等于服务端未执行。重试 XADD、INCR 等命令可能造成重复结果。写入协议必须设计幂等性或可确认性。
错误七:用 Redis 锁代替数据库约束
锁可能因租约过期、进程暂停、故障切换而失效。最终业务正确性仍应由数据库条件更新、唯一约束、版本号和幂等记录等机制共同保证。
15. 选择 Redis、数据库和 Streams 的边界
Redis 适合低延迟访问、临时状态、计数、集合运算、短期缓存、协调信息和事件流转。关系数据库更适合需要持久约束、复杂查询、多表事务和审计的核心业务数据。
可以用以下因果关系判断:
- 数据丢失后可以重建:适合缓存;
- 数据丢失后必须审计恢复:需要持久化系统和明确备份;
- 需要严格唯一性:使用数据库唯一约束或等价的强约束;
- 需要至少一次异步处理:Streams 消费组可行,但必须幂等;
- 需要跨数据库事务发布事件:使用 Outbox 等事务消息模式;
- 需要全局强一致锁:先定义故障模型,再评估 Redis 锁是否满足,不应直接套用代码。
Lettuce 只负责把 Java 程序接入 Redis。真正可靠的 Redis 工程,依赖的是清晰的键和值协议、正确的连接生命周期、明确的原子边界、可验证的租约和所有权、可恢复的消息消费流程,以及对网络和故障的显式建模。
系列导航与关联阅读
- 系列入口:Java 完整学习路线:从 Java 25 语言与 JVM 到 Spring、微服务和生产交付
- 上一篇:Spring 事务深入:传播、隔离、代理、自调用和事件边界
- 下一篇:Java Elasticsearch 工程:客户端、Mapping、查询、Bulk 和重试
官方资料
本文依据 Java、Spring 与相关项目官方文档重新梳理;正文、示例与生产清单由 WR BLOG 编写。

评论
0 条讨论