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

Java Elasticsearch 工程:客户端、Mapping、查询、Bulk 和重试

Elasticsearch 是一个分布式搜索与分析引擎。Java 工程通常通过官方 Elasticsearch Java API Client 访问它,使用 Mapping 描述索引中的字段结构,使用 Query DSL 表达查询条件,通过 Bulk 批量写入,最后根据错误类型实施有限且可观测的重试。

这几个概念并不是相互独立的:

  • Mapping 决定 JSON 字段如何被索引,以及查询时采用什么语义。
  • 查询请求依赖 Mapping;textkeyword、数值和日期字段的查询行为不同。
  • Bulk 只是多个写操作的批处理协议,不等于数据库事务。
  • 重试必须理解客户端异常、HTTP 状态、Bulk 的逐项结果和写入幂等性。
  • Java 客户端负责类型安全的请求构造和响应解析,但不替应用程序承担数据模型、并发控制和失败恢复。

本文以 Java 25 LTS 为编译和运行环境,以 Elasticsearch 8.x Java API Client 风格为例。Java 25 是客户端运行时版本,不能据此推断 Elasticsearch 服务端版本;客户端、服务端和 Java 运行时仍需分别确认兼容性。示例使用本地、关闭安全认证的 Elasticsearch,仅适合开发环境。


一、先建立整体模型:Java 应用如何访问 Elasticsearch

一次典型的写入和查询链路如下:

sequenceDiagram
    participant App as Java 应用
    participant Client as ElasticsearchClient
    participant Node as Elasticsearch 节点
    participant Shard as 主分片
    participant Replica as 副本分片

    App->>Client: 构造 Index/Bulk/Search 请求
    Client->>Node: HTTP JSON 请求
    Node->>Shard: 路由到目标主分片
    Shard->>Shard: 解析 Mapping、写入倒排索引或 Doc Values
    Shard->>Replica: 复制操作
    Replica-->>Shard: 副本确认
    Shard-->>Node: 返回操作结果
    Node-->>Client: HTTP 响应
    Client-->>App: 类型化响应或异常

这里有几个容易混淆的边界:

  1. Java 客户端不是 Elasticsearch 服务端。
    客户端只负责构造 HTTP 请求、序列化 JSON、解析响应和暴露 Java API。

  2. 索引不是数据库表的完全等价物。
    Elasticsearch 的索引包含分片、副本、倒排索引、Doc Values 等搜索结构;它不提供关系数据库那种通用事务、跨表 Join 和外键约束。

  3. 写入成功不必然意味着搜索立即可见。
    Elasticsearch 通过 refresh 让最近写入的内容对搜索可见。默认 refresh 周期下,写入响应返回后,搜索可能仍暂时查不到文档。

  4. Bulk 响应可能整体 HTTP 成功,但部分文档失败。
    这与 JDBC 批处理或单条 SQL 的直觉不同。应用必须检查每一个 Bulk item。


二、Java 客户端:传输层、API 层和生命周期

2.1 推荐使用官方 Java API Client

当前 Java 工程通常使用:

co.elastic.clients:elasticsearch-java

它与旧的 High Level REST Client 不是同一套 API。新工程不应把两套客户端的代码风格混合使用。

Java API Client 大致分为三层:

ElasticsearchClient
        ↓
ElasticsearchTransport
        ↓
Apache HTTP client / HTTP connection
        ↓
Elasticsearch REST API
  • ElasticsearchClient:提供 indexsearchbulkindices().create 等类型化 API。
  • ElasticsearchTransport:负责 HTTP 传输和 JSON 映射。
  • Apache HTTP Client 5:通常作为底层 HTTP 实现。

客户端实例应当复用,而不是每次请求重新创建。创建客户端通常涉及连接池、线程和连接配置;每次方法调用都创建并关闭客户端,会浪费连接复用能力,甚至造成连接泄漏或线程膨胀。

2.2 Maven 依赖

下面的版本只是一个 8.x 示例。实际工程应选择与集群版本兼容、并经过验证的版本;不要把版本号硬编码为“Java 25 对应版本”,因为 Java 版本和 Elasticsearch 客户端版本没有这种一一对应关系。

<properties>
    <maven.compiler.release>25</maven.compiler.release>
    <elasticsearch.version>8.19.0</elasticsearch.version>
</properties>

