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

Java Redis 与 Elasticsearch:客户端、缓存一致性、检索和故障降级

Redis 和 Elasticsearch 经常同时出现在 Java 服务中:Redis 负责低延迟数据访问、缓存、分布式协调或短期状态;Elasticsearch 负责倒排索引、全文检索、聚合和复杂筛选。两者都能“保存数据”,但数据模型、持久化语义、一致性边界和故障表现完全不同。

如果把 Redis 当数据库、把 Elasticsearch 当强一致主库,或者把缓存更新和数据库更新写在同一个本地事务里,系统通常会在并发、重试、节点故障或索引延迟下暴露问题。正确的设计首先要明确:

  1. 谁是权威数据源;
  2. Redis 和 Elasticsearch 分别承担什么查询;
  3. 客户端如何管理连接、超时和错误;
  4. 缓存、数据库、索引之间允许多长时间的不一致;
  5. 检索失败时,业务如何降级;
  6. 重试是否会造成重复写入或流量放大。

本文以 Java 25 LTS 为语言运行时背景,示例采用 Spring Boot 生态中常见的 Spring Data Redis 和 Elasticsearch Java API Client。具体依赖版本应由项目使用的 Spring Boot 和 Elasticsearch 版本管理,不能把不同大版本的客户端 API 混用。


一、先划分职责:数据库、Redis 和 Elasticsearch 不是三份等价数据

一个典型商品服务可以这样划分:

写入请求
   |
   v
关系数据库 ──事务提交──> Outbox 事件 ──异步消费──> Elasticsearch 索引
   |                                      |
   |                                      v
   └──────────────缓存失效/更新──────────> Redis

关系数据库通常是商品名称、价格、库存、状态等业务事实的权威来源。Redis 保存可重建的缓存或短期状态;Elasticsearch 保存用于检索的派生索引。

这里的“权威”不是由组件的性能决定,而是由业务事实的生命周期决定:

  • 数据库中的商品价格是业务事实;
  • Redis 中的价格是缓存副本,丢失后可重新加载;
  • Elasticsearch 中的价格是检索文档中的副本,可能因为索引异步更新而短暂过期。

如果一次商品更新只写 Redis 而没有写权威数据库,Redis 发生故障或过期后,数据就无法恢复。这个设计不是缓存,而是把数据库责任交给了一个不应承担该责任的组件。

1. Redis 的适用边界

Redis 是内存优先的键值数据库,常见数据结构包括:

  • String:计数器、JSON 字符串、分布式锁的值;
  • Hash:对象字段;
  • List:队列或列表;
  • Set:集合关系;
  • Sorted Set:排行榜、延迟任务;
  • Stream:带消费组的消息流;
  • Bitmap、HyperLogLog 等特殊结构。

Redis 的核心访问模型是通过 key 定位值:

GET product:1001

它擅长的是已知 key 的快速访问,而不是对任意文本进行分词、相关性排序和复杂查询。

2. Elasticsearch 的适用边界

Elasticsearch 是面向搜索和分析的分布式文档系统。文档写入后会被分析器处理,文本字段通常进入倒排索引。查询时,搜索词可以映射到包含这些词项的文档集合,再结合相关性、过滤条件、排序和聚合返回结果。

它适合:

  • 商品名、文章正文等全文检索;
  • 分词、同义词、拼写或相关性搜索;
  • 多条件过滤和排序;
  • 聚合统计;
  • 按游标导出大量搜索结果。

它不适合直接替代关系数据库的事务主库,也不应被当作 Redis 的低延迟 key-value 缓存。

3. 三者的典型职责

需求 首选组件 原因
按主键读取事实数据 数据库或 Redis 缓存 事务事实由数据库维护,Redis 加速读取
商品名全文检索 Elasticsearch 倒排索引和分析器
热门商品计数 Redis 原子计数和高吞吐
订单状态变更 数据库 需要明确的事务边界
搜索结果中的展示字段 Elasticsearch 文档 允许短暂索引延迟
不能丢失的业务事件 数据库 Outbox、消息系统 Redis 和搜索索引不应作为唯一事件存储

二、Redis 客户端:连接、线程模型和错误边界

1. Lettuce 与 Jedis

Spring Data Redis 可以通过不同客户端访问 Redis。工程中常见的是 Lettuce 和 Jedis。

Lettuce 基于 Netty,支持同步、异步和响应式访问,连接通常可以在多个线程之间共享。Jedis 也提供同步访问,现代版本支持连接池等模式。两者都能完成普通 Redis 操作,但连接生命周期和并发使用方式不同。

不要把底层客户端对象的使用方式混在一起:

  • 使用 Lettuce 时,应遵循 Spring Data Redis 对连接复用和线程安全的封装;
  • 使用 Jedis 时,应正确配置连接池,借出的连接必须归还;
  • 不要在每个请求中手动创建和销毁 Redis 客户端;
  • 不要把单个非线程安全的底层连接对象随意放进单例 Bean 中并发使用。

