第1章:Redis 与 Java 集成概述
1.1 什么是 Redis
| 概念名称 | 说明 | 注意事项 |
|---|---|---|
| Redis | Remote Dictionary Server,是一个开源的内存数据结构存储系统,可用作数据库、缓存和消息中间件。支持字符串、哈希、列表、集合、有序集合等数据类型。 | 数据主要存储在内存中,因此读写性能极高,但需注意持久化策略以防数据丢失。 |
| 内存存储 | Redis 将数据保存在内存中,实现高速读写。 | 内存成本较高,不适合存储超大规模数据;可通过持久化机制将数据写入磁盘。 |
| 单线程模型 | Redis 服务端采用单线程处理命令(除后台线程外),通过非阻塞 I/O 和多路复用实现高并发。 | 避免执行耗时命令(如 KEYS *),否则会阻塞其他请求。 |
| 持久化机制 | 提供 RDB(快照)和 AOF(日志追加)两种持久化方式,保证重启后数据可恢复。 | 可根据业务需求选择合适的持久化策略,或结合使用。 |
| 支持多种数据结构 | 原生支持 String、Hash、List、Set、ZSet、Bitmap、HyperLogLog、Geo 等。 | 不同数据结构适用于不同场景,合理选择可提升性能与可维护性。 |
1.2 Redis 在 Java 中的应用场景
| 应用场景 | 说明 | 注意事项 |
|---|---|---|
| 高速缓存 | 将数据库查询结果缓存到 Redis,减少数据库压力,提升响应速度。 | 注意设置合理的过期时间,避免缓存堆积;防范缓存穿透、击穿、雪崩。 |
| 分布式会话共享 | 在集群部署中,将用户 Session 存储于 Redis,实现多节点共享。 | 需配合 Web 容器(如 Spring Session)实现透明化管理。 |
| 分布式锁 | 利用 SETNX 或 Redlock 算法实现跨 JVM 的互斥锁。 | 注意锁的超时释放与可重入性,防止死锁。 |
| 计数器 | 利用原子操作(如 INCR)实现访问统计、点赞数、限流计数等。 | 适合高并发读写场景,性能远优于数据库。 |
| 消息队列 | 使用 List 或 Pub/Sub 实现轻量级消息发布与订阅。 | Pub/Sub 为”发后即忘”模式,不保证消息可靠投递;List 更适合可靠队列。 |
| 排行榜 | 使用 ZSet 实现按分数排序的排行榜功能(如积分榜、热搜榜)。 | 支持范围查询、排名查询、分数更新,操作高效。 |
| 限流器 | 基于计数器或滑动窗口算法,控制接口调用频率。 | 可结合 Lua 脚本保证原子性,避免超限。 |
1.3 Java 操作 Redis 的主流客户端对比(Jedis vs Lettuce)
| 对比项 | Jedis | Lettuce | 注意事项 |
|---|---|---|---|
| 连接方式 | 基于 Socket 直连,每个线程需独立连接(推荐使用连接池) | 基于 Netty 的 NIO 框架,支持异步、响应式编程 | Lettuce 的连接是线程安全的,适合高并发场景 |
| 线程安全性 | Jedis 实例不线程安全,多线程需使用 JedisPool | Lettuce 的连接(StatefulRedisConnection)是线程安全的 | 使用 Jedis 时务必通过连接池获取实例 |
| 功能支持 | 支持基本命令、事务、Pipeline、集群等 | 支持同步、异步、响应式 API,支持 Redis Streams、集群拓扑动态更新 | Lettuce 更适合现代响应式编程模型(如 Spring WebFlux) |
| 依赖大小 | 轻量级,依赖少 | 依赖 Netty,包体积较大 | 若项目已用 Netty,Lettuce 是更好选择 |
| 集群支持 | 支持 Redis Cluster,但拓扑更新需手动刷新 | 自动监听集群拓扑变化,支持自适应重定向 | Lettuce 在集群环境下更稳定 |
| 社区与维护 | 早期主流,目前维护较慢 | Spring Data Redis 默认客户端,活跃维护中 | 推荐新项目优先使用 Lettuce |
1.4 开发环境准备与依赖引入
| 项目 | 说明 | 示例配置/代码 | 注意事项 |
|---|---|---|---|
| Redis 服务 | 需本地或远程运行 Redis 服务 | 下载地址:https://redis.io/download 启动命令: redis-server | 确保端口(默认 6379)可访问,关闭防火墙或配置安全组 |
| Java 环境 | JDK 8 或以上版本 | java -version | Redis 客户端多基于 JDK 8+ 开发 |
| Maven 依赖(Jedis) | 添加 Jedis 客户端依赖 | <dependency><groupId>redis.clients</groupId><artifactId>jedis</artifactId><version>4.3.1</version></dependency> | 版本建议使用 4.x 系列 |
| Maven 依赖(Lettuce) | 添加 Lettuce 客户端依赖 | <dependency><groupId>io.lettuce.core</groupId><artifactId>lettuce-core</artifactId><version>6.2.3.RELEASE</version></dependency> | 需注意与 Netty 版本兼容 |
| 构建工具 | 推荐使用 Maven 或 Gradle 管理依赖 | 使用 IDE(如 IntelliJ IDEA)导入项目 | 确保依赖下载成功,无冲突 |
第2章:连接 Redis 服务器
2.1 使用 Jedis 连接单机 Redis
| 方法名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| 构造函数 | new Jedis(String host, int port) | 创建连接到指定 Redis 服务器的 Jedis 实例 | Jedis jedis = new Jedis("localhost", 6379); | 必须手动关闭连接(jedis.close()),否则资源泄漏 |
| auth | jedis.auth(String password) | 认证密码(若 Redis 配置了 requirepass) | jedis.auth("123456"); | 必须在连接后立即调用 |
| connect | jedis.connect() | 显式建立连接(构造函数已自动连接) | jedis.connect(); | 一般无需手动调用 |
| ping | jedis.ping() | 测试连接是否正常 | String response = jedis.ping();System.out.println(response); // 返回 PONG | 可用于健康检查 |
| close | jedis.close() | 关闭连接,释放资源 | jedis.close(); | 每次使用后必须调用,建议用 try-with-resources |
2.2 使用连接池(JedisPool)管理连接
| 方法/类 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| JedisPool 构造 | new JedisPool(JedisPoolConfig config, String host, int port) | 创建连接池实例 | JedisPoolConfig config = new JedisPoolConfig();config.setMaxTotal(20);JedisPool pool = new JedisPool(config, "localhost", 6379); | 建议将 pool 设为单例 |
| getResource | pool.getResource() | 从池中获取一个 Jedis 连接 | Jedis jedis = pool.getResource(); | 获取的 jedis 实例使用后必须 close(),否则连接不会归还 |
| close | jedis.close() | 关闭连接(实际归还给池) | jedis.close(); | 必须调用,close() 会将连接返回池中 |
| close | pool.close() | 关闭整个连接池,释放所有资源 | pool.close(); | 应用关闭时调用一次即可 |
| JedisPoolConfig | 提供连接池配置参数 | 配置最大连接数、空闲数、超时等 | config.setMaxIdle(10);config.setMinIdle(5);config.setMaxWaitMillis(3000); | 合理设置参数避免资源浪费或连接不足 |
2.3 使用 Lettuce 连接 Redis(同步/异步)
| 方法/接口 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| RedisClient.create | RedisClient.create(String uri) | 创建 Redis 客户端实例 | RedisClient client = RedisClient.create("redis://localhost:6379"); | URI 格式:redis://host:port 或 redis://password@host:port |
| connect | client.connect() | 获取同步连接 | RedisConnection<String, String> connection = client.connect();SyncCommands<String, String> sync = connection.sync(); | sync 接口用于同步操作 |
| connectAsync | client.connectAsync(RedisCodec) | 获取异步连接 | RedisAsyncCommands<String, String> async = client.connect().async(); | 返回 CompletableFuture,支持非阻塞调用 |
| set/get | sync.set(String key, String value)sync.get(String key) | 同步执行读写命令 | sync.set("name", "Alice");String value = sync.get("name"); | 同步阻塞当前线程 |
| setAsync/getAsync | async.set(K key, V value)async.get(K key) | 异步执行命令 | CompletableFuture future = async.get("name");future.thenAccept(System.out::println); | 非阻塞,适合高并发 |
| close | connection.close()client.shutdown() | 关闭连接和客户端 | connection.close();client.shutdown(); | 应用结束时关闭客户端 |
2.4 连接 Redis 集群与哨兵模式
| 连接模式 | 语法/类 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| Redis 集群(Jedis) | new JedisCluster(Set<HostAndPort> hosts) | 连接 Redis Cluster | Set nodes = new HashSet<>();nodes.add(new HostAndPort("127.0.0.1", 7000));JedisCluster jc = new JedisCluster(nodes); | 不支持多 key 操作(除非在同个 slot);需关闭时调用 jc.close() |
| Redis 哨兵(Jedis) | new JedisSentinelPool(String masterName, Set<String> sentinels) | 通过哨兵发现主节点并连接 | Set sentinels = new HashSet<>();sentinels.add("127.0.0.1:26379");JedisSentinelPool pool = new JedisSentinelPool("mymaster", sentinels); | 自动切换主从,适合高可用场景 |
| Redis 集群(Lettuce) | RedisClusterClient.create(RedisURI... uris) | 连接 Redis 集群 | RedisURI uri = RedisURI.create("redis://127.0.0.1:7000");RedisClusterClient client = RedisClusterClient.create(uri);StatefulRedisClusterConnection<String, String> conn = client.connect(); | 自动发现集群拓扑,支持同步与异步操作 |
| Redis 哨兵(Lettuce) | RedisClient.create(RedisURI) 其中 RedisURI 指向哨兵 | 通过哨兵连接主节点 | RedisURI uri = RedisURI.Builder.sentinel("127.0.0.1", 26379, "mymaster").withPassword("123").build();RedisClient client = RedisClient.create(uri); | Lettuce 自动监听主节点变更,无需额外配置 |
| 跨 slot 操作限制 | - | Redis Cluster 中多 key 操作必须在同个 hash slot | 使用 {} 包裹 key 前缀,如 user:{1}:name 和 user:{1}:age 会被分配到同 slot | 否则会抛出 CROSSSLOT 错误 |
第3章:Redis 核心数据类型与 Java 操作(基础篇)
3.1 String 类型:基本读写与数值操作
| 方法名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| set | jedis.set(String key, String value) | 设置字符串值 | jedis.set("name", "Alice"); | 若键已存在则覆盖原值 |
| get | jedis.get(String key) | 获取字符串值 | String name = jedis.get("name"); | 键不存在返回 null |
| setex | jedis.setex(String key, int seconds, String value) | 设置带过期时间的字符串 | jedis.setex("token", 3600, "abc123"); | 单位为秒,常用于缓存 |
| setnx | jedis.setnx(String key, String value) | 仅当键不存在时设置(分布式锁基础) | long result = jedis.setnx("lock", "1"); | 返回 1 表示成功,0 表示已存在 |
| incr | jedis.incr(String key) | 将字符串值视为整数并加 1 | jedis.set("count", "10");jedis.incr("count"); // → 11 | 值必须为整数格式字符串 |
| decr | jedis.decr(String key) | 将字符串值减 1 | jedis.decr("count"); // → 10 | 同上,支持负数 |
| incrBy | jedis.incrBy(String key, long increment) | 指定增量进行增加 | jedis.incrBy("count", 5); // +5 | 支持任意长整型增量 |
| decrBy | jedis.decrBy(String key, long decrement) | 指定减量进行减少 | jedis.decrBy("count", 3); // -3 | 同上 |
| append | jedis.append(String key, String value) | 在原有字符串末尾追加内容 | jedis.append("msg", " World"); | 若键不存在等价于 set |
| strlen | jedis.strlen(String key) | 获取字符串长度 | long len = jedis.strlen("msg"); | 中文字符按 UTF-8 字节计算 |
3.2 Hash 类型:字段映射操作
| 方法名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| hset | jedis.hset(String key, String field, String value) | 设置哈希字段的值 | jedis.hset("user:1", "name", "Bob"); | 返回 1 表示新增,0 表示更新 |
| hget | jedis.hget(String key, String field) | 获取指定字段的值 | String name = jedis.hget("user:1", "name"); | 字段不存在返回 null |
| hmset | jedis.hmset(String key, Map<String,String> hash) | 批量设置多个字段 | Map<String,String> user = new HashMap<>();user.put("name","Tom");user.put("age","25");jedis.hmset("user:1", user); | 原子性操作,推荐批量写入 |
| hmget | jedis.hmget(String key, String... fields) | 批量获取多个字段值 | List values = jedis.hmget("user:1", "name", "age"); | 返回顺序与传入字段一致 |
| hgetAll | jedis.hgetAll(String key) | 获取所有字段及其值 | Map<String,String> map = jedis.hgetAll("user:1"); | 返回整个 Hash 结构 |
| hdel | jedis.hdel(String key, String... fields) | 删除一个或多个字段 | jedis.hdel("user:1", "age"); | 返回被删除字段的数量 |
| hexists | jedis.hexists(String key, String field) | 判断字段是否存在 | boolean exists = jedis.hexists("user:1", "name"); | 返回 true/false |
| hkeys | jedis.hkeys(String key) | 获取所有字段名 | Set fields = jedis.hkeys("user:1"); | 不包含值 |
| hvals | jedis.hvals(String key) | 获取所有字段值 | List vals = jedis.hvals("user:1"); | 不包含字段名 |
| hincrBy | jedis.hincrBy(String key, String field, long increment) | 对数字字段执行原子增减 | jedis.hincrBy("user:1", "score", 10); | 字段值必须为整数 |
3.3 List 类型:列表结构操作
| 方法名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| lpush | jedis.lpush(String key, String... strings) | 从左侧插入一个或多个元素 | jedis.lpush("tasks", "task1", "task2"); | 插入后元素在列表头部 |
| rpush | jedis.rpush(String key, String... strings) | 从右侧插入一个或多个元素 | jedis.rpush("queue", "msg1", "msg2"); | 常用于实现队列 |
| lpop | jedis.lpop(String key) | 从左侧弹出一个元素 | String task = jedis.lpop("tasks"); | 元素被移除,线程安全 |
| rpop | jedis.rpop(String key) | 从右侧弹出一个元素 | String msg = jedis.rpop("queue"); | 可用于实现栈或队列 |
| lrange | jedis.lrange(String key, long start, long end) | 获取指定范围内的元素 | List list = jedis.lrange("tasks", 0, -1); | 支持负索引(-1 表末尾) |
| llen | jedis.llen(String key) | 获取列表长度 | long size = jedis.llen("tasks"); | 空列表返回 0 |
| lindex | jedis.lindex(String key, long index) | 获取指定索引位置的元素 | String first = jedis.lindex("tasks", 0); | 越界返回 null |
| lset | jedis.lset(String key, long index, String value) | 设置指定索引位置的值 | jedis.lset("tasks", 0, "newTask"); | 索引必须存在,否则报错 |
| lrem | jedis.lrem(String key, long count, String value) | 删除指定数量的元素 | jedis.lrem("tasks", 1, "task1"); | count > 0 从头删,< 0 从尾删,= 0 删除全部 |
| ltrim | jedis.ltrim(String key, long start, long end) | 保留指定范围内的元素,其余删除 | jedis.ltrim("tasks", 0, 9); | 常用于限制列表大小 |
3.4 Set 类型:无序集合操作
| 方法名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| sadd | jedis.sadd(String key, String... members) | 添加一个或多个成员 | jedis.sadd("tags", "java", "redis"); | 已存在成员不会重复添加 |
| smembers | jedis.smembers(String key) | 获取集合中所有成员 | Set tags = jedis.smembers("tags"); | 无序返回 |
| sismember | jedis.sismember(String key, String member) | 判断成员是否存在于集合中 | boolean hasTag = jedis.sismember("tags", "java"); | 返回 true/false |
| srem | jedis.srem(String key, String... members) | 删除一个或多个成员 | jedis.srem("tags", "redis"); | 返回实际删除的数量 |
| spop | jedis.spop(String key) | 随机弹出并移除一个成员 | String tag = jedis.spop("tags"); | 集合为空返回 null |
| srandmember | jedis.srandmember(String key)jedis.srandmember(String key, int count) | 随机获取一个或多个成员(不删除) | String rand = jedis.srandmember("tags");List rands = jedis.srandmember("tags", 2); | 可用于抽奖等场景 |
| scard | jedis.scard(String key) | 获取集合成员数量 | long count = jedis.scard("tags"); | 类似 size() |
| sinter | jedis.sinter(String... keys) | 计算多个集合的交集 | Set inter = jedis.sinter("tags1", "tags2"); | 返回新集合 |
| sunion | jedis.sunion(String... keys) | 计算多个集合的并集 | Set union = jedis.sunion("tags1", "tags2"); | 返回所有唯一成员 |
| sdiff | jedis.sdiff(String... keys) | 计算第一个集合与其他集合的差集 | Set diff = jedis.sdiff("tags1", "tags2"); | 属于 tags1 但不属于 tags2 的成员 |
3.5 ZSet(Sorted Set)类型:有序集合操作
| 方法名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| zadd | jedis.zadd(String key, double score, String member)jedis.zadd(String key, Map<String,Double> members) | 添加带分数的成员 | jedis.zadd("scores", 95.5, "Alice");Map<String,Double> m = new HashMap<>();m.put("Bob", 87.0);jedis.zadd("scores", m); | 分数决定排序顺序 |
| zrange | jedis.zrange(String key, long start, long end) | 按分数升序获取成员(不含分数) | Set top3 = jedis.zrange("scores", 0, 2); | 支持负索引 |
| zrevrange | jedis.zrevrange(String key, long start, long end) | 按分数降序获取成员 | Set worst = jedis.zrevrange("scores", -1, -3); | 用于排行榜倒序 |
| zrangeWithScores | jedis.zrangeWithScores(String key, long start, long end) | 获取成员及其分数(升序) | Set tuples = jedis.zrangeWithScores("scores", 0, -1); | Tuple 包含 value 和 score |
| zrevrangeWithScores | jedis.zrevrangeWithScores(String key, long start, long end) | 获取成员及其分数(降序) | Set desc = jedis.zrevrangeWithScores("scores", 0, 4); | 排行榜常用 |
| zscore | jedis.zscore(String key, String member) | 获取指定成员的分数 | Double score = jedis.zscore("scores", "Alice"); | 成员不存在返回 null |
| zrank | jedis.zrank(String key, String member) | 获取成员按分数升序的排名(从 0 开始) | Long rank = jedis.zrank("scores", "Alice"); | 不存在返回 null |
| zrevrank | jedis.zrevrank(String key, String member) | 获取成员按分数降序的排名 | Long revRank = jedis.zrevrank("scores", "Alice"); | 常用于”您当前排名第 X” |
| zcard | jedis.zcard(String key) | 获取有序集合的成员数量 | long total = jedis.zcard("scores"); | 类似 size() |
| zrem | jedis.zrem(String key, String... members) | 删除一个或多个成员 | jedis.zrem("scores", "Charlie"); | 返回删除成功的数量 |
| zincrby | jedis.zincrby(String key, double increment, String member) | 对成员分数执行增量操作 | jedis.zincrby("scores", 5.0, "Alice"); | 原子性操作,适合实时更新 |
| zcount | jedis.zcount(String key, double min, double max) | 统计指定分数范围内成员数量 | long highScorers = jedis.zcount("scores", 90, 100); | 支持区间查询 |
第4章:Redis 高级功能与 Java 实现
4.1 键的过期与生存时间(TTL)操作
| 方法名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| expire | jedis.expire(String key, int seconds) | 设置键的过期时间(秒) | jedis.expire("token", 1800); | 过期后自动删除 |
| pexpire | jedis.pexpire(String key, long milliseconds) | 设置键的过期时间(毫秒) | jedis.pexpire("temp", 5000); | 更精细控制 |
| expireAt | jedis.expireAt(String key, long unixTime) | 在指定 Unix 时间戳过期 | jedis.expireAt("data", System.currentTimeMillis()/1000 + 60); | 单位为秒 |
| pexpireAt | jedis.pexpireAt(String key, long millisecondsTimestamp) | 在指定毫秒时间戳过期 | jedis.pexpireAt("log", System.currentTimeMillis() + 10000); | 精确到毫秒 |
| ttl | jedis.ttl(String key) | 获取键的剩余生存时间(秒) | long ttl = jedis.ttl("token"); | -2 表示键不存在,-1 表示永不过期 |
| pttl | jedis.pttl(String key) | 获取剩余生存时间(毫秒) | long pttl = jedis.pttl("token"); | 同上,单位不同 |
| persist | jedis.persist(String key) | 移除键的过期时间,使其永久有效 | jedis.persist("data"); | 若原本无过期时间则无影响 |
4.2 事务操作(MULTI/EXEC)
| 方法名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| multi | jedis.multi() | 开启事务,后续命令进入队列 | Transaction tx = jedis.multi();tx.set("a", "1");tx.incr("counter"); | 命令不会立即执行 |
| exec | tx.exec() | 提交事务,执行所有队列中的命令 | List<Object> results = tx.exec(); | 返回每条命令的执行结果列表 |
| discard | tx.discard() | 取消事务,清空命令队列 | tx.discard(); | 用于放弃事务 |
| watch | jedis.watch(String... keys) | 监视一个或多个键,在事务执行前若被修改则事务失败 | jedis.watch("balance");Transaction tx = jedis.multi();tx.decrBy("balance", 100);tx.exec(); | 实现乐观锁机制 |
| unwatch | jedis.unwatch() | 取消对所有键的监视 | jedis.unwatch(); | 通常在事务失败后调用 |
注意:Redis 事务不支持回滚(Rollback),exec 失败时已执行的命令无法撤销。watch 是实现条件更新的关键。
4.3 Lua 脚本执行(EVAL/EVALSHA)
| 方法名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| eval | jedis.eval(String script, int keyCount, String... params) | 执行 Lua 脚本 | String script = "return redis.call('GET', KEYS[1])";String result = (String) jedis.eval(script, 1, "name"); | 支持复杂逻辑原子执行 |
| evalsha | jedis.evalsha(String sha1, int keyCount, String... params) | 使用脚本 SHA1 哈希执行缓存脚本 | String sha = jedis.scriptLoad(script);jedis.evalsha(sha, 1, "name"); | 提升性能,避免重复传输脚本 |
| scriptLoad | jedis.scriptLoad(String script) | 预加载脚本并返回其 SHA1 值 | String sha = jedis.scriptLoad(script); | 用于后续 evalsha 调用 |
| exists | jedis.scriptExists(String... sha1s) | 检查脚本是否已缓存 | boolean[] exists = jedis.scriptExists(sha); | 可批量检查 |
| flush | jedis.scriptFlush() | 清空服务器端所有 Lua 脚本缓存 | jedis.scriptFlush(); | 一般无需手动调用 |
原子性实现限流器(令牌桶)示例:
String script = "local tokens = redis.call('GET', KEYS[1]) " +
"if tokens < tonumber(ARGV[1]) then " +
" return 0 " +
"else " +
" redis.call('DECRBY', KEYS[1], ARGV[1]) " +
" return 1 " +
"end";
Object result = jedis.eval(script, 1, "tokens", "1");
4.4 发布/订阅模式(Pub/Sub)
| 方法/类 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| subscribe | jedis.subscribe(JedisPubSub listener, String... channels) | 订阅一个或多个频道 | jedis.subscribe(new JedisPubSub(){...}, "news"); | 阻塞当前线程,需单独线程运行 |
| publish | jedis.publish(String channel, String message) | 向频道发布消息 | long receivers = jedis.publish("news", "Hello!"); | 返回收到消息的订阅者数量 |
| JedisPubSub | 抽象类 | 用户需继承并重写回调方法 | class MyListener extends JedisPubSub {public void onMessage(String ch, String msg) { ... }} | 必须实现 onMessage 等方法 |
| onMessage | public void onMessage(String channel, String message) | 收到消息时回调 | @Overridepublic void onMessage(String ch, String msg) {System.out.println(ch + ": " + msg);} | 核心处理逻辑在此 |
| unsubscribe | listener.unsubscribe() | 取消订阅 | listener.unsubscribe("news"); | 可指定频道或全部取消 |
注意:Jedis 的 Pub/Sub 是阻塞式的,一旦 subscribe 调用,该连接不能再执行其他命令,建议使用独立连接或 Lettuce 的异步模型。
4.5 Pipeline 批量操作提升性能
| 方法名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| pipelined | jedis.pipelined() | 获取 Pipeline 实例,开启管道模式 | Pipeline pipe = jedis.pipelined();pipe.set("a", "1");pipe.get("b");List<Object> results = pipe.syncAndReturnAll(); | 所有命令暂存本地,一次发送 |
| sync | pipe.sync() | 同步等待所有命令执行完成 | pipe.sync(); | 不返回结果,仅等待 |
| syncAndReturnAll | pipe.syncAndReturnAll() | 同步执行并返回所有命令结果 | List<Object> results = pipe.syncAndReturnAll(); | 结果顺序与发送顺序一致 |
| close | pipe.close() | 关闭管道连接 | try (Pipeline p = jedis.pipelined()) { ... } | 推荐使用 try-with-resources |
批量插入示例:
Pipeline pipe = jedis.pipelined();
for (int i = 1; i <= 10000; i++) {
pipe.set("key:" + i, "value" + i);
}
pipe.sync(); // 发送所有命令
性能对比:普通操作耗时数秒,Pipeline 可降至几十毫秒。
第5章:Redis 持久化与高可用在 Java 中的应用
5.1 RDB 与 AOF 原理简介
| 概念 | 说明 | 注意事项 |
|---|---|---|
| RDB(Redis Database) | 快照持久化方式,在指定时间间隔内生成数据集的二进制快照(dump.rdb)。通过 fork 子进程进行,主进程继续处理请求。 | 适合备份和灾难恢复;恢复速度快;可能丢失最后一次快照后的数据。 |
| 触发条件(RDB) | 配置 save <seconds> <changes>,如 save 900 1 表示 900 秒内至少有 1 次修改则触发。手动执行 SAVE 或 BGSAVE 命令。 | SAVE 阻塞主进程,生产环境应使用 BGSAVE。 |
| AOF(Append Only File) | 记录每一条写命令到日志文件(appendonly.aof),重启时重放命令恢复数据。支持三种同步策略:no、everysec、always。 | 数据更安全,最多丢失 1 秒数据;文件通常比 RDB 大;恢复速度较慢。 |
| AOF 重写(Rewrite) | 定期重写 AOF 文件以压缩体积(如将多次 set 合并为最终值)。不会读取旧 AOF 文件,而是从内存中读取当前状态重新生成。 | 通过 BGREWRITEAOF 命令或自动触发(auto-aof-rewrite-percentage)。 |
| 混合持久化(3.2+) | 开启 aof-use-rdb-preamble yes 后,AOF 文件前半部分为 RDB 格式,后半部分为增量命令。 | 兼顾恢复速度与数据安全性,推荐开启。 |
5.2 Java 客户端对持久化状态的监控
| 方法名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| info | jedis.info() | 获取 Redis 服务器运行信息 | String info = jedis.info("persistence");System.out.println(info); | 可获取 RDB/AOF 状态、上次保存时间、错误信息等 |
| getLastSave | jedis.lastsave() | 获取最后一次成功保存 RDB 的 Unix 时间戳 | long lastSaveTime = jedis.lastsave();Date date = new Date(lastSaveTime * 1000); | 单位为秒,可用于判断是否正常持久化 |
| bgsave | jedis.bgsave() | 异步触发 RDB 快照保存 | String status = jedis.bgsave(); // 返回 "Background saving started" | 非阻塞,推荐用于定时备份 |
| bgrewriteaof | jedis.bgrewriteaof() | 异步触发 AOF 重写 | jedis.bgrewriteaof(); | 减少 AOF 文件大小,提升性能 |
| isSaving | jedis.isbgsaving() | 判断是否正在执行 BGSAVE | boolean saving = jedis.isbgsaving(); | 结合 lastsave 可实现备份进度监控 |
定期检查持久化状态示例:
String persistenceInfo = jedis.info("persistence");
if (persistenceInfo.contains("rdb_last_bgsave_status:ok")) {
System.out.println("RDB 持久化正常");
} else {
System.err.println("RDB 持久化失败");
}
5.3 哨兵模式下的自动故障转移处理
| 概念/方法 | 说明 | 代码示例 | 注意事项 |
|---|---|---|---|
| 哨兵(Sentinel)角色 | 监控主从节点健康状态,主节点宕机时自动选举新主节点并通知客户端。 | - | 至少部署 3 个哨兵实例以避免脑裂 |
| JedisSentinelPool | 支持通过哨兵自动发现主节点的连接池 | Set sentinels = new HashSet<>(Arrays.asList("192.168.1.101:26379"));JedisSentinelPool pool = new JedisSentinelPool("mymaster", sentinels);Jedis jedis = pool.getResource(); | 构造时传入 masterName 和哨兵地址集合 |
| 故障转移流程 | 主节点宕机 → 哨兵投票 → 选举从节点为新主 → 更新配置 → 通知客户端重新连接 | - | 客户端需支持自动重连机制 |
| onSwitchMaster | 重写 Sentinel 连接监听器回调 | pool.getSentinels().iterator().next().getHostAndPort(); | JedisSentinelPool 内部自动处理,无需手动干预 |
| 超时设置 | 建议设置合理的 connectionTimeout 和 soTimeout | new JedisSentinelPool("mymaster", sentinels, new JedisPoolConfig(), 2000, 2000); | 避免因网络抖动导致频繁切换 |
注意:JedisSentinelPool 在主从切换后会自动更新内部主节点地址,但已获取的 Jedis 实例不会自动刷新,建议每次使用后
close()归还连接池。
5.4 Redis 集群模式下的分片与重定向处理
| 概念/方法 | 说明 | 代码示例 | 注意事项 |
|---|---|---|---|
| 数据分片(Sharding) | Redis Cluster 将键空间划分为 16384 个槽(slot),每个节点负责一部分槽。key 的 slot = CRC16(key) % 16384 | 客户端可直接计算目标节点 | - |
| MOVED 重定向 | 当请求的 key 不属于当前节点时,返回 MOVED <slot> <ip:port>,客户端需重试到正确节点 | JedisCluster 自动处理 MOVED 重定向 | 应用无需关心具体节点位置 |
| ASK 重定向 | 在集群重新分片期间,临时指向目标节点的提示 | JedisCluster 也自动处理 ASK | 属于临时状态 |
| JedisCluster | 支持集群模式的客户端封装,内置重定向、重试机制 | Set nodes = new HashSet<>();nodes.add(new HostAndPort("127.0.0.1", 7000));JedisCluster jc = new JedisCluster(nodes); | 所有操作基于该实例 |
| 多 key 操作限制 | 只有所有 key 的 slot 相同时才能执行 multi-key 命令(如 mget、union) | 使用 {user1000} 包裹 key 实现同 slot 分配 | 否则抛出 CROSSSLOT 错误 |
| clusterInfo / clusterNodes | 获取集群状态信息 | String info = jc.clusterInfo();Set nodes = jc.clusterNodes().keySet(); | 用于监控和诊断 |
确保多 key 操作在同一 slot:
// 使用 {} 包裹共同部分,Redis 只对 {} 内的内容计算 slot
jc.set("{session}:uid:123", "alice");
jc.set("{session}:token:123", "abc");
// 上述两个 key 的 slot 相同,可安全执行 mget
List<String> result = jc.mget("{session}:uid:123", "{session}:token:123");
第6章:Spring 与 Redis 集成
6.1 Spring Data Redis 简介与配置
| 项目 | 说明 | 配置示例 | 注意事项 |
|---|---|---|---|
| Spring Data Redis | Spring 提供的 Redis 抽象层,统一操作接口,支持多种客户端(Jedis/Lettuce) | Maven 依赖:<dependency><groupId>org.springframework.data</groupId><artifactId>spring-data-redis</artifactId></dependency> | 推荐搭配 Lettuce 使用 |
| RedisConnectionFactory | 连接工厂接口,创建与 Redis 的连接 | 对于 Jedis:JedisConnectionFactory对于 Lettuce: LettuceConnectionFactory | Spring Boot 自动配置 |
| XML 配置方式 | 使用 XML 定义 Bean | <bean id="connectionFactory" class="org.springframework.data.redis.connection.jedis.JedisConnectionFactory"/> | 已较少使用 |
| Java Config 方式 | 使用 @Configuration 类配置 | @Configuration@EnableRedisRepositoriespublic class RedisConfig {@Beanpublic LettuceConnectionFactory connectionFactory() {return new LettuceConnectionFactory(new RedisStandaloneConfiguration("localhost", 6379));}} | 推荐方式 |
| Spring Boot 自动装配 | 添加 spring-boot-starter-data-redis 后自动配置 | application.yml:spring:redis:host: localhostport: 6379 | 简化配置,开箱即用 |
6.2 使用 RedisTemplate 操作数据
| 方法名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| opsForValue | redisTemplate.opsForValue() | 获取 ValueOperations,操作 String 类型 | ValueOperations<String, Object> ops = redisTemplate.opsForValue();ops.set("name", "Alice");String name = (String) ops.get("name"); | 默认序列化使用 JDK,可能导致乱码 |
| opsForHash | redisTemplate.opsForHash() | 获取 HashOperations | HashOperations<String, String, Object> hOps = redisTemplate.opsForHash();hOps.put("user:1", "age", 25); | 支持泛型 |
| opsForList | redisTemplate.opsForList() | 获取 ListOperations | ListOperations<String, Object> lOps = redisTemplate.opsForList();lOps.leftPush("queue", "task1"); | - |
| opsForSet | redisTemplate.opsForSet() | 获取 SetOperations | SetOperations<String, Object> sOps = redisTemplate.opsForSet();sOps.add("tags", "java", "redis"); | - |
| opsForZSet | redisTemplate.opsForZSet() | 获取 ZSetOperations | ZSetOperations<String, Object> zOps = redisTemplate.opsForZSet();zOps.add("scores", "Alice", 95); | - |
| execute | redisTemplate.execute(RedisCallback<T> action) | 执行底层原生命令 | Long size = redisTemplate.execute((RedisCallback) con -> con.dbSize()); | 高级用法,绕过模板封装 |
注意:默认情况下
RedisTemplate使用JdkSerializationRedisSerializer,存储的 key/value 会出现不可读的二进制格式,建议自定义序列化器。
6.3 使用 StringRedisTemplate 处理字符串
| 特性 | 说明 | 代码示例 | 注意事项 |
|---|---|---|---|
| 默认序列化 | 使用 StringRedisSerializer,key 和 value 均为 UTF-8 字符串 | StringRedisTemplate template = new StringRedisTemplate(connectionFactory);template.opsForValue().set("msg", "Hello Redis");String msg = template.opsForValue().get("msg"); | 存储内容可读性强,适合纯字符串场景 |
| 与 RedisTemplate 区别 | 继承自 RedisTemplate,但泛型固定为 <String, String>,且序列化器为字符串类型 | - | 推荐用于缓存简单字符串、JSON 文本等 |
| JSON 存储场景 | 常用于存储序列化后的 JSON 字符串 | User user = new User("Bob", 30);template.opsForValue().set("user:1", objectMapper.writeValueAsString(user)); | 需配合 ObjectMapper 手动序列化 |
6.4 自定义序列化策略(JSON、JDK、String)
| 序列化器类 | 说明 | 配置示例 | 注意事项 |
|---|---|---|---|
| JdkSerializationRedisSerializer | 默认序列化器,使用 Java 原生序列化 | redisTemplate.setValueSerializer(new JdkSerializationRedisSerializer()); | 需对象实现 Serializable;结果不可读;跨语言不兼容 |
| StringRedisSerializer | 将对象转换为 UTF-8 字符串 | redisTemplate.setKeySerializer(new StringRedisSerializer());redisTemplate.setValueSerializer(new StringRedisSerializer()); | 仅适用于字符串或已序列化的文本(如 JSON) |
| GenericJackson2JsonRedisSerializer | 使用 Jackson 将对象序列化为 JSON 字符串,并保留类型信息 | redisTemplate.setValueSerializer(new GenericJackson2JsonRedisSerializer()); | 支持复杂对象,存储可读,推荐用于 POJO 缓存 |
| Jackson2JsonRedisSerializer | Jackson 序列化器,需指定类型 | redisTemplate.setValueSerializer(new Jackson2JsonRedisSerializer<>(User.class)); | 类型固定,灵活性较低 |
| 配置完整示例 | 设置 key 和 value 的序列化方式 | redisTemplate.setKeySerializer(new StringRedisSerializer());redisTemplate.setValueSerializer(new GenericJackson2JsonRedisSerializer());redisTemplate.setHashKeySerializer(new StringRedisSerializer());redisTemplate.setHashValueSerializer(new GenericJackson2JsonRedisSerializer()); | 建议统一配置 hash 的序列化方式 |
推荐配置(Java Config):
@Bean
public RedisTemplate<String, Object> redisTemplate(RedisConnectionFactory factory) {
RedisTemplate<String, Object> template = new RedisTemplate<>();
template.setConnectionFactory(factory);
template.setKeySerializer(new StringRedisSerializer());
template.setValueSerializer(new GenericJackson2JsonRedisSerializer());
template.setHashKeySerializer(new StringRedisSerializer());
template.setHashValueSerializer(new GenericJackson2JsonRedisSerializer());
template.afterPropertiesSet();
return template;
}
6.5 缓存注解:@Cacheable、@CachePut、@CacheEvict
| 注解 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| @Cacheable | @Cacheable(value="users", key="#id") | 缓存方法返回值,下次相同参数直接从缓存读取 | @Cacheable("users")public User findById(Long id) {return userRepository.findById(id);} | 默认使用方法参数生成 key |
| @CachePut | @CachePut(value="users", key="#user.id") | 执行方法并更新缓存,无论是否存在 | @CachePut("users")public User update(User user) {return userRepository.save(user);} | 常用于更新操作,保证缓存最新 |
| @CacheEvict | @CacheEvict(value="users", key="#id")@CacheEvict(value="users", allEntries=true) | 清除指定或全部缓存条目 | @CacheEvict("users")public void deleteById(Long id) {userRepository.deleteById(id);} | allEntries=true 清空整个缓存区 |
| @Caching | @Caching(evict = {...}, put = {...}) | 组合多个缓存操作 | @Caching(evict = @CacheEvict("users"),put = @CachePut("logs")) | 复杂场景使用 |
| SpEL 表达式 | 支持在 key、condition 中使用 SpEL | key="#username"condition="#age > 18"unless="#result == null" | 动态控制缓存行为 |
启用缓存注解:
@Configuration
@EnableCaching // 开启缓存支持
public class CacheConfig {
@Bean
public CacheManager cacheManager(RedisConnectionFactory factory) {
RedisCacheConfiguration config = RedisCacheConfiguration.defaultCacheConfig()
.entryTtl(Duration.ofMinutes(10)); // 设置默认过期时间
return RedisCacheManager.builder(factory).cacheDefaults(config).build();
}
}
第7章:Redis 实际应用场景与 Java 实现
7.1 分布式锁的实现(SETNX + Lua)
| 方法/机制 | 说明 | 代码示例 | 注意事项 |
|---|---|---|---|
| SETNX + EXPIRE | 使用 setnx 设置锁,expire 设置超时防止死锁 | jedis.setnx("lock:order", "1");jedis.expire("lock:order", 10); | 非原子操作,存在竞态条件 |
| SET 带 NX EX 参数 | 原子性设置锁并设置过期时间 | jedis.set("lock:order", "client_1", "NX", "EX", 10); | 推荐方式,避免上述问题 |
| Lua 脚本释放锁 | 使用脚本保证”判断-删除”原子性 | String script = "if redis.call('get', KEYS[1]) == ARGV[1] then return redis.call('del', KEYS[1]) else return 0 end";jedis.eval(script, 1, "lock:order", "client_1"); | 防止误删其他客户端的锁 |
| 锁续期(Watchdog) | 使用后台线程定期延长锁过期时间 | ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(1);scheduler.scheduleAtFixedRate(() -> {jedis.expire("lock:order", 10);}, 5, 5, TimeUnit.SECONDS); | 实现类似 Redisson 的看门狗机制 |
| Redlock 算法 | 多实例部署下提高锁可靠性 | 使用 Redisson 的 RLock 实现 Redlock | 单节点存在脑裂风险,生产环境建议多节点 |
推荐实现(原子性加锁 + Lua 释放):
// 加锁
String result = jedis.set(lockKey, clientId, "NX", "EX", 30);
if ("OK".equals(result)) {
// 获取锁成功
}
// 释放锁(Lua)
String script = "if redis.call('get', KEYS[1]) == ARGV[1] then return redis.call('del', KEYS[1]) else return 0 end";
jedis.eval(script, 1, lockKey, clientId);
7.2 限流器(Rate Limiter)设计
| 实现方式 | 说明 | 代码示例 | 注意事项 |
|---|---|---|---|
| 固定窗口计数器 | 使用 INCR 实现单位时间内的访问计数 | String key = "rate_limit:" + userId;jedis.incr(key);jedis.expire(key, 60); // 60秒内最多100次 | 存在临界问题(短时间内翻倍) |
| 滑动窗口(ZSet) | 利用 ZSet 存储时间戳,计算窗口内请求数 | long now = System.currentTimeMillis();jedis.zadd("sliding:win", now, String.valueOf(now));jedis.zremrangeByScore("sliding:win", 0, now - 60000);Long count = jedis.zcard("sliding:win"); | 精确控制,适合高并发场景 |
| 令牌桶算法(Lua) | Redis 中维护令牌生成与消费逻辑 | String script = "local tokens = redis.call('get', KEYS[1]); if not tokens then tokens = tonumber(ARGV[1]) end; if tokens >= tonumber(ARGV[2]) then redis.call('set', KEYS[1], tokens - tonumber(ARGV[2])); return 1; else return 0; end";jedis.eval(script, 1, "tokens", "100", "1"); | 原子性操作,支持突发流量 |
| Redisson RRateLimiter | 使用 Redisson 封装的限流器 | RRateLimiter rateLimiter = redisson.getRateLimiter("myRateLimiter");rateLimiter.trySetRate(RateType.OVERALL, 10, 1, RateIntervalUnit.SECONDS);boolean canPass = rateLimiter.tryAcquire(); | 开箱即用,推荐生产环境使用 |
注意:避免在 Lua 脚本中执行耗时操作,影响 Redis 性能。
7.3 缓存穿透、击穿、雪崩的应对策略
| 问题 | 原因 | 解决方案 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| 缓存穿透 | 查询不存在的数据,绕过缓存直接打到 DB | 布隆过滤器拦截非法请求;缓存空值并设置短过期时间 | jedis.setex("user:999", 60, ""); // 缓存空结果 | 空值过期时间不宜过长,避免占用内存 |
| 缓存击穿 | 热点 key 过期瞬间大量请求涌入 DB | 设置永不过期或逻辑过期;使用互斥锁重建缓存 | synchronized (this) {if (jedis.get(key) == null) {User user = db.find(id);jedis.setex(key, 3600, serialize(user));}} | 分布式环境下应使用分布式锁 |
| 缓存雪崩 | 大量 key 同时过期或 Redis 宕机 | 过期时间加随机值(如 30±10 分钟);高可用部署(哨兵/集群);多级缓存(本地+Redis) | int expire = 1800 + new Random().nextInt(600);jedis.setex("data", expire, value); | 避免集中过期,提升系统韧性 |
7.4 用户会话(Session)共享存储
| 方案 | 说明 | 实现方式 | 注意事项 |
|---|---|---|---|
| Spring Session + Redis | 使用 Spring Session 替换默认 HttpSession | 1. 添加 spring-session-data-redis2. 配置 @EnableRedisHttpSession3. 自动将 Session 存入 Redis | 支持透明迁移,推荐方案 |
| 自定义 Session 存储 | 手动将 Session 数据存入 Redis | String sessionId = UUID.randomUUID().toString();Map<String, String> sessionData = new HashMap<>();sessionData.put("userId", "123");jedis.hmset("session:" + sessionId, sessionData);jedis.expire("session:" + sessionId, 1800); | 需处理过期、清理、安全性等问题 |
| Cookie + Redis | 将 Session ID 存于 Cookie,数据存于 Redis | response.addCookie(new Cookie("JSESSIONID", sessionId)); | 注意 Cookie 安全属性(HttpOnly、Secure) |
| TTL 设置 | 控制会话生命周期 | jedis.expire("session:" + id, 30 * 60); // 30分钟 | 可结合用户活动延长有效期 |
Spring Session 配置示例:
@Configuration
@EnableRedisHttpSession(maxInactiveIntervalInSeconds = 1800)
public class SessionConfig {
@Bean
public LettuceConnectionFactory connectionFactory() {
return new LettuceConnectionFactory();
}
}
7.5 排行榜与计分系统(ZSet 应用)
| 操作 | 方法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| 添加用户得分 | zadd | 初始化或更新用户分数 | jedis.zadd("leaderboard", 95.5, "Alice"); | 分数决定排名顺序 |
| 获取 Top N | zrevrangeWithScores | 获取排名前 N 的用户 | Set top10 = jedis.zrevrangeWithScores("leaderboard", 0, 9); | 降序排列,第一名在前 |
| 查询用户排名 | zrevrank | 获取用户按分数降序的排名 | Long rank = jedis.zrevrank("leaderboard", "Alice"); | 排名从 0 开始 |
| 查询用户分数 | zscore | 获取用户当前得分 | Double score = jedis.zscore("leaderboard", "Alice"); | 用于展示详情 |
| 分数更新 | zincrby | 原子性增加用户分数 | jedis.zincrby("leaderboard", 5.0, "Alice"); | 支持负数(扣分) |
| 分页查询 | zrevrange | 实现分页排行榜 | Set page = jedis.zrevrange("leaderboard", (page-1)*size, page*size-1); | 注意边界处理 |
| 删除用户 | zrem | 从排行榜移除用户 | jedis.zrem("leaderboard", "Bob"); | 用户不再参与排名 |
场景扩展:带条件的排行榜(如按部门):
// 使用 key 隔离不同维度
jedis.zadd("leaderboard:sales", 88.0, "Tom");
jedis.zadd("leaderboard:tech", 92.0, "Jerry");
第8章:性能优化与最佳实践
8.1 连接池参数调优(JedisPool)
| 参数 | 说明 | 推荐值 | 注意事项 |
|---|---|---|---|
| maxTotal | 最大连接数 | 200~500 | 根据并发量调整,避免过多连接耗尽资源 |
| maxIdle | 最大空闲连接数 | 50~100 | 避免频繁创建/销毁连接 |
| minIdle | 最小空闲连接数 | 20~50 | 保证一定数量的热连接 |
| maxWaitMillis | 获取连接最大等待时间(毫秒) | 5000~10000 | 超时抛出异常,避免线程阻塞 |
| testOnBorrow | 借出时验证连接有效性 | true | 保证连接可用,但增加开销 |
| testOnReturn | 归还时验证连接 | false | 一般设为 false 提升性能 |
| blockWhenExhausted | 池耗尽时是否阻塞 | true | 配合 maxWaitMillis 使用 |
| timeBetweenEvictionRunsMillis | 空闲连接检测周期 | 30000(30秒) | 定期清理无效连接 |
配置示例:
JedisPoolConfig config = new JedisPoolConfig();
config.setMaxTotal(300);
config.setMaxIdle(100);
config.setMinIdle(50);
config.setMaxWaitMillis(5000);
config.setTestOnBorrow(true);
config.setBlockWhenExhausted(true);
JedisPool pool = new JedisPool(config, "localhost", 6379);
8.2 批量操作与 Pipeline 使用建议
| 场景 | 推荐方式 | 说明 | 注意事项 |
|---|---|---|---|
| 多 key 读写 | 使用 Pipeline | 减少网络往返延迟 | 单次 Pipeline 不宜过大(建议 < 1000 条) |
| 大数据量导入 | 分批次 Pipeline | 避免 OOM 和超时 | 每批 100~500 条 |
| 原子性要求高 | Lua 脚本 | 在服务端执行复杂逻辑 | 脚本执行时间不宜过长 |
| 读多写少 | mget / mset | 原生支持批量操作 | key 应尽量分布在同节点(集群模式) |
| 异步处理 | Lettuce + CompletableFuture | 非阻塞 I/O,更高吞吐 | 适合高并发场景 |
Pipeline 批量插入优化:
List<String> data = getData(); // 10000条
int batchSize = 500;
for (int i = 0; i < data.size(); i += batchSize) {
try (Pipeline p = jedis.pipelined()) {
for (int j = i; j < i + batchSize && j < data.size(); j++) {
p.set("key:" + j, data.get(j));
}
p.sync(); // 提交一批
}
}
8.3 避免大 Key 与热 Key 的设计
| 问题 | 风险 | 解决方案 | 检测方法 |
|---|---|---|---|
| 大 Key(Big Key) | 阻塞 Redis、网络超时、内存溢出 | 拆分结构(如 Hash 分片)、压缩存储、使用外部存储 | redis-cli --bigkeysMEMORY USAGE key |
| 热 Key(Hot Key) | 单节点负载过高、CPU 瓶颈 | 本地缓存(Caffeine)+ Redis、读写分离、Key 拆分 | 监控命令耗时、使用 Redis 监视器 |
| 大 Value | 网络传输慢、序列化开销大 | 存储摘要或引用,数据存于 DB/OSS | 使用 strlen 检查字符串长度 |
| 集合过大 | ZSet/Hash 包含百万级元素 | 分页存储、按时间分片(如按天) | zcard / hlen 统计大小 |
热 Key 本地缓存示例:
LoadingCache<String, String> localCache = Caffeine.newBuilder()
.maximumSize(1000)
.expireAfterWrite(10, TimeUnit.MINUTES)
.build(key -> jedis.get(key));
String value = localCache.get("hot:config");
8.4 监控 Redis 性能指标(INFO 命令与客户端监控)
| 指标类别 | 关键指标 | 获取方式 | 告警阈值 |
|---|---|---|---|
| 内存 | used_memory, used_memory_rss, mem_fragmentation_ratio | jedis.info("memory") | 碎片率 > 1.5 需关注 |
| CPU | used_cpu_sys, used_cpu_user | jedis.info("cpu") | 持续 > 80% 可能存在问题 |
| 命令统计 | total_commands_processed, instantaneous_ops_per_sec | jedis.info("stats") | 突增可能为异常流量 |
| 连接 | connected_clients, maxclients | jedis.info("clients") | 接近 maxclients 需扩容 |
| 持久化 | rdb_last_bgsave_status, aof_enabled | jedis.info("persistence") | RDB 失败需立即处理 |
| 集群 | cluster_state, keys_slot | jedis.clusterInfo() | fail 状态表示异常 |
| 延迟 | latency, slowlog | redis-cli --latencyjedis.slowlogGet() | 慢查询日志需分析优化 |
建议:定期采集 INFO 信息并接入监控系统(如 Prometheus + Grafana)。
8.5 异常处理与重试机制
| 异常类型 | 常见原因 | 处理策略 | 代码示例 |
|---|---|---|---|
| JedisConnectionException | 连接超时、断开 | 重试 + 熔断(如 Hystrix) | try { ... } catch (JedisConnectionException e) { retry(); } |
| JedisDataException | 命令语法错误、类型错误 | 日志记录 + 修复逻辑 | 检查 key 类型是否匹配操作 |
| TimeoutException | 响应超时 | 缩短超时时间、降级处理 | 设置 soTimeout ≤ 500ms |
| RedisException | 一般 Redis 错误 | 统一捕获并分类处理 | 结合监控告警 |
| 重试机制 | 网络抖动、临时故障 | 指数退避重试(Exponential Backoff) | Thread.sleep(100 * (1 << retryCount)) |
| 降级策略 | Redis 不可用 | 降级到 DB、返回默认值、限流 | 缓存不可用时直接查库 |
简单重试逻辑示例:
int retries = 0;
while (retries < 3) {
try {
return jedis.get("key");
} catch (JedisConnectionException e) {
retries++;
if (retries >= 3) throw e;
try {
Thread.sleep(100 * (1 << retries));
} catch (InterruptedException ie) {
Thread.currentThread().interrupt();
}
}
}
第9章:常见问题与调试技巧
9.1 连接超时与拒绝连接的排查
| 问题现象 | 可能原因 | 排查步骤 | 解决方案 | 注意事项 |
|---|---|---|---|---|
java.net.SocketTimeoutException: Read timed out | 网络延迟高、Redis 负载大、客户端超时设置过短 | 1. 检查网络延迟(ping、telnet) 2. 查看 Redis INFO stats 中 total_commands_processed 和 instantaneous_ops_per_sec 3. 使用 redis-cli --latency 测试延迟 | - 增加 soTimeout - 优化查询(避免大 key) - 升级 Redis 实例配置 | 超时时间建议设置为 500~2000ms,避免过长导致线程阻塞 |
java.net.ConnectException: Connection refused | Redis 未启动、端口错误、防火墙拦截、绑定 IP 限制 | 1. 确认 Redis 是否运行:ps aux | grep redis2. 检查 redis.conf 中 port 和 bind 配置 3. 使用 telnet 测试连通性 4. 查看防火墙规则(iptables/firewalld) | - 启动 Redis 服务 - 修改 bind 0.0.0.0 允许远程访问(生产环境谨慎) - 开放防火墙端口 | |
JedisConnectionException: Could not get a resource from the pool | 连接池耗尽、连接未正确归还 | 1. 检查 maxTotal 是否过小 2. 确保每次使用后调用 jedis.close() 或 pool.returnResource()3. 监控 connected_clients 数量 | - 增加 maxTotal - 使用 try-with-resources - 启用 testOnReturn | 推荐使用连接池并确保资源释放 |
| SSL/TLS 连接失败 | Redis 启用 TLS,但客户端未配置 | 检查 Redis 是否启用 tls-port,客户端是否使用 SSLSocketFactory | 使用 Lettuce 并配置 SSL 选项:LettuceClientConfiguration.builder().useSsl().build(); | Jedis 不支持原生 SSL,建议使用 Lettuce |
最佳实践:
try (Jedis jedis = pool.getResource()) {
return jedis.get("key");
} // 自动归还连接,避免泄露
9.2 序列化异常与乱码问题
| 问题现象 | 可能原因 | 排查方法 | 解决方案 | 注意事项 |
|---|---|---|---|---|
存储后出现 \xac\xed\x00\x05t\x00\x04key 等二进制前缀 | 使用默认 JdkSerializationRedisSerializer 序列化 | 使用 redis-cli 查看 key 值是否包含不可读字符 | 自定义序列化器:redisTemplate.setValueSerializer(new StringRedisSerializer());或使用 JSON 序列化 | JDK 序列化仅适用于实现 Serializable 的类 |
| SerializationException:无法反序列化对象 | 类结构变更、类未实现 Serializable、序列化器不匹配 | 检查类是否修改了 serialVersionUID 确认读写使用相同序列化策略 | 统一使用 GenericJackson2JsonRedisSerializer固定 serialVersionUID | 避免频繁修改实体类结构 |
| 中文乱码(如 æ\x9f³è\x8a¹) | 编码不一致(非 UTF-8) | 检查序列化器是否使用 UTF-8 | 显式指定字符集:new StringRedisSerializer(StandardCharsets.UTF_8) | 确保客户端、Redis、序列化器编码一致 |
| JSON 反序列化失败(类型丢失) | 使用 Jackson2JsonRedisSerializer 时未指定类型 | 反序列化为 LinkedHashMap 而非原类型 | 改用 GenericJackson2JsonRedisSerializer或注册子类型 | GenericJackson2JsonRedisSerializer 会写入 @class 字段保留类型信息 |
| Hash 结构 key/value 乱码 | 未设置 hashKeySerializer 和 hashValueSerializer | 使用 HGETALL 查看字段是否可读 | 统一设置 hash 序列化策略:template.setHashKeySerializer(new StringRedisSerializer());template.setHashValueSerializer(new GenericJackson2JsonRedisSerializer()); | RedisTemplate 需单独设置 hash 的序列化器 |
推荐配置(避免乱码):
@Bean
public RedisTemplate<String, Object> redisTemplate(RedisConnectionFactory factory) {
RedisTemplate<String, Object> template = new RedisTemplate<>();
template.setConnectionFactory(factory);
StringRedisSerializer stringSerializer = new StringRedisSerializer();
GenericJackson2JsonRedisSerializer jsonSerializer = new GenericJackson2JsonRedisSerializer();
template.setKeySerializer(stringSerializer);
template.setValueSerializer(jsonSerializer);
template.setHashKeySerializer(stringSerializer);
template.setHashValueSerializer(jsonSerializer);
template.afterPropertiesSet();
return template;
}
9.3 集群环境下 MOVED/ASK 重定向处理
| 问题现象 | 可能原因 | 排查方法 | 解决方案 | 注意事项 |
|---|---|---|---|---|
MOVED 1234 192.168.1.100:7001 | 请求的 key 所在 slot 不属于当前节点 | 使用 CLUSTER KEYSLOT <key> 计算 slot使用 CLUSTER NODES 查看 slot 分配 | 使用 JedisCluster 或 Lettuce 集群客户端,自动处理重定向 | 原生 Jedis 不支持自动重定向 |
ASK 1234 192.168.1.100:7001 | 集群正在迁移 slot,临时重定向 | 查看 redis-cli --cluster check 是否有迁移任务 | 客户端应先发送 ASKING 命令再执行请求 | JedisCluster 自动处理 ASK 流程 |
多 key 操作报错 CROSSSLOT Keys in request don't hash to the same slot | 多个 key 的 slot 不一致 | 使用 CRC128(key) % 16384 计算各 key 的 slot | 使用 {} 包裹 key 实现哈希标签(Hash Tag):jedis.mget("{user}:1", "{user}:2") | {} 内的内容用于计算 slot,外部不影响 |
| 客户端连接失败,提示节点不可达 | 集群拓扑变更、节点宕机 | 使用 redis-cli -c -h <node> -p <port> 测试连接执行 CLUSTER INFO 查看 cluster_state | 确保客户端能访问所有主从节点 配置正确的集群节点列表 | JedisCluster 初始化时传入部分节点即可,会自动发现其他节点 |
使用 Hash Tag 确保同槽:
// 以下两个 key 的 slot 相同,可执行 mget
String key1 = "{session}:user:123";
String key2 = "{session}:token:123";
List<String> values = jedisCluster.mget(key1, key2);
9.4 内存溢出与 Redis 内存分析
| 问题现象 | 可能原因 | 分析工具/命令 | 解决方案 | 注意事项 |
|---|---|---|---|---|
| Redis 进程占用内存持续增长,接近 maxmemory | 大 key、未设置 TTL、缓存堆积 | INFO memory:- used_memory - mem_fragmentation_ratio redis-cli --bigkeysMEMORY USAGE <key> | - 设置合理的 maxmemory-policy(如 allkeys-lru) - 清理无用 key - 为缓存 key 设置 TTL | 避免使用 del 删除大 key,改用 UNLINK(异步删除) |
| Java 客户端 OOM | 加载大 value 到内存(如 100MB 的 String) | JVM 堆分析(jmap + MAT) 监控 jedis.get() 返回对象大小 | - 分页读取大集合 - 使用 HSCAN / ZSCAN - 增加 JVM 堆内存 | 单个 key 建议不超过 10KB,最大不超过 1MB |
| 内存碎片率高(> 1.5) | 频繁增删 key 导致内存碎片 | INFO memory 中 mem_fragmentation_ratio | - 重启 Redis(释放碎片) - 启用 activedefrag(Redis 4.0+) | 碎片率 = used_memory_rss / used_memory |
| 持久化时内存翻倍 | RDB bgsave 或 AOF rewrite 时 fork 子进程 | INFO persistence系统内存监控 | - 确保物理内存充足 - 避免在高峰时段触发持久化 - 使用 copy-on-write 优化 | fork 失败会导致持久化失败 |
| 集群节点内存不均 | 数据分布不均、热 key 集中 | redis-cli --cluster info <node>统计各节点 keyspace | - 使用 Hash Tag 均匀分布 - 拆分大 key - 重新分片(rebalance) | 定期使用 --cluster rebalance 均衡负载 |
内存分析命令示例:
# 查找大 key
redis-cli --bigkeys
# 查看某个 key 的内存占用
redis-cli MEMORY USAGE "user:profile:1000"
# 查看内存碎片率
redis-cli INFO memory | grep mem_fragmentation_ratio
配置建议(redis.conf):
# 设置最大内存
maxmemory 4gb
# 内存满时使用 LRU 回收
maxmemory-policy allkeys-lru
# 启用主动碎片整理
activedefrag yes
第10章:Redis 安全与权限控制
| 问题/功能 | 风险或目标 | 配置方法 | Java 实现建议 | 注意事项 |
|---|---|---|---|---|
| 未授权访问 | Redis 默认无密码,可被任意连接读写 | 在 redis.conf 中设置:requirepass yourpassword或运行时执行 CONFIG SET requirepass "yourpassword" | Jedis/Lettuce 连接时传入密码:jedis.auth("password");new RedisStandaloneConfiguration("host", 6379).setPassword(RedisPassword.of("password")); | 密码应使用强密码策略,避免硬编码 |
| 绑定 IP 限制 | 允许外部网络扫描和攻击 | 修改 bind 指令:bind 127.0.0.1 192.168.1.100仅监听内网或本地接口 | 客户端配置对应 IP 地址 | 生产环境禁止 bind 0.0.0.0(除非有防火墙保护) |
| 禁用危险命令 | 如 FLUSHALL、KEYS *、SHUTDOWN 可能被滥用 | 使用 rename-command 重命名或禁用:rename-command FLUSHALL ""rename-command CONFIG "config_cmd" | 应用层避免调用高危命令;运维侧统一管理 | 建议在哨兵/集群模式下统一配置 |
| ACL(Access Control List)(Redis 6+) | 细粒度用户权限控制(推荐方式) | 配置 users.acl 文件或使用命令:ACL SETUSER alice on >secret ~cached:* +get +set +hmgetACL SETUSER monitor on >ro ~* +info +client | Lettuce 支持 ACL 用户登录:StatefulRedisConnection<String, String> connection = client.connect(RedisURI.create("redis://alice:secret@localhost:6379")); | 每个应用使用独立账号,遵循最小权限原则 |
| TLS 加密通信 | 数据传输明文,存在窃听风险 | 启用 TLS:tls-port 6379tls-cert-file /path/to/cert.pemtls-key-file /path/to/key.pem | 使用 Lettuce 并启用 SSL:LettuceConnectionFactory factory = new LettuceClientConfigurationBuilder().useSsl().build(); | Jedis 不原生支持 TLS,建议升级到 Lettuce |
| 防火墙规则 | 开放不必要的端口 | 使用 iptables/firewalld 限制访问:iptables -A INPUT -p tcp --dport 6379 -s 192.168.1.0/24 -j ACCEPT | - | 结合 bind 和防火墙实现双重防护 |
推荐安全配置模板(redis.conf):
bind 192.168.1.100
protected-mode yes
port 6379
tcp-backlog 511
timeout 300
tcp-keepalive 300
requirepass StrongPassw0rd!2025
rename-command FLUSHALL ""
rename-command FLUSHDB ""
rename-command DEBUG ""
aclfile /etc/redis/users.acl
# TLS 配置(可选)
# tls-port 6380
# tls-cert-file /ssl/redis.crt
# tls-key-file /ssl/redis.key
第11章:Redis 6+ 新特性(ACL、IO 多线程)
| 特性 | 说明 | 配置与使用 | Java 应用影响 | 注意事项 |
|---|---|---|---|---|
| ACL(访问控制列表) | 替代旧版密码机制,支持多用户、命令权限、key 权限、密码哈希 | 创建用户:ACL SETUSER app_user on >P@ssw0rd ~cart:* +get +set +expire查看用户: ACL LIST / ACL WHOAMI | Spring Data Redis 支持通过 URI 指定用户名密码:redis://app_user:P@ssw0rd@localhost:6379 | 推荐用于微服务架构中不同模块的权限隔离 |
| IO 多线程(Threaded I/O) | 将网络读写操作并行化,提升高并发吞吐量(命令执行仍在主线程) | 启用多线程:io-threads-do-reads yesio-threads 4(建议 ≤ CPU 核心数) | 性能透明提升,无需修改客户端代码 | 禁止对磁盘操作(如 AOF)使用多线程;性能测试验证效果 |
| 客户端缓存(Client-side Caching) | 客户端可缓存数据,服务端通过 CLIENT TRACKING 通知失效 | 启用追踪:CLIENT TRACKING ON REDIRECT 10086(需另一个连接接收推送) | Lettuce 支持 tracking 模式:connection.setClientTracking(true, Consumer.of(redirectChannel)); | 适用于热点数据场景,减少网络往返 |
| RESP3 协议 | 新一代通信协议,支持更多数据类型(如 Map、Set、Stream) | 启用 RESP3:proto-max-bulk-len 512mb客户端协商升级 | Lettuce 支持 RESP3,Jedis 当前不支持 | 提升协议扩展性,未来主流方向 |
| Swappable Value Storage(实验性) | 将冷数据交换到磁盘,节省内存 | 使用 LFU 策略标记冷数据,配合外部存储 | - | 目前为实验功能,生产环境慎用 |
| Modules 扩展能力增强 | 更强大的模块 API(如 RedisJSON、RedisBloom) | 加载模块:loadmodule /usr/lib/redis/modules/rejson.so | Java 可通过 JSON 模块直接操作文档:jedis.sendCommand(Command.JSON_SET, "user:1", ".", "{\"name\":\"Alice\"}"); | 极大拓展 Redis 使用场景 |
IO 多线程性能优化建议:
# redis.conf
io-threads 4
io-threads-do-reads yes
# 禁用 write 多线程(可能降低性能)
# io-threads-do-writes no
第12章:综合案例
12.1 电商购物车系统设计
12.1.1 业务需求分析
| 功能 | 描述 |
|---|---|
| 添加商品 | 用户将商品加入购物车,支持数量调整 |
| 查看购物车 | 获取所有商品列表及总价 |
| 更新数量 | 修改某商品数量 |
| 删除商品 | 从购物车移除商品 |
| 批量操作 | 全选、清空、删除多个商品 |
| 跨设备同步 | 登录后购物车数据自动加载 |
| 过期清理 | 未登录用户购物车保留 30 天 |
12.1.2 技术选型与数据结构设计
| 数据项 | 存储结构 | Key 设计 | 说明 |
|---|---|---|---|
| 已登录用户购物车 | Hash | cart:user:{userId} | field=商品ID, value=数量(JSON 字符串) |
| 未登录用户购物车 | Hash | cart:guest:{sessionId} | 同上,session 过期即失效 |
| 商品信息缓存 | String(JSON) | product:{productId} | 缓存商品名称、价格等,减少 DB 查询 |
| 购物车总数量 | String | cart:count:{userId} | 用于首页展示小红点 |
| 热门商品 Top 10 | ZSet | cart:hot | score=添加次数,用于推荐 |
示例数据结构:
HSET cart:user:1001 "10086" "{\"count\":2,\"price\":599.0}"
HSET cart:user:1001 "20001" "{\"count\":1,\"price\":89.9}"
GET product:10086 → {"name":"iPhone Case","price":599.0}
12.1.3 Java 核心代码实现(Spring Boot + RedisTemplate)
@Service
public class ShoppingCartService {
@Autowired
private RedisTemplate<String, Object> redisTemplate;
@Value("${cart.expiration.minutes:43200}") // 30天
private long expirationMinutes;
private static final String CART_PREFIX = "cart:user:";
private static final String GUEST_CART_PREFIX = "cart:guest:";
private static final String PRODUCT_PREFIX = "product:";
// 添加商品
public void addItem(Long userId, Long productId, int count) {
String key = CART_PREFIX + userId;
String productKey = PRODUCT_PREFIX + productId;
// 缓存商品信息(异步)
cacheProductInfo(productId);
// 获取当前数量并更新
BoundHashOperations<String, String, String> ops = redisTemplate.boundHashOps(key);
String existing = (String) ops.get(productId.toString());
CartItem item = StringUtils.hasText(existing) ?
parseCartItem(existing) : new CartItem();
item.setCount(item.getCount() + count);
item.setProductId(productId);
item.setPrice(getProductPrice(productId)); // 应从缓存获取
ops.put(productId.toString(), toJson(item));
redisTemplate.expire(key, expirationMinutes, TimeUnit.MINUTES);
// 更新计数器
updateCartCount(userId, ops.size());
}
// 获取购物车列表
public List<CartVO> getCartItems(Long userId) {
String key = CART_PREFIX + userId;
BoundHashOperations<String, String, String> ops = redisTemplate.boundHashOps(key);
List<CartVO> result = new ArrayList<>();
double totalPrice = 0.0;
for (Map.Entry<String, String> entry : ops.entries().entrySet()) {
CartItem item = parseCartItem(entry.getValue());
CartVO vo = new CartVO();
vo.setProductId(Long.valueOf(entry.getKey()));
vo.setCount(item.getCount());
vo.setName(getProductName(vo.getProductId()));
vo.setPrice(item.getPrice());
vo.setSubtotal(item.getPrice() * item.getCount());
totalPrice += vo.getSubtotal();
result.add(vo);
}
return result;
}
// 删除商品
public void removeItem(Long userId, Long productId) {
String key = CART_PREFIX + userId;
redisTemplate.opsForHash().delete(key, productId.toString());
}
// 清空购物车
public void clearCart(Long userId) {
String key = CART_PREFIX + userId;
redisTemplate.delete(key);
}
// 更新数量
public void updateItemCount(Long userId, Long productId, int count) {
if (count <= 0) {
removeItem(userId, productId);
} else {
addItem(userId, productId, count - getCurrentCount(userId, productId));
}
}
// 私有辅助方法...
private double getProductPrice(Long productId) { /* 从缓存或 DB 获取 */ }
private String getProductName(Long productId) { /* 获取名称 */ }
private void cacheProductInfo(Long productId) { /* 异步缓存 */ }
private CartItem parseCartItem(String json) { /* JSON 反序列化 */ }
private String toJson(Object obj) { /* JSON 序列化 */ }
private void updateCartCount(Long userId, long size) { /* 更新计数器 */ }
private int getCurrentCount(Long userId, Long productId) { /* 获取当前数量 */ }
}
12.1.4 高可用与性能优化策略
| 优化点 | 实现方式 | 说明 |
|---|---|---|
| 本地缓存加速 | 使用 Caffeine 缓存热门商品信息 | 减少 Redis 访问次数 |
| 异步持久化 | 用户操作后异步同步到 DB(Kafka + 消费者) | 保证最终一致性 |
| 批量操作 Pipeline | 批量添加/删除使用 Pipeline | 减少网络开销 |
| 分页加载 | 超过 100 件商品分页展示 | 避免大 value 传输 |
| 缓存预热 | 启动时加载热销商品到 Redis | 提升首屏速度 |
| 监控告警 | 接入 Prometheus + Grafana 监控 cart:* 指标 | 及时发现异常 |
12.1.5 异常处理与边界情况
| 场景 | 处理方案 |
|---|---|
| 商品已下架 | 从购物车移除并提示用户 |
| 库存不足 | 提示”库存紧张”,限制购买数量 |
| 价格变动 | 显示”价格已更新”,以结算页为准 |
| 用户登录合并购物车 | 将 guest cart 合并到 user cart,去重累加 |
| 并发修改冲突 | 使用 Lua 脚本保证原子性更新 |
Lua 脚本示例(原子更新数量):
-- KEYS[1]: cart key, ARGV[1]: product id, ARGV[2]: delta
local current = redis.call('HGET', KEYS[1], ARGV[1])
local count = tonumber(current) or 0
local newCount = count + tonumber(ARGV[2])
if newCount <= 0 then
redis.call('HDEL', KEYS[1], ARGV[1])
else
redis.call('HSET', KEYS[1], ARGV[1], newCount)
end
return newCount
12.2 Redis 与消息队列实现订单超时取消
在电商系统中,用户下单后通常有 15分钟 的支付有效期。若超时未支付,系统需自动取消订单并释放库存。本节介绍基于 Redis + 消息队列 的高效、可靠实现方案。
12.2.1 业务需求分析
| 功能 | 描述 |
|---|---|
| 订单创建 | 用户下单后生成待支付订单,设置超时时间(如 15 分钟) |
| 超时检测 | 在指定时间后触发取消逻辑 |
| 支付成功 | 用户支付后立即取消超时任务 |
| 高可用 | 系统重启或节点宕机后,未完成任务不丢失 |
| 低延迟 | 超时后尽可能快地执行取消操作(秒级延迟) |
12.2.2 技术选型对比
| 方案 | 优点 | 缺点 | 适用场景 |
|---|---|---|---|
| 定时轮询 DB(如 Quartz) | 实现简单,数据一致 | 轮询压力大,延迟高(分钟级) | 小型系统,低并发 |
| Redis ZSet 延迟队列 | 延迟低(秒级),轻量 | 无持久化保障,宕机可能丢失任务 | 中小型系统,可接受少量丢失 |
| Redis Streams | 支持消费者组、持久化、ACK 机制 | Redis 5.0+ 才支持 | 推荐:中大型系统 |
| Kafka 延迟消息 | 高吞吐、高可靠 | 延迟控制复杂,需额外组件 | 超大型系统,已有 Kafka 架构 |
| RabbitMQ TTL + 死信队列 | 延迟精确,支持重试 | 运维复杂,性能低于 Redis | 已使用 RabbitMQ 的系统 |
推荐方案:Redis Streams(兼顾性能、可靠性和实现复杂度)
12.2.3 基于 Redis Streams 的实现方案
1. 数据结构设计
| 结构 | Key | 说明 |
|---|---|---|
| 订单延迟队列 | stream:order:delay | 使用 XADD 写入待取消订单,score 为超时时间戳(毫秒) |
| 消费者组 | group:order:cancel | 创建消费者组处理取消任务 |
| 订单状态缓存 | order:status:{orderId} | 存储订单当前状态(PAID / CANCELLED / UNPAID) |
| 已支付订单集合 | set:order:paid | 用于快速判断订单是否已支付,防止重复取消 |
2. 流程图解:
用户下单
↓
[服务A] 保存订单到DB
↓
[服务A] XADD stream:order:delay * orderId 1500000000000 // 15分钟后
↓
[服务B] XREADGROUP GROUP group:order:cancel consumer1 STREAMS stream:order:delay >
↓
检查 order:status:{orderId} 是否为 UNPAID
↓ 是
释放库存 + 更新订单状态为 CANCELLED
↓
XACK stream:order:delay group:order:cancel <msg-id>
12.2.4 Java 核心代码实现(Spring Boot + Lettuce + Redis Streams)
1. 添加依赖(Maven):
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-redis</artifactId>
</dependency>
<dependency>
<groupId>io.lettuce</groupId>
<artifactId>lettuce-core</artifactId>
</dependency>
2. 订单创建时写入延迟队列:
@Service
public class OrderService {
@Autowired
private StringRedisTemplate redisTemplate;
@Value("${order.timeout.minutes:15}")
private long timeoutMinutes;
// 创建订单
public void createOrder(Long orderId, Long userId, BigDecimal amount) {
// 1. 保存订单到数据库
saveOrderToDB(orderId, userId, amount);
// 2. 设置订单状态为待支付
String statusKey = "order:status:" + orderId;
redisTemplate.opsForValue().set(statusKey, "UNPAID", timeoutMinutes * 60, TimeUnit.SECONDS);
// 3. 写入 Redis Streams 延迟队列
String streamKey = "stream:order:delay";
long expireAt = System.currentTimeMillis() + timeoutMinutes * 60 * 1000; // 毫秒
Map<String, Object> message = new HashMap<>();
message.put("orderId", orderId.toString());
message.put("expireAt", String.valueOf(expireAt));
redisTemplate.opsForStream().add(
StreamRecords.string(message)
.withStreamKey(streamKey)
);
}
}
3. 启动消费者组监听超时订单:
@Component
public class OrderTimeoutConsumer {
@Autowired
private StringRedisTemplate redisTemplate;
private static final String STREAM_KEY = "stream:order:delay";
private static final String GROUP_NAME = "group:order:cancel";
private static final String CONSUMER_NAME = "consumer:timeout";
@PostConstruct
public void startConsumer() {
// 创建消费者组(如果不存在)
try {
redisTemplate.opsForStream().createGroup(STREAM_KEY, ReadOffset.from("0"), GROUP_NAME);
} catch (Exception e) {
// 组已存在,忽略
}
// 启动监听
StreamMessageListenerContainerOptions<String, ObjectRecord<String, String>> options =
StreamMessageListenerContainerOptions.builder()
.pollTimeout(Duration.ofSeconds(1))
.build();
StreamMessageListenerContainer<String, ObjectRecord<String, String>> container =
StreamMessageListenerContainer.create(redisTemplate.getConnectionFactory(), options);
Consumer consumer = Consumer.from(GROUP_NAME, CONSUMER_NAME);
StreamOffset<String> offset = StreamOffset.create(STREAM_KEY, ReadOffset.lastConsumed());
container.receive(consumer, offset, (record) -> {
String orderIdStr = record.getValue().get("orderId");
Long orderId = Long.valueOf(orderIdStr);
try {
// 检查订单是否已支付
String statusKey = "order:status:" + orderId;
String status = redisTemplate.opsForValue().get(statusKey);
if ("UNPAID".equals(status)) {
// 执行取消逻辑
cancelOrder(orderId);
// 标记为已处理
redisTemplate.opsForStream().acknowledge(STREAM_KEY, GROUP_NAME, record.getId());
}
} catch (Exception e) {
// 记录日志,消息将重新投递(最多3次)
log.error("处理订单取消失败,orderId: {}", orderId, e);
}
});
container.start();
}
private void cancelOrder(Long orderId) {
// 1. 释放库存
inventoryService.releaseStock(orderId);
// 2. 更新订单状态
orderMapper.updateStatus(orderId, "CANCELLED");
// 3. 清理缓存
redisTemplate.delete("order:status:" + orderId);
log.info("订单 {} 已自动取消(超时未支付)", orderId);
}
}
4. 支付成功后移除任务:
@Service
public class PaymentService {
@Autowired
private StringRedisTemplate redisTemplate;
public void paySuccess(Long orderId) {
// 1. 更新订单状态
orderMapper.updateStatus(orderId, "PAID");
// 2. 清除 Redis 中的状态
redisTemplate.delete("order:status:" + orderId);
// 3. 【可选】从 Stream 中删除消息(实际不可行,改为标记)
// 由于 Stream 消息无法删除,建议通过状态判断避免重复取消
// 更佳做法:在消费者中检查状态
}
}
12.2.5 异常处理与高可用保障
| 场景 | 处理策略 |
|---|---|
| 消费者宕机 | Redis Streams 支持 Pending Entries,重启后继续处理未 ACK 的消息 |
| 重复消费 | 在取消逻辑前检查订单状态(幂等性) |
| 消息堆积 | 监控 XPENDING 数量,增加消费者实例 |
| Redis 宕机 | 启用 AOF + RDB 持久化;生产环境建议 Redis 集群 |
| 时间精度 | 使用 XREADGROUP 阻塞读取,延迟可控制在秒级 |
监控命令:
# 查看待处理消息
XPENDING stream:order:delay group:order:cancel - + 10
# 查看流信息
XINFO STREAM stream:order:delay
# 查看消费者组
XINFO GROUPS stream:order:delay
12.2.6 优化建议
| 优化点 | 说明 |
|---|---|
| 批量消费 | 使用 count 参数一次处理多条消息,提升吞吐量 |
| 动态超时 | 不同商品可设置不同超时时间(如秒杀商品 5 分钟) |
| 降级策略 | Redis 不可用时,降级为 DB 轮询 + 本地缓存 |
| 监控告警 | 接入 Prometheus,监控 pending_messages、consume_latency |
| 压力测试 | 模拟 10w+ 订单写入,验证消费延迟和系统稳定性 |
12.2.7 替代方案:纯 Redis ZSet 实现(轻量级)
若不使用 Streams,可基于 ZSet 实现简单延迟队列:
// 添加任务
redisTemplate.opsForZSet().add("zset:order:delay", orderId.toString(), System.currentTimeMillis() + 900000); // 15分钟后
// 轮询处理(定时任务)
Set<String> expired = redisTemplate.opsForZSet().rangeByScore("zset:order:delay", 0, now);
for (String orderIdStr : expired) {
Long orderId = Long.valueOf(orderIdStr);
if ("UNPAID".equals(redisTemplate.opsForValue().get("order:status:" + orderId))) {
cancelOrder(orderId);
}
redisTemplate.opsForZSet().remove("zset:order:delay", orderIdStr);
}
缺点:需定时轮询,不支持 ACK,宕机可能丢失任务。