<dependencies>
    <dependency>
        <groupId>co.elastic.clients</groupId>
        <artifactId>elasticsearch-java</artifactId>
        <version>${elasticsearch.version}</version>
    </dependency>

    <dependency>
        <groupId>org.apache.httpcomponents.client5</groupId>
        <artifactId>httpclient5</artifactId>
        <version>5.4.3</version>
    </dependency>
</dependencies>

依赖版本必须以实际发布版本和组织的依赖管理策略为准。生产项目还应通过 Maven Enforcer、依赖锁定或 SBOM 检查传递依赖,而不是只检查直接依赖。

2.3 创建客户端

开发环境关闭 Elasticsearch 安全认证时,可以使用 HTTP:

import co.elastic.clients.elasticsearch.ElasticsearchClient;
import co.elastic.clients.json.jackson.JacksonJsonpMapper;
import co.elastic.clients.transport.rest_client.RestClientTransport;
import org.apache.http.HttpHost;
import org.elasticsearch.client.RestClient;

public final class EsClients {

    private EsClients() {
    }

    public static ElasticsearchClient createLocalClient() {
        RestClient lowLevelClient = RestClient.builder(
                new HttpHost("localhost", 9200, "http")
        ).build();

        RestClientTransport transport = new RestClientTransport(
                lowLevelClient,
                new JacksonJsonpMapper()
        );

        return new ElasticsearchClient(transport);
    }
}

这里的生命周期是:

  1. 创建底层 RestClient
  2. 用它创建 RestClientTransport
  3. 用 transport 创建 ElasticsearchClient
  4. 应用关闭时关闭 transport 或底层客户端。

一个更完整的资源管理示例:

import co.elastic.clients.elasticsearch.ElasticsearchClient;
import co.elastic.clients.transport.rest_client.RestClientTransport;
import org.apache.http.HttpHost;
import org.elasticsearch.client.RestClient;

public final class ClientHolder implements AutoCloseable {

    private final RestClient lowLevelClient;
    private final RestClientTransport transport;
    private final ElasticsearchClient client;

    public ClientHolder() {
        this.lowLevelClient = RestClient.builder(
                new HttpHost("localhost", 9200, "http")
        ).build();

        this.transport = new RestClientTransport(
                lowLevelClient,
                new co.elastic.clients.json.jackson.JacksonJsonpMapper()
        );

        this.client = new ElasticsearchClient(transport);
    }

    public ElasticsearchClient client() {
        return client;
    }

    @Override
    public void close() throws Exception {
        transport.close();
    }
}

生产环境通常还要配置 HTTPS、认证、连接超时、读取超时、连接池和节点发现。关闭安全认证的本地配置不能直接复制到生产环境,否则用户名、密码和文档内容都可能以明文传输。


三、Mapping:JSON 字段如何变成可搜索结构

3.1 Mapping 的定义

Mapping 是 Elasticsearch 对索引字段的类型和索引方式描述。它至少回答三个问题:

  1. 字段是什么类型?
  2. 字段是否建立索引?
  3. 字段如何被分析、排序、聚合或存储?

例如:

{
  "properties": {
    "name": {
      "type": "text",
      "fields": {
        "keyword": {
          "type": "keyword"
        }
      }
    },
    "price": {
      "type": "double"
    },
    "category": {
      "type": "keyword"
    },
    "createdAt": {
      "type": "date"
    }
  }
}

这里的 name 同时具有两个视图:

  • nametext,适合全文搜索。
  • name.keywordkeyword,适合精确匹配、排序和聚合。

3.2 textkeyword 的区别

假设文档中有:

{
  "name": "Java Elasticsearch 工程"
}

text 字段进行标准分析后,可能得到类似:

java
elasticsearch
工程

具体分词结果取决于 analyzer。对 name.keyword,整个字符串通常作为一个不可拆分的词:

Java Elasticsearch 工程

因此:

{
  "match": {
    "name": "Elasticsearch 工程"
  }
}

是在分析后进行全文匹配;而:

{
  "term": {
    "name.keyword": "Java Elasticsearch 工程"
  }
}

要求整个 keyword 值精确相等。

常见错误是用 term 查询普通 text 字段,并期待它执行分词后的全文检索。term 不会像 match 那样按字段 analyzer 分析输入,它适合已知的精确词项。反过来,用 match 查询状态码、用户 ID 或枚举值,也会引入不必要的分析语义。

3.3 数值、日期和禁用索引

数值字段应使用数值类型:

{
  "price": {
    "type": "double"
  },
  "stock": {
    "type": "integer"
  }
}

不要把价格保存成字符串再依赖字符串排序。字符串排序中 "100" 可能排在 "20" 之前,因为比较的是字符,而不是数值。