Spring Boot 通常会自动配置 RedisConnectionFactoryRedisTemplate。一个最小配置如下:

spring:
  data:
    redis:
      host: localhost
      port: 6379
      timeout: 800ms
      connect-timeout: 300ms

这里的两个超时含义不同:

  • connect-timeout:建立连接最多等待多久;
  • timeout:执行 Redis 命令或读取响应最多等待多久,具体行为受客户端和 Spring Data Redis 版本影响。

超时不是越大越好。假设接口线程池有 200 个线程,每个请求都允许 Redis 阻塞 5 秒,那么 Redis 故障时最多可能有大量线程同时等待,最终拖垮业务线程池。超时应小于该接口的整体预算,并配合熔断或快速降级。

2. RedisTemplate 的序列化边界

一个可运行的缓存服务可以使用字符串 key 和 JSON value:

package com.example.product;

import java.time.Duration;

import org.springframework.data.redis.core.RedisTemplate;
import org.springframework.stereotype.Service;

@Service
public class ProductCache {

    private final RedisTemplate<String, ProductView> redisTemplate;

    public ProductCache(RedisTemplate<String, ProductView> redisTemplate) {
        this.redisTemplate = redisTemplate;
    }

    public ProductView get(String id) {
        return redisTemplate.opsForValue().get(key(id));
    }

    public void put(String id, ProductView value, Duration ttl) {
        redisTemplate.opsForValue().set(key(id), value, ttl);
    }

    public Boolean delete(String id) {
        return redisTemplate.delete(key(id));
    }

    private String key(String id) {
        return "product:" + id;
    }
}

record ProductView(
        String id,
        String name,
        long priceInCents,
        boolean enabled
) {}

要使这段代码可用,必须配置 RedisTemplate<String, ProductView> 的 key 和 value serializer。常见选择是:

  • key 使用 StringRedisSerializer
  • value 使用 JSON serializer;
  • 不直接依赖 Java 原生序列化,以免类名变化、反序列化风险和跨版本兼容问题。

Redis 中的 value 是字节序列。ProductView 是否能正确读取,不取决于 Java record 本身,而取决于写入端和读取端是否使用兼容的序列化协议。

缓存对象发生结构变更时,还要考虑:

  • 老数据是否仍能被新代码反序列化;
  • 是否需要在 key 中加入版本,例如 product:v2:1001
  • 是否需要灰度期间同时兼容旧格式;
  • 过期时间是否足以让旧格式自然淘汰。

3. Redis 原子操作与复合操作

以下代码不是原子的:

Long count = redisTemplate.opsForValue().get("view:1001");
redisTemplate.opsForValue().set("view:1001", count + 1);

两个请求可能都读到 10,最终都写回 11,丢失一次更新。

应使用 Redis 原子命令:

Long count = redisTemplate.opsForValue().increment("view:1001");

如果操作包含多个条件和多个 key,例如“只有库存大于 0 才扣减”,不能仅依赖多条普通命令。可以使用 Lua 脚本、Redis 事务或把扣减逻辑放回权威数据库。Redis 事务提供命令队列和执行顺序,但不会自动回滚已经执行的命令,也不等价于关系数据库事务。


三、缓存一致性:先定义“正确”,再选择模式

1. 一致性不是单一概念

“缓存一致性”至少涉及三个对象:

数据库值 D
Redis 缓存值 C
读取请求看到的值 R

常见的要求有:

  • 读己之写:同一个用户刚成功写入后,随后读取不能看到旧值;
  • 单调读:同一个会话后续读取不能从新值退回旧值;
  • 最终一致:经过有限时间,缓存最终与数据库一致;
  • 强一致:每个成功读取都符合严格的最新值约束。

大多数旁路缓存系统只承诺最终一致,不能因为使用了 Redis 就自动得到强一致。

设数据库更新在时间 tdt_d 提交,缓存失效在时间 tct_c 生效。如果:

tctdt_c \ge t_d

那么在不考虑并发读回填的情况下,失效动作不会早于数据库提交。若数据库提交后缓存删除失败,则缓存可能永久保存旧值,直到 TTL 到期。因此“删除失败后依靠 TTL”只能提供有界但不确定的一致性恢复。

2. 为什么“先更新数据库,再删除缓存”仍可能不一致

这是常见的 Cache-Aside(旁路缓存)流程:

写请求:
1. 更新数据库
2. 删除 Redis 缓存

读请求:
1. 查询 Redis
2. 未命中时查询数据库
3. 写回 Redis

看似合理,但存在以下并发时序:

初始:数据库 = A,缓存 = A

读线程 R:缓存未命中
写线程 W:数据库更新为 B
写线程 W:删除缓存
读线程 R:读取数据库,恰好读到 A
读线程 R:把 A 写入缓存

最终数据库是 B,缓存却是 A。

这个反例说明,单纯改变“更新数据库”和“删除缓存”的顺序不能消除所有竞争。常用缓解方式包括:

  1. 数据库更新后删除缓存;
  2. 删除失败通过重试、消息或 Outbox 补偿;
  3. 回填缓存时使用版本号;
  4. 删除后延迟再次删除;
  5. 对关键读写使用短期一致性策略,而不是把所有读取都交给缓存。

“延迟双删”只能降低某些竞态概率,不能在没有版本控制或协调机制的情况下证明强一致。它还会增加删除流量,延迟时间也很难凭经验固定。

3. 带版本号的缓存回填

可以为数据库记录增加单调递增版本 v

数据库:
  product:1001 = (value=B, version=8)

缓存:
  product:1001 = (value=A, version=7)

回填或更新缓存时,只允许更高版本覆盖更低版本:

write(vnew) succeeds only if vnewvold\text{write}(v_{\text{new}}) \text{ succeeds only if } v_{\text{new}} \ge v_{\text{old}}

实际实现需要 Redis Lua 脚本保证“比较版本并写入”是原子的。否则两个客户端仍可能分别读取旧版本并互相覆盖。

版本机制能解决“旧回填覆盖新回填”,但不能凭空产生数据库版本。版本必须由数据库或事件流可靠地产生,而且读取数据库时必须能获得与数据一致的版本。

4. Cache-Aside 的完整实现

下面是一个简化的读取流程:

@Service
public class ProductService {

    private final ProductRepository repository;
    private final ProductCache cache;

    public ProductService(ProductRepository repository, ProductCache cache) {
        this.repository = repository;
        this.cache = cache;
    }

    public ProductView get(String id) {
        ProductView cached = cache.get(id);
        if (cached != null) {
            return cached;
        }

        ProductView loaded = repository.findViewById(id)
                .orElseThrow(() -> new ProductNotFoundException(id));

        cache.put(id, loaded, Duration.ofMinutes(5));
        return loaded;
    }

    @Transactional
    public void update(ProductCommand command) {
        repository.update(command);
        cache.delete(command.id());
    }
}

这里的关键不是代码本身,而是事务边界:

  • repository.update 所在事务提交前,删除缓存可能过早;
  • Spring 的 @Transactional 只管理当前数据库事务,不能自动把 Redis 删除纳入同一个原子事务;
  • 如果数据库提交成功而 Redis 删除失败,必须有补偿路径;
  • 如果数据库回滚,缓存不应被错误地删除或更新到未提交数据。

更稳妥的做法是注册事务提交后的动作,或者写入 Outbox:

数据库事务:
  1. 更新 product
  2. 插入 product_changed 事件
  3. 一起提交

异步消费者:
  4. 删除 product:1001
  5. 删除失败则重试

Outbox 不能保证 Redis 永远可用,但可以避免“数据库更新成功、事件完全丢失”。

5. 缓存击穿、穿透和雪崩

这三个词描述的是不同问题。

缓存击穿

某个热点 key 过期,大量请求同时回源数据库:

10000 个请求
     |
     +-- Redis miss
     |
     +-- 10000 次数据库查询

可以使用:

  • 单飞(single flight):同一 key 只有一个请求回源;
  • 分布式锁;
  • 热点 key 逻辑不过期,后台异步刷新;
  • TTL 加随机抖动,减少同时过期。

但分布式锁本身也有故障边界。锁必须设置过期时间,释放时验证 value,不能无条件删除其他请求持有的锁。Redis 锁不能替代数据库事务。

缓存穿透

请求持续访问不存在的 key,每次都命中不到缓存并查询数据库。可以缓存短 TTL 的空值:

product:does-not-exist -> NOT_FOUND,TTL 30 秒

空值缓存的风险是:如果同一个 key 后来被创建,短时间内仍可能读到“不存在”。因此空值 TTL 应短于正常对象,并在创建时主动删除对应 key。

缓存雪崩

大量 key 在同一时间过期,或 Redis 集群整体不可用,导致请求集中访问数据库。随机 TTL 可以降低集中过期,但 Redis 整体故障仍需要限流、熔断和降级。


四、Elasticsearch 建模:文档、映射、分析器和刷新

1. 文档不是关系表的一行

Elasticsearch 的基本存储单位是 JSON 文档。文档有 _id、字段和索引名。字段的 mapping 决定它如何被索引。

例如:

PUT products-v1
{
  "settings": {
    "analysis": {
      "analyzer": {
        "product_text": {
          "type": "custom",
          "tokenizer": "standard",
          "filter": ["lowercase"]
        }
      }
    }
  },
  "mappings": {
    "properties": {
      "id":       { "type": "keyword" },
      "name":     { "type": "text", "analyzer": "product_text" },
      "category": { "type": "keyword" },
      "price":    { "type": "long" },
      "enabled":  { "type": "boolean" },
      "updatedAt": { "type": "date" }
    }
  }
}

字段类型的语义不同:

  • text:用于全文检索,通常会被分析器分词;
  • keyword:作为完整值,用于精确匹配、过滤、排序和聚合;
  • long:数值比较和聚合;
  • date:时间过滤和排序;
  • nested:保留数组对象内部字段的关联关系。