如果字段只需要出现在 _source 中、不需要搜索,可以设置:

{
  "internalNote": {
    "type": "keyword",
    "index": false
  }
}

这会减少索引结构,但该字段不能再通过普通查询条件搜索。Mapping 的每个选择都是能力和存储成本之间的取舍。

3.4 使用 Java API Client 创建 Mapping

定义一个 Java record:

public record Product(
        String id,
        String name,
        String category,
        double price
) {
}

创建索引:

import co.elastic.clients.elasticsearch.ElasticsearchClient;
import co.elastic.clients.elasticsearch.indices.CreateIndexResponse;

public final class ProductIndex {

    public static final String INDEX = "products-v1";

    public static CreateIndexResponse create(ElasticsearchClient client)
            throws Exception {

        return client.indices().create(c -> c
                .index(INDEX)
                .mappings(m -> m
                        .properties("name",
                                p -> p.text(t -> t
                                        .fields("keyword",
                                                k -> k.keyword(kb -> kb))))
                        .properties("category",
                                p -> p.keyword(k -> k))
                        .properties("price",
                                p -> p.double_(d -> d))
                )
        );
    }
}

前提条件是客户端能访问 Elasticsearch,并且 products-v1 尚不存在。成功响应通常包含:

{
  "acknowledged": true,
  "shards_acknowledged": true,
  "index": "products-v1"
}

这两个确认字段含义不同:

  • acknowledged:集群状态已接受索引创建请求。
  • shards_acknowledged:分片启动确认也完成。

创建成功后不能随意把已有字段从 text 改成 keyword,也不能把 double 直接改成 integer。生产中的典型迁移流程是:

  1. 创建新索引 products-v2
  2. products-v2 定义正确 Mapping。
  3. 使用 _reindex 或应用程序重新写入数据。
  4. 验证文档数量、抽样查询和聚合结果。
  5. 通过别名把读写流量切换到新索引。
  6. 确认无回滚需求后删除旧索引。

直接删除旧索引是不可逆的破坏性操作,必须先确认快照、恢复方案和别名切换状态。

3.5 动态 Mapping 的边界

默认动态 Mapping 能让简单 JSON 快速可用,但它也会把输入数据结构变化自动转化为字段变化。例如某次写入:

{
  "value": 10
}

随后另一条文档写入:

{
  "value": "unknown"
}

同一字段可能产生类型冲突。更严重的是,用户自定义字段不断增加会造成 mapping explosion,影响集群状态和资源使用。

工程上通常采用:

  • 核心字段显式 Mapping。
  • 外部输入字段使用白名单或显式对象结构。
  • 对不需要搜索的对象使用 enabled: false 或关闭索引。
  • 在测试环境用真实样例验证日期、数字和空值行为。

四、查询:从查询语义到 Java 请求

4.1 查询、过滤和评分

Elasticsearch 查询通常由两类逻辑组成:

  • Query context:关注相关性评分,例如 match
  • Filter context:关注是否满足条件,例如 term、范围条件,通常不需要评分。

例如“名称中包含 Elasticsearch,同时分类为 book,价格不低于 100”可以拆成:

全文条件:name match "Elasticsearch"
精确条件:category = "book"
范围条件:price >= 100

其中名称影响相关性,分类和价格只是筛选条件。用 bool.filter 表达后两个条件,比把所有条件都放在 must 中更准确地表达意图。

4.2 Java 查询示例

import co.elastic.clients.elasticsearch.ElasticsearchClient;
import co.elastic.clients.elasticsearch._types.SortOrder;
import co.elastic.clients.elasticsearch.core.SearchResponse;

public final class ProductSearch {

    public static SearchResponse<Product> search(
            ElasticsearchClient client,
            String text,
            String category,
            double minimumPrice
    ) throws Exception {

        return client.search(s -> s
                        .index(ProductIndex.INDEX)
                        .query(q -> q
                                .bool(b -> b
                                        .must(m -> m
                                                .match(mm -> mm
                                                        .field("name")
                                                        .query(text)))
                                        .filter(f -> f
                                                .term(t -> t
                                                        .field("category")
                                                        .value(category)))
                                        .filter(f -> f
                                                .range(r -> r
                                                        .number(n -> n
                                                                .field("price")
                                                                .gte(minimumPrice))))
                                )
                        )
                        .sort(so -> so
                                .field(f -> f
                                        .field("price")
                                        .order(SortOrder.Asc)))
                        .size(20),
                Product.class
        );
    }
}