例如商品名称 "Java Redis 实战" 作为 text 可能被分析为多个 token,而 category = "book" 作为 keyword 必须整体匹配。

如果把 category 错误映射成 text,下面的精确过滤和聚合可能无法按预期工作:

{
  "term": {
    "category": "book"
  }
}

全文搜索使用 match,精确过滤使用 term。二者不能因为都叫“查询”而互换。

2. objectnested 的边界

假设文档包含:

{
  "variants": [
    { "color": "red", "size": "S" },
    { "color": "blue", "size": "L" }
  ]
}

如果 variants 是普通 object,查询:

color = red AND size = L

可能把不同数组元素的字段组合起来,错误地命中文档。

若业务要求颜色和尺码必须来自同一个变体,就应使用 nested

"variants": {
  "type": "nested",
  "properties": {
    "color": { "type": "keyword" },
    "size":  { "type": "keyword" }
  }
}

nested 查询会在同一个嵌套对象范围内判断条件,但索引成本和查询复杂度更高。不能为了“结构看起来更准确”而无条件使用。

3. Refresh 与近实时搜索

Elasticsearch 写入成功不等于搜索请求立即可见。文档先写入事务日志和内存缓冲,之后通过 refresh 使其对搜索可见。因此 Elasticsearch 通常是 near real-time,即近实时搜索。

需要区分:

  • 写入响应成功:节点接受并持久化了写请求的相关状态;
  • refresh 完成:后续搜索可以看到文档;
  • replica 同步:副本是否已完成同步;
  • refresh=wait_for:请求等待下一次 refresh,使搜索可见,但会增加写请求等待时间。

不要在每一次业务写入后强制 refresh=true。这会显著增加索引段生成和合并压力。更合理的是依赖正常 refresh 周期,只有测试、低频管理操作或确有读后写要求的场景才考虑等待 refresh。


五、使用 Elasticsearch Java API Client 进行检索

现代 Elasticsearch Java API Client 的核心包通常是:

co.elastic.clients.elasticsearch.ElasticsearchClient

具体构造方式会因 Spring Boot、Elasticsearch Java API Client 和 Spring Data Elasticsearch 版本而变化。下面给出客户端 API 的核心调用形式,依赖和连接 Bean 应由项目版本对应的 Spring Boot 自动配置或官方客户端配置完成。

1. 查询代码

import co.elastic.clients.elasticsearch.ElasticsearchClient;
import co.elastic.clients.elasticsearch._types.query_dsl.Query;
import co.elastic.clients.elasticsearch.core.SearchResponse;
import org.springframework.stereotype.Service;

import java.io.IOException;
import java.util.List;

@Service
public class ProductSearchService {

    private final ElasticsearchClient client;

    public ProductSearchService(ElasticsearchClient client) {
        this.client = client;
    }

    public List<ProductDocument> search(String keyword, String category)
            throws IOException {

        Query textQuery = Query.of(q -> q.multiMatch(m -> m
                .query(keyword)
                .fields("name^3", "description")));

        SearchResponse<ProductDocument> response = client.search(s -> s
                        .index("products-read")
                        .query(q -> q.bool(b -> b
                                .must(textQuery)
                                .filter(f -> f.term(t -> t
                                        .field("category")
                                        .value(category)))
                                .filter(f -> f.term(t -> t
                                        .field("enabled")
                                        .value(true)))))
                        .from(0)
                        .size(20),
                ProductDocument.class);

        return response.hits().hits().stream()
                .map(hit -> hit.source())
                .filter(java.util.Objects::nonNull)
                .toList();
    }
}

record ProductDocument(
        String id,
        String name,
        String description,
        String category,
        long price,
        boolean enabled
) {}

查询语义如下:

  • multi_match 在多个文本字段搜索;
  • "name^3" 表示名称字段的相关性权重高于描述字段;
  • must 参与匹配和相关性评分;
  • filter 只判断是否满足条件,通常适合不需要评分的结构化条件;
  • size(20) 限制本页返回数量;
  • products-read 是读别名,不直接绑定具体物理索引。

如果 categorytext 而不是 keyword,这个 term 查询就可能不符合预期。因此查询语句和 mapping 必须一起设计和测试。

2. mustfiltershould

一个查询可以写成:

{
  "query": {
    "bool": {
      "must": [
        { "multi_match": {
          "query": "redis 缓存",
          "fields": ["name^3", "description"]
        }}
      ],
      "filter": [
        { "term":  { "category": "book" }},
        { "range": { "price": { "lte": 10000 }}},
        { "term":  { "enabled": true }}
      ],
      "should": [
        { "term": { "brand": "official" }}
      ],
      "minimum_should_match": 0
    }
  }
}

直觉上:

  • must:必须匹配,且通常影响 _score
  • filter:必须匹配,但不参与相关性评分;
  • should:提高匹配偏好,是否必须满足取决于上下文和 minimum_should_match

价格范围、启用状态、类别等通常不需要相关性评分,放在 filter 中更清晰。品牌偏好可以放在 should 中,让它影响排序但不排除其他结果。