读取结果:

SearchResponse<Product> response =
        ProductSearch.search(client, "Elasticsearch", "book", 100.0);

System.out.println("total hits = " + response.hits().total().value());

for (var hit : response.hits().hits()) {
    Product product = hit.source();
    if (product != null) {
        System.out.printf(
                "%s: %s, %.2f%n",
                hit.id(),
                product.name(),
                product.price()
        );
    }
}

这里的 Product.class 告诉 Java 客户端把 _source 反序列化为 Product。如果 Elasticsearch 文档缺少字段、字段类型和 Java 类型不匹配,解析阶段可能失败;这不是“没有命中”,而是客户端无法把响应转换成目标 Java 对象。

4.3 matchtermmatch_phrase

match

{
  "match": {
    "name": "Java Elasticsearch"
  }
}

输入会根据字段 analyzer 处理,适合自然语言搜索。多个词如何组合还受到 operator、minimum_should_match 等参数影响。

term

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

适合状态、枚举、ID 和 keyword 字段。若 category 映射为 text,这个查询通常不会得到你期待的结果。

match_phrase

{
  "match_phrase": {
    "name": "Java Elasticsearch"
  }
}

要求分析后的词项按短语位置匹配,适合搜索连续短语,但仍然依赖 analyzer;它不是对原始字符串做简单的 equals

4.4 分页不是只有 fromsize

小结果集可以使用:

{
  "from": 0,
  "size": 20
}

但是深分页会带来成本。对于第 p 页、每页 s 条数据,协调节点需要处理大致前 p × s 条候选结果,再返回当前页。from 越大,内存和排序压力越大。

连续翻页更适合 search_after

  1. 第一次查询按稳定排序返回结果。
  2. 取最后一条 hit 的 sort 值。
  3. 下一次请求把该值放入 search_after
  4. 重复直到没有结果。

如果需要长时间一致的遍历视图,可以结合 Point-in-Time。否则在分页期间发生新增或删除,结果边界可能发生变化。

批量导出或重建数据时,还可以使用 Scroll,但 Scroll 更偏向长时间遍历,不应把它当作面向用户的普通分页机制。


五、Bulk:批处理,不是事务

5.1 Bulk 请求的结构

Bulk API 在一个 HTTP 请求中携带多个操作。逻辑结构类似:

action metadata
source document

action metadata
source document

例如:

{ "index": { "_index": "products-v1", "_id": "p-1" } }
{ "name": "Java", "category": "book", "price": 99.0 }
{ "delete": { "_index": "products-v1", "_id": "p-2" } }

Bulk 减少了 HTTP 往返次数,但不会把所有操作变成一个原子事务。第一条成功、第二条失败、第三条成功是合法结果。

Bulk 中常见动作:

  • index:写入文档,指定同一 ID 时通常会覆盖已有文档。
  • create:只允许创建,不允许覆盖;ID 已存在时失败。
  • update:对已有文档做部分更新或脚本更新。
  • delete:删除文档。

5.2 Java Bulk 示例

import co.elastic.clients.elasticsearch.ElasticsearchClient;
import co.elastic.clients.elasticsearch.core.BulkRequest;
import co.elastic.clients.elasticsearch.core.BulkResponse;

import java.util.List;

public final class ProductBulk {

    public static BulkResponse index(
            ElasticsearchClient client,
            List<Product> products
    ) throws Exception {

        BulkRequest.Builder request = new BulkRequest.Builder();

        for (Product product : products) {
            request.operations(op -> op.index(i -> i
                    .index(ProductIndex.INDEX)
                    .id(product.id())
                    .document(product)
            ));
        }

        return client.bulk(request.build());
    }
}

调用后必须逐项检查:

BulkResponse response = ProductBulk.index(client, products);

if (response.errors()) {
    for (var item : response.items()) {
        if (item.error() != null) {
            System.err.printf(
                    "bulk item failed: id=%s, status=%s, type=%s, reason=%s%n",
                    item.id(),
                    item.status(),
                    item.error().type(),
                    item.error().reason()
            );
        }
    }
}

response.errors() 只是一个快速判断。即使 HTTP 状态为 200,也必须检查每个 item 的 error()。把“HTTP 200”误认为“所有文档成功”是 Bulk 最常见的错误之一。

5.3 Bulk 的批量大小