3. 分页:from/size 不是无限可扩展

传统分页:

第 1 页:from = 0,size = 20
第 1000 页:from = 19980,size = 20

深页查询需要维护并跳过大量前置结果,成本会增加。Elasticsearch 还通常限制 from + size 的窗口大小,避免单次请求消耗过多内存。

对于用户滚动浏览,使用 search_after

SearchResponse<ProductDocument> response = client.search(s -> s
        .index("products-read")
        .size(20)
        .sort(so -> so.field(f -> f.field("updatedAt").order(
                co.elastic.clients.elasticsearch._types.SortOrder.Desc)))
        .sort(so -> so.field(f -> f.field("id").order(
                co.elastic.clients.elasticsearch._types.SortOrder.Asc)))
        .searchAfter(lastUpdatedAt, lastId),
        ProductDocument.class);

实际代码中的 searchAfter 值必须与上一页最后一条文档的排序值完全对应,并且排序需要稳定。仅按 updatedAt 排序时,如果多条文档时间相同,分页可能重复或遗漏,因此常加入唯一的 id 作为 tie-breaker。

如果需要在长时间导出期间保持一致的搜索视图,可以使用 Point in Time(PIT)配合 search_after。PIT 会占用集群资源,必须设置合理的 keep-alive,并在完成后关闭。

4. 聚合不是业务事务统计

Elasticsearch 聚合适合搜索页面上的分类计数、价格区间统计和趋势图,但它们是索引视图上的统计结果:

  • 可能落后于数据库;
  • 受 refresh 和副本状态影响;
  • 大基数字段聚合会消耗内存;
  • 不能直接等价于财务或库存的强一致统计。

订单金额结算、库存扣减等结果应来自权威事务系统,而不是搜索聚合。


六、索引写入:幂等、批量和零停机映射变更

1. 用稳定 _id 实现幂等

将数据库主键作为 Elasticsearch 文档 _id

数据库 product.id = 1001
Elasticsearch document _id = 1001

同一事件重试时,使用相同 _id 执行 index 或带条件的更新,通常可以避免产生重复文档。

但“相同 ID”不代表所有操作都自动幂等:

  • index 同 ID 通常是覆盖写;
  • update 可能执行脚本或部分字段更新;
  • increment 类逻辑若重试可能重复累加;
  • 批量请求可能部分成功,必须逐项检查失败结果。

消费者应记录事件版本或序列号,例如:

product 1001:
事件 v=7:价格 100
事件 v=8:价格 120

如果 v=8 先到,之后 v=7 重试,消费者必须拒绝旧版本覆盖新版本。否则异步乱序会把索引回退。

2. Bulk 的错误处理

批量写入能降低请求往返开销,但响应中的单项结果可能不同:

Bulk 请求:
  item 1 -> 成功
  item 2 -> mapping 错误
  item 3 -> 429 rejected

因此不能只判断 HTTP 请求是否返回成功。处理逻辑至少要区分:

  • 可重试错误:临时拒绝、节点不可用、网络超时;
  • 不可重试错误:字段类型冲突、非法文档、映射错误;
  • 业务无效事件:源数据不存在或版本过旧。

重试必须有上限和退避。无限重试会把 Elasticsearch 的压力再次放大。常见策略是指数退避并加入随机抖动:

dn=min(dmax,d0×2n)+random(0,j)d_n = \min(d_{\max}, d_0 \times 2^n) + \text{random}(0, j)

其中 d0d_0 是初始等待时间,nn 是重试次数,dmaxd_{\max} 是上限,jj 是抖动范围。

3. 别名和重建索引

Elasticsearch 的 mapping 一旦上线,修改字段类型通常不能直接完成。常见流程是:

products-v1
    |
    | 创建新索引、重建数据、验证
    v
products-v2

原子切换 products-read:
products-v1  -> products-read
products-v2  -> products-read

读请求始终访问 products-read,写入别名可以单独管理。切换别名后,新请求立即指向新索引,避免应用同时维护物理索引名。

验证不能只看索引文档数量,还应检查:

  • mapping 是否符合预期;
  • 中文、英文、数字和特殊字符的分词结果;
  • 过滤、排序、聚合是否正确;
  • 典型查询的相关性;
  • 错误文档是否被记录;
  • 新旧索引结果差异是否在可接受范围。

七、数据库、缓存和索引的更新顺序

1. 不要把跨组件操作误认为一个事务

Spring 的 @Transactional 默认管理某个事务资源,例如关系数据库。以下方法并不天然具备“数据库和 Redis 一起提交”的语义:

@Transactional
public void updateProduct(ProductCommand command) {
    repository.update(command);  // 数据库
    cache.delete(command.id());  // Redis
}

如果 Redis 删除成功、数据库随后回滚,缓存可能被删除,但这通常只是一次可接受的缓存未命中;如果数据库提交成功、Redis 删除失败,缓存可能保留旧值。二者没有一个通用的本地原子提交协议。