Bulk 批次不应只按文档条数决定,因为文档大小可能差异很大。更合理的边界至少考虑:

  • 文档数量。
  • 序列化后的请求体大小。
  • 单批处理耗时。
  • Elasticsearch 节点的写入队列和内存压力。

批次过小,HTTP 开销高;批次过大,单次失败影响范围大,重试请求也更重。具体上限应通过压测和监控确认,不能使用脱离文档大小、分片数和集群配置的固定“最佳数字”。

5.4 refresh 的取舍

Bulk 请求支持 refresh 行为,例如:

client.bulk(b -> b
        .refresh(co.elastic.clients.elasticsearch._types.Refresh.WaitFor)
        // operations...
);

常见语义包括:

  • false:不主动等待 refresh,吞吐通常更好,搜索可见性延迟更高。
  • wait_for:等待下一次 refresh 后返回,适合写后读要求较强的场景。
  • true:主动 refresh,可能增加刷新压力,不应无条件用于高吞吐写入。

如果业务要求“写入响应返回后立刻能搜索到”,需要明确采用 wait_for 或显式 refresh,并接受相应成本。不要在每条文档写入后都强制 refresh。


六、重试:先判断失败是否值得重试

重试不是“捕获异常后再执行一次”。一次正确的重试决策至少需要回答:

  1. 失败发生在请求发送前、发送中还是服务端处理后?
  2. 服务端是否已经执行了这次写入?
  3. 这次操作是否幂等?
  4. 错误是临时资源不足,还是永久数据错误?
  5. 重试是否会放大流量和故障?

6.1 典型错误分类

通常可重试的错误

  • 连接建立失败。
  • 连接被重置。
  • 读取超时。
  • HTTP 429 Too Many Requests
  • HTTP 502 Bad Gateway
  • HTTP 503 Service Unavailable
  • HTTP 504 Gateway Timeout

这些错误可能表示网络抖动、节点暂时不可用或线程池/写入队列繁忙。

通常不应自动重试的错误

  • HTTP 400:请求结构或参数错误。
  • Mapping 冲突。
  • JSON 序列化或反序列化错误。
  • 字段类型不匹配。
  • 查询语法错误。
  • 认证失败。
  • 索引不存在且应用没有明确的创建策略。

重复提交这些错误不会修复根因,反而会增加日志和集群压力。

需要业务判断的错误

  • 409 Conflict:可能是版本冲突,也可能是并发更新导致的失败。
  • 写入超时:服务端可能没有及时返回,但操作可能已经成功。
  • Bulk 中的单个 429:可以只重试该 item。
  • create 的 ID 已存在:这通常是业务冲突,不是临时故障。

6.2 幂等性是重试的前提

设一次写操作为 W,重试次数为 n。只有当重复执行满足:

W(W(state)) = W(state)

时,重复请求才不会继续改变最终状态,这就是幂等性。

使用固定文档 ID 的完整覆盖写入通常容易做到幂等:

PUT /products-v1/p-1

无论请求执行一次还是因超时执行两次,最终文档内容都应是同一份。

相反,自动生成 ID 的创建请求可能不幂等:

POST /products-v1/_doc

第一次请求可能已经在服务端成功,只是响应返回途中断开;客户端再次提交会生成第二个文档。

因此工程上常使用业务主键作为 Elasticsearch _id。这既减少重复文档,也让超时后的重试具备更清晰的最终状态。

6.3 指数退避和抖动

k 次重试的等待时间可以定义为:

delay(k) = min(cap, base × 2^k) + random(0, jitter)

其中:

  • base:基础等待时间。
  • k:从 0 开始的重试序号。
  • cap:单次等待上限。
  • random(0, jitter):随机抖动。
  • min:防止等待时间无限增长。

例如 base = 100 mscap = 5 s 时,理论退避序列可能是:

100 ms
200 ms
400 ms
800 ms
1600 ms
3200 ms
5000 ms

如果所有客户端在同一时刻失败并按完全相同的固定时间重试,会形成惊群效应。抖动让重试请求分散到不同时间点。

6.4 单条写入重试示例

import co.elastic.clients.elasticsearch.ElasticsearchClient;
import co.elastic.clients.elasticsearch.core.IndexResponse;

import java.io.IOException;
import java.util.concurrent.ThreadLocalRandom;

public final class RetryableIndexer {

    private static final int MAX_RETRIES = 4;
    private static final long BASE_DELAY_MILLIS = 100;
    private static final long MAX_DELAY_MILLIS = 5_000;