更重要的是,Elasticsearch 通常是异步索引目标,不应强行加入业务数据库的同步事务。同步双写的典型失败路径是:

数据库提交成功
Elasticsearch 请求超时
客户端重试
应用不知道第一次请求是否已经成功

这要求写入具备幂等性和补偿机制,而不是简单地在一个方法里连续调用两个客户端。

2. Outbox 的因果关系

推荐的因果链是:

同一个数据库事务:
  更新商品
  插入 outbox(product_changed, version=8)

事务提交后:
  发布器读取 outbox
  投递到消息系统
  搜索消费者更新 Elasticsearch
  缓存消费者删除或刷新 Redis

这样至少能保证:

数据库事实提交变更事件不会因应用进程崩溃而完全丢失\text{数据库事实提交} \Rightarrow \text{变更事件不会因应用进程崩溃而完全丢失}

但仍需处理:

  • outbox 重复投递;
  • 消费者重复消费;
  • 消息乱序;
  • 消费者长期积压;
  • Elasticsearch 暂时不可用;
  • Redis 删除失败。

所以事件必须包含稳定事件 ID、实体 ID、版本、事件类型和发生时间。消费者要以幂等方式处理。

3. 搜索结果和详情读取的组合

一种常见设计是:

搜索请求
  1. Elasticsearch 返回商品 ID、标题、价格快照
  2. 直接展示搜索文档中的字段
  3. 点击详情时按 ID 读取 Redis/数据库

另一种设计是:

搜索请求
  1. Elasticsearch 只返回 ID
  2. 批量从 Redis/数据库读取最新详情

第二种能降低展示旧价格的风险,但会增加 N+1 或批量回源成本。不能在每个搜索结果上逐个查询数据库。若必须回源,应使用批量查询,并明确搜索结果中的排序和详情字段可能来自不同时间点。


八、故障降级:先划分故障类型,再决定返回什么

“降级”不是简单返回空列表。它是故障发生时,按照业务优先级选择可接受的替代路径。

1. Redis 故障路径

Redis 超时或连接失败时,缓存读取通常可以降级为数据库读取:

Redis 成功:
  返回缓存

Redis 超时/连接失败:
  记录指标
  进入数据库读取

但数据库回源必须有保护:

  • 限制同时回源的并发数;
  • 对热点 key 使用单飞;
  • 设置更短的数据库超时;
  • 使用熔断器快速拒绝过量请求;
  • 对非核心接口返回降级结果。

如果 Redis 保存的是会话、限流计数或幂等键,不能总是“忽略 Redis 故障”:

  • 会话服务可能无法验证身份;
  • 限流失效可能造成下游过载;
  • 幂等键丢失可能导致重复业务操作。

此时降级策略可能是拒绝请求,而不是继续执行。

2. Elasticsearch 故障路径

搜索索引不可用时,常见降级方式有:

  1. 返回热门商品或预先生成的推荐结果;
  2. 按结构化条件查询数据库,但不提供全文相关性;
  3. 返回明确的“搜索暂不可用”,保留详情和其他业务功能;
  4. 使用备用只读索引集群;
  5. 对无关键词请求切换为数据库分页。

不能把“搜索失败”伪装成“搜索成功但结果为空”。两者语义不同:

  • 空结果:系统确认没有匹配数据;
  • 搜索失败:系统没有完成查询,结果未知。

API 可以使用明确的响应字段:

{
  "items": [],
  "degraded": true,
  "message": "搜索服务暂时不可用"
}

是否向最终用户暴露具体错误文字,取决于产品设计;但内部日志和指标必须区分“无结果”和“查询失败”。

3. 熔断、超时和重试的关系

一次远程调用应有明确的时间预算:

HTTP 总预算 800 ms
  ├─ Redis 100 ms
  ├─ Elasticsearch 500 ms
  └─ 序列化和业务处理 200 ms

如果 Elasticsearch 客户端超时设置为 2 秒,HTTP 接口却只有 800 ms,那么请求必然在客户端返回之前已经超时。超时预算必须自上而下传递。

重试会改变实际流量。若一次请求初始发送 100 QPS,失败后每次最多重试 2 次,理论上可能产生:

100×(1+2)=300 QPS100 \times (1 + 2) = 300\ \text{QPS}

如果每一层都重试,多个微服务会形成乘法放大。因此通常应:

  • 只在最接近故障源的一层进行有限重试;
  • 只重试明确的临时错误;
  • 对写操作要求幂等后再重试;
  • 使用指数退避和抖动;
  • 达到阈值后由熔断器快速失败。

4. 熔断状态

熔断器通常有三个逻辑状态:

CLOSED
  请求正常通过
  失败率超过阈值
      |
      v
OPEN
  快速拒绝,不访问故障组件
  等待恢复时间
      |
      v
HALF_OPEN
  放行少量探测请求
  成功 -> CLOSED
  失败 -> OPEN