    public static IndexResponse indexWithRetry(
            ElasticsearchClient client,
            Product product
    ) throws Exception {

        for (int attempt = 0; ; attempt++) {
            try {
                return client.index(i -> i
                        .index(ProductIndex.INDEX)
                        .id(product.id())
                        .document(product)
                );
            } catch (IOException e) {
                if (attempt >= MAX_RETRIES) {
                    throw e;
                }

                long delay = backoffMillis(attempt);
                Thread.sleep(delay);
            }
        }
    }

    private static long backoffMillis(int attempt) {
        long exponential = Math.min(
                MAX_DELAY_MILLIS,
                BASE_DELAY_MILLIS * (1L << attempt)
        );

        long jitter = ThreadLocalRandom.current()
                .nextLong(0, BASE_DELAY_MILLIS + 1);

        return Math.min(MAX_DELAY_MILLIS, exponential + jitter);
    }
}

这个简化示例只捕获 IOException,它不能区分所有 HTTP 状态,因此生产实现应根据客户端异常对象中的 HTTP 状态和 Elasticsearch 错误类型决定是否重试。不能仅仅因为出现了 IOException 就无限重试。

Thread.sleep 适合展示算法,不一定适合高并发异步服务。在线程池中大量 sleep 会占用工作线程;实际服务可使用异步调度、限流队列或专门的重试组件。无论使用哪种实现,都必须设置最大尝试次数和总超时时间。


七、Bulk 重试必须只重试失败项

7.1 为什么不能重试整个 Bulk

假设一个 Bulk 包含四项:

item 0: 成功
item 1: 429
item 2: Mapping 错误
item 3: 成功

如果把整个 Bulk 原样重试:

  • item 0 和 item 3 会被重复发送。
  • item 2 仍然会失败。
  • item 1 可能恢复,也可能继续失败。
  • 若使用非幂等操作,重复发送可能产生重复数据。

正确流程是:

  1. 保留原始操作与业务 ID 的对应关系。
  2. 检查每个 item 的状态。
  3. 只提取可重试的失败项。
  4. 丢弃永久失败项并记录原因。
  5. 对失败项重新组成更小的 Bulk。
  6. 达到上限后进入死信队列或人工处理。

7.2 Bulk 结果的状态模型

可以把每个 item 看作以下状态:

待发送
  ↓
已发送
  ├── 成功 → 完成
  ├── 永久失败 → 死信/告警
  └── 临时失败 → 等待退避 → 重试

超时是特殊情况:

已发送
  ↓
客户端未收到响应
  ├── 服务器实际未执行 → 重试
  └── 服务器已经执行 → 重试可能重复

因此,超时后的写入尤其依赖固定 _id 和幂等操作。对于带有自增副作用的脚本更新、计数操作,必须设计业务级去重或采用版本条件,否则客户端无法仅凭超时判断服务端是否已经成功。

7.3 伪代码

List<Operation> pending = originalOperations;

for (int attempt = 0; attempt <= maxRetries && !pending.isEmpty(); attempt++) {
    BulkResponse response = send(pending);

    List<Operation> next = new ArrayList<>();

    for (int i = 0; i < pending.size(); i++) {
        Operation original = pending.get(i);
        BulkResponseItem item = response.items().get(i);

        if (item.error() == null) {
            markSuccess(original);
        } else if (isRetryable(item.status(), item.error())) {
            next.add(original);
        } else {
            moveToDeadLetter(original, item.error());
        }
    }

    pending = next;

    if (!pending.isEmpty()) {
        sleepWithExponentialBackoff(attempt);
    }
}

for (Operation operation : pending) {
    moveToDeadLetter(operation, "retry limit exceeded");
}

关键前提是:response.items() 与请求操作按顺序对应。应用不得在中间随意排序或丢失原始操作索引,否则会把一个文档的错误错误地归因给另一个文档。


八、并发更新:重试不能替代版本控制

两个 Java 应用同时修改同一个文档时,仅靠固定 ID 只能避免重复 ID,不能避免“后写覆盖先写”。

例如初始文档:

{
  "_id": "p-1",
  "_source": {
    "price": 100,
    "stock": 10
  }
}

两个请求同时读取:

请求 A 读取 stock=10,准备写 stock=9
请求 B 读取 stock=10,准备写 stock=8

如果都执行整文档覆盖,最终结果取决于到达顺序,其中一个修改会丢失。

需要区分两种策略:

覆盖式写入

适用于事件最终状态明确的场景:

以事件中的完整对象覆盖 Elasticsearch 文档