熔断器不是恢复机制。它只减少对故障组件的压力,并保护调用方线程。恢复仍依赖 Redis 或 Elasticsearch 自身恢复、连接重建、索引修复等操作。


九、Spring 中的错误处理和资源生命周期

1. 不要把远程异常全部包装成业务“空值”

下面的写法会隐藏故障:

public List<ProductDocument> searchSafely(String keyword) {
    try {
        return searchService.search(keyword, null);
    } catch (Exception e) {
        return List.of();
    }
}

它把以下情况都变成空列表:

  • Elasticsearch 没有匹配;
  • 连接超时;
  • mapping 错误;
  • 程序序列化失败;
  • 查询参数非法。

正确做法是区分异常类型,并在边界处转换:

public SearchResult searchWithFallback(String keyword) {
    try {
        return SearchResult.success(searchService.search(keyword, null));
    } catch (java.net.SocketTimeoutException e) {
        metrics.increment("search.timeout");
        return SearchResult.degraded(popularProducts());
    } catch (IOException e) {
        metrics.increment("search.io_error");
        return SearchResult.degraded(popularProducts());
    }
}

示例中的 SearchResultpopularProducts() 需要由业务实现;重点是降级结果应携带状态,而不是伪装成正常空结果。

2. 连接池和线程池不是越大越好

Redis 客户端连接、Elasticsearch HTTP 连接和 Web 请求线程池之间存在资源链路:

请求线程
  -> Elasticsearch 连接池
      -> ES 节点线程池
          -> 磁盘、CPU、segment merge

连接池过大可能让更多请求同时进入已经过载的 Elasticsearch,增加排队和超时;过小则可能限制吞吐。连接数应结合:

  • 请求并发;
  • 单请求平均耗时;
  • 节点数和分片数;
  • 客户端连接池限制;
  • 下游线程池拒绝策略;
  • 实测延迟分位数。

不能根据“CPU 有很多核”直接推导出应配置多少连接。

3. 配置客户端超时

不同 Spring Boot 和客户端版本的配置键可能不同,不能复制未经版本验证的参数。原则上至少应配置:

  • 连接建立超时;
  • socket/read 超时;
  • 连接池获取超时;
  • 每个请求的业务超时;
  • 最大响应体大小;
  • TLS、认证和证书校验。

生产环境必须验证配置是否真正生效。一个常见错误是配置写在了旧版本的属性路径下,应用启动成功但客户端仍使用默认值。

验证方法包括:

  • 启动日志确认客户端连接设置;
  • 通过故障注入制造延迟;
  • 观察请求是否在预期时间失败;
  • 检查连接池活动数和等待数;
  • 使用管理端点或指标确认熔断器状态。

十、缓存和搜索的测试方式

1. 缓存一致性并发测试

构造如下测试:

初始:数据库=A,缓存=A

线程 W:
  更新数据库为 B
  模拟提交后暂停
  删除缓存

线程 R:
  缓存 miss
  数据库读取
  模拟延迟
  回填缓存

测试需要验证:

  • 是否出现数据库为 B、缓存为 A;
  • 删除失败后是否能补偿;
  • 回填是否有版本保护;
  • 读取是否会在写成功后看到旧值;
  • TTL 到期后是否能恢复。

单元测试只能验证局部逻辑,并发竞态通常需要集成测试、延迟注入和多线程压力测试。

2. Elasticsearch 查询测试

测试数据至少应覆盖:

  • 空字符串和只含空白字符的关键词;
  • 中文、英文、数字和混合词;
  • 同义词或大小写处理;
  • 未映射字段;
  • 价格边界;
  • 相同排序值;
  • 深页和 search_after
  • 索引不存在;
  • 部分分片失败;
  • refresh 前后的可见性。

对检索结果不要只断言“数量正确”,还要断言:

  • 精确过滤是否排除错误文档;
  • 相关性排序是否符合业务预期;
  • 同一分页游标是否不会重复;
  • 聚合桶是否来自正确字段类型;
  • 降级响应是否与正常空结果区分。

3. 故障注入

可以在测试或预生产环境中执行:

1. 停止 Redis
2. 注入 Redis 读延迟
3. 停止 Elasticsearch 节点
4. 返回 429 或模拟连接拒绝
5. 让 outbox 消费者暂停
6. 恢复组件并观察积压是否清空

每个实验都应有:

  • 预期用户行为;
  • 预期错误码或降级状态;
  • 预期指标变化;
  • 恢复步骤;
  • 数据一致性验证。

如果 Redis 停止后所有请求都打爆数据库,说明缓存降级只有“路径”,没有容量保护。


十一、监控指标必须对应故障机制

Redis 指标

至少观察:

  • 命中率;
  • 命中和未命中数量;
  • 命令延迟分位数;
  • 连接失败和超时;
  • key 过期数量;
  • 内存使用和淘汰;
  • 热 key;
  • 慢查询;
  • 主从复制延迟或集群重定向错误。

命中率高不代表缓存正确。一个长期保存旧值的缓存可能有很高命中率,因此还需要记录缓存版本、回源校验或业务数据新鲜度。