要求事件中带有稳定版本或时间序列,避免旧事件晚到后覆盖新状态。

乐观并发控制

读取文档时记录 _seq_no_primary_term,更新时带上条件:

只有文档仍处于读取时的版本,更新才成功

如果版本已变化,服务端返回冲突。此时应用可以重新读取、合并业务字段,再决定是否重试。不能把所有 409 无脑重试,因为这可能无限覆盖其他请求的更新。

对于计数器、库存扣减等场景,通常还要使用服务端脚本或业务事务设计;简单的“读取后在 Java 中加一再写回”存在竞态。


九、端到端示例:创建索引、写入、刷新可见性和查询

下面的程序串联前面的流程:

import co.elastic.clients.elasticsearch.ElasticsearchClient;
import co.elastic.clients.elasticsearch._types.Refresh;
import co.elastic.clients.elasticsearch.core.BulkResponse;
import co.elastic.clients.elasticsearch.core.SearchResponse;
import co.elastic.clients.elasticsearch.indices.ExistsResponse;

import java.util.List;

public final class ProductApplication {

    public static void main(String[] args) throws Exception {
        try (ClientHolder holder = new ClientHolder()) {
            ElasticsearchClient client = holder.client();

            ExistsResponse exists = client.indices()
                    .exists(e -> e.index(ProductIndex.INDEX));

            if (!exists.value()) {
                ProductIndex.create(client);
            }

            List<Product> products = List.of(
                    new Product("p-1", "Java Elasticsearch 工程", "book", 129.0),
                    new Product("p-2", "Distributed Systems", "book", 159.0),
                    new Product("p-3", "Mechanical Keyboard", "hardware", 299.0)
            );

            BulkResponse bulkResponse = client.bulk(b -> {
                for (Product product : products) {
                    b.operations(op -> op.index(i -> i
                            .index(ProductIndex.INDEX)
                            .id(product.id())
                            .document(product)
                    ));
                }
                return b.refresh(Refresh.WaitFor);
            });

            if (bulkResponse.errors()) {
                bulkResponse.items().forEach(item -> {
                    if (item.error() != null) {
                        System.err.println(item.error().reason());
                    }
                });
                throw new IllegalStateException("bulk contains failures");
            }

            SearchResponse<Product> searchResponse =
                    ProductSearch.search(client, "Elasticsearch", "book", 100.0);

            searchResponse.hits().hits().forEach(hit ->
                    System.out.println(hit.id() + " -> " + hit.source()));
        }
    }
}

每一步的因果关系是:

  1. 先检查索引是否存在,避免重复创建。
  2. 索引不存在时创建显式 Mapping。
  3. Bulk 使用业务 ID,重复执行时不会生成新的随机文档 ID。
  4. 设置 wait_for,让示例后续查询等待 refresh。
  5. 检查 Bulk 的逐项错误,而不是只看请求是否返回。
  6. 查询使用 name 的全文语义、category 的精确语义和 price 的数值范围语义。

这个程序仍不是完整生产实现,因为它没有处理认证、TLS、连接池参数、Bulk 分批、失败项重试、死信、指标和优雅停机。这些属于生命周期和故障处理的一部分,而不是客户端 API 调用本身能够自动完成的功能。


十、常见失败表现与诊断路径

10.1 “写入成功但查询不到”

优先检查:

  1. 查询是否打到了正确的索引或别名。
  2. 查询字段是否写成了 name.keywordname
  3. term 是否错误地用于 text 字段。
  4. 是否只是 refresh 尚未发生。
  5. _source 中是否真的包含预期值。
  6. 是否有路由参数导致查询没有覆盖目标分片。

如果只是短暂不可见,应检查 refresh 语义;如果长时间不可见,应检查索引名、Mapping、查询 DSL 和写入响应。

10.2 mapper_parsing_exception

这通常表示文档值与 Mapping 不兼容,例如:

Mapping: price 是 double
文档:    price = "unknown"

重试不会解决类型错误。正确处理是修正输入、修正 Mapping 并进行索引迁移,或者把不稳定字段设计成字符串/未索引对象。

10.3 Bulk 返回 HTTP 200 但有失败项

这是正常的 Bulk 处理模型。诊断时输出:

  • 文档业务 ID。
  • item 的 HTTP 状态。
  • Elasticsearch error type。
  • reason。
  • 原始操作类型。
  • 重试次数。

如果只记录“Bulk 成功”,会丢失部分数据失败的事实。

10.4 大量 429

429 通常意味着节点暂时无法接收更多工作,可能与写入线程池、分片热点、批次过大或并发生产者过多有关。增加重试并不能无限提升容量;如果生产速率持续高于消费速率,队列最终仍会堆积。

应同时观察:

  • Bulk 请求耗时和大小。
  • 429 比例。
  • 各节点写入线程池队列。
  • 分片分布和热点。
  • JVM 堆、GC 和磁盘 I/O。
  • 应用端待处理队列长度。

退避、限流和降低并发是缓解手段;根本问题可能需要调整索引分片设计或写入架构。

10.5 查询结果相关性异常

可能原因包括:

  • text 字段使用了不合适的 analyzer。
  • 本应精确匹配的字段被映射为 text
  • 应使用 match_phrase 却使用了普通 match
  • 过滤条件错误地放在影响评分的 query context。
  • 多字段查询没有为不同字段设置合理权重。

诊断时先用最小 DSL 验证单字段,再逐步加入 bool、过滤、排序和聚合。不要直接从复杂 Java Builder 代码猜测最终 JSON,可以记录请求或在 Kibana/REST API 中独立验证 DSL。


十一、与 JDBC 和 JPA 的边界

JDBC 面向关系数据库连接、SQL 执行、结果集和事务;Jakarta Persistence 面向实体、持久化上下文、关系映射和对象生命周期。Elasticsearch Java 客户端则面向 REST API、索引文档、搜索 DSL 和分布式搜索结果。

因此,以下迁移思路通常是不成立的:

数据库表      → Elasticsearch 索引
数据库列      → 直接等价的搜索字段
JPA save      → Elasticsearch index
数据库事务    → 一个 Bulk 请求
SQL WHERE     → 任意 Query DSL

可以建立概念上的对应,但不能忽略语义差异:

  • JDBC/JPA 依赖数据库事务保证多条更新的一致性;Bulk 不提供同等的跨文档原子性。
  • JPA 实体关系可以通过外键和 Join 表达;Elasticsearch 通常通过文档冗余、嵌套对象或应用侧组合查询。
  • 数据库唯一约束可以拒绝重复值;Elasticsearch 使用固定 _id 能防止同一 ID 的重复文档,但不能自动为任意字段提供关系数据库式唯一约束。
  • 数据库精确字符串比较和 Elasticsearch 分词全文搜索不是同一种语义。

常见架构是:关系数据库作为事务事实源,Elasticsearch 作为搜索投影。应用先在事务数据库中完成业务变更,再通过消息、CDC 或可靠任务把数据同步到 Elasticsearch。此时必须接受并处理最终一致性、重复事件、乱序事件和索引重建。


十二、生产取舍的核心原则

客户端

  • 全局复用客户端和连接池。
  • 配置 HTTPS 和认证,不使用本地明文配置。
  • 设置连接、读取和总请求超时。
  • 记录请求耗时、状态码和异常类型。
  • 不要在每个业务方法中创建和关闭客户端。

Mapping

  • 让核心字段显式定义类型。
  • 全文字段使用 text,精确字段使用 keyword
  • 数值和日期使用真实数值、日期类型。
  • 预计变化的 Mapping 使用版本化索引和别名迁移。
  • 控制动态字段数量,防止 Mapping 无限制增长。

查询

  • match 表达全文语义。
  • term 表达精确词项语义。
  • filter 表达不需要评分的筛选条件。
  • 深分页使用 search_after 或适合场景的遍历机制。
  • 对用户输入设置查询超时、结果大小和字段范围。

Bulk

  • 按文档数量和请求体大小共同切分。
  • 把 Bulk 当作批处理,不当作事务。
  • 检查每一个 item。
  • 失败时只重试可重试项。
  • 为重试和死信保留业务 ID、原始数据和错误原因。

重试

  • 只对临时错误重试。
  • 使用有限次数、指数退避和随机抖动。
  • 优先使用固定 _id 保证写入幂等。
  • 对超时假设“服务端可能已经执行”,不要盲目追加操作。
  • 对版本冲突使用乐观并发控制,而不是无限重试。

Java Elasticsearch 工程的可靠性,不在于把 API 调用封装得多复杂,而在于让 Mapping、查询语义、批量结果、并发版本和失败恢复形成一致的数据流。只要把“请求成功”“文档成功”“搜索可见”和“业务最终完成”区分开,客户端、Bulk 与重试的设计就会从偶然可用变成可验证、可诊断的工程实现。


系列导航与关联阅读

官方资料

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