Elasticsearch 指标

至少观察:

  • 查询延迟分位数;
  • 429、5xx、超时;
  • 分片失败;
  • refresh 延迟;
  • indexing pressure;
  • JVM heap 和 GC;
  • segment merge;
  • 搜索线程池拒绝;
  • 索引文档数量;
  • outbox 消费延迟;
  • 事件版本落后数量。

只看平均延迟会掩盖长尾。用户感知通常由 p95、p99 和超时率决定。

一致性指标

可以定义:

index_lag = 当前数据库版本 - Elasticsearch 已消费版本
cache_invalidation_lag = 数据库提交时间 - 缓存失效完成时间

如果事件按实体有版本号,就能直接发现某个商品的索引是否落后,而不是等用户投诉搜索结果过期。


十二、常见错误及其失败表现

错误一:把 Redis 当作唯一数据源

失败表现:

  • Redis 重启后业务数据消失;
  • 淘汰策略触发后出现大量“数据不存在”;
  • 无法通过事务恢复历史状态。

原因:

缓存没有权威来源,重建条件不存在。

修复:

把 Redis 限定为可重建副本,或明确将其作为数据库使用并为持久化、备份、恢复和数据模型承担相应责任。

错误二:把 Elasticsearch 写入成功当作业务事务提交

失败表现:

  • 搜索结果存在,但数据库订单状态未提交;
  • 索引有脏文档;
  • 双写失败时无法判断哪一边成功。

原因:

Elasticsearch 请求不在关系数据库本地事务的原子提交范围内。

修复:

数据库事务写事实和 Outbox,异步、幂等地构建搜索索引。

错误三:使用 text 字段执行 term 过滤

失败表现:

  • 查询返回零结果;
  • 聚合桶异常;
  • 大小写、分词导致精确匹配失效。

原因:

text 字段经过分析,索引中的 token 不一定等于原始字符串。

修复:

需要全文检索的字段使用 text,需要精确过滤的字段使用 keyword,必要时同时建立多字段。

错误四:搜索失败直接返回空列表

失败表现:

  • 监控显示接口成功,但用户误以为没有数据;
  • Elasticsearch 长期故障没有触发告警;
  • 业务方无法区分无结果和系统错误。

原因:

错误语义被丢弃。

修复:

保留降级标记、错误指标和日志;对核心流程选择备用数据源或明确错误响应。

错误五:无限重试远程调用

失败表现:

  • 下游已经过载,调用方继续发送请求;
  • 线程池耗尽;
  • 一个用户请求产生多次写入;
  • 整条调用链级联超时。

原因:

没有区分可重试错误、没有上限,也没有幂等语义。

修复:

设置总预算、有限重试、指数退避、熔断和幂等键。


十三、一个可落地的决策模型

面对一个新需求,可以按以下顺序判断。

第一步:确定权威事实

询问“如果 Redis 和 Elasticsearch 都丢失,能否从某处完整重建”。如果答案是否定的,应重新审查数据归属。

第二步:确定查询类型

  • 已知主键:优先缓存加数据库;
  • 全文、相关性、多过滤和聚合:使用 Elasticsearch;
  • 强一致事务查询:使用关系数据库;
  • 热点计数和短期状态:使用 Redis。

第三步:确定一致性预算

例如:

商品搜索结果允许落后 10 秒;
商品详情价格要求成功写入后立即读取最新值;
库存扣减不允许依赖缓存副本。

不同字段可以有不同边界,不要用一个“全局最终一致”掩盖所有业务要求。

第四步:设计失败路径

对每个依赖分别回答:

Redis 超时怎么办?
Elasticsearch 返回 429 怎么办?
数据库提交成功但缓存删除失败怎么办?
事件重复消费怎么办?
事件乱序怎么办?
搜索没有结果和搜索失败如何区分?

如果只能回答“重试”,说明故障设计还不完整。

第五步:定义验证证据

设计必须能被观测和验证:

  • 事件版本是否连续;
  • 索引延迟是否超过预算;
  • 缓存失效是否失败;
  • 降级次数是否增加;
  • 重试是否造成流量放大;
  • 恢复后积压是否清空。

结语

Java 服务同时使用 Redis 和 Elasticsearch 时,最重要的不是记住某个客户端 API,而是建立清晰的因果链:

权威数据库提交事实
        |
        v
可靠事件记录
        |
        +--> Redis 缓存失效或刷新
        |
        +--> Elasticsearch 幂等索引

Redis 解决的是快速按 key 读取以及部分短期状态问题;Elasticsearch 解决的是面向索引的检索和分析问题;缓存一致性解决的是副本何时、以什么方式追上权威事实;故障降级解决的是依赖不可用时系统还能保证什么。

当客户端超时、事务边界、事件版本、索引刷新、分页游标和降级语义都被明确写进设计后,Redis 与 Elasticsearch 才是两个可控的基础设施,而不是两个需要“出问题再重试”的远程服务。


系列导航与关联阅读

官方资料

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