Java 基础体系 · 第 88/100 篇。示例统一以 Java 25 LTS 为语言和 JVM 基线;框架示例使用与其兼容的现代稳定版本。
Java Elasticsearch 工程:客户端、Mapping、查询、Bulk 和重试
Elasticsearch 是一个分布式搜索与分析引擎。Java 工程通常通过官方 Elasticsearch Java API Client 访问它,使用 Mapping 描述索引中的字段结构,使用 Query DSL 表达查询条件,通过 Bulk 批量写入,最后根据错误类型实施有限且可观测的重试。
这几个概念并不是相互独立的:
- Mapping 决定 JSON 字段如何被索引,以及查询时采用什么语义。
- 查询请求依赖 Mapping;
text、keyword、数值和日期字段的查询行为不同。 - 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: 类型化响应或异常
这里有几个容易混淆的边界:
-
Java 客户端不是 Elasticsearch 服务端。
客户端只负责构造 HTTP 请求、序列化 JSON、解析响应和暴露 Java API。 -
索引不是数据库表的完全等价物。
Elasticsearch 的索引包含分片、副本、倒排索引、Doc Values 等搜索结构;它不提供关系数据库那种通用事务、跨表 Join 和外键约束。 -
写入成功不必然意味着搜索立即可见。
Elasticsearch 通过 refresh 让最近写入的内容对搜索可见。默认 refresh 周期下,写入响应返回后,搜索可能仍暂时查不到文档。 -
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:提供index、search、bulk、indices().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);
}
}
这里的生命周期是:
- 创建底层
RestClient。 - 用它创建
RestClientTransport。 - 用 transport 创建
ElasticsearchClient。 - 应用关闭时关闭 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 对索引字段的类型和索引方式描述。它至少回答三个问题:
- 字段是什么类型?
- 字段是否建立索引?
- 字段如何被分析、排序、聚合或存储?
例如:
{
"properties": {
"name": {
"type": "text",
"fields": {
"keyword": {
"type": "keyword"
}
}
},
"price": {
"type": "double"
},
"category": {
"type": "keyword"
},
"createdAt": {
"type": "date"
}
}
}
这里的 name 同时具有两个视图:
name是text,适合全文搜索。name.keyword是keyword,适合精确匹配、排序和聚合。
3.2 text 和 keyword 的区别
假设文档中有:
{
"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。生产中的典型迁移流程是:
- 创建新索引
products-v2。 - 为
products-v2定义正确 Mapping。 - 使用
_reindex或应用程序重新写入数据。 - 验证文档数量、抽样查询和聚合结果。
- 通过别名把读写流量切换到新索引。
- 确认无回滚需求后删除旧索引。
直接删除旧索引是不可逆的破坏性操作,必须先确认快照、恢复方案和别名切换状态。
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 match、term 和 match_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 分页不是只有 from 和 size
小结果集可以使用:
{
"from": 0,
"size": 20
}
但是深分页会带来成本。对于第 p 页、每页 s 条数据,协调节点需要处理大致前 p × s 条候选结果,再返回当前页。from 越大,内存和排序压力越大。
连续翻页更适合 search_after:
- 第一次查询按稳定排序返回结果。
- 取最后一条 hit 的 sort 值。
- 下一次请求把该值放入
search_after。 - 重复直到没有结果。
如果需要长时间一致的遍历视图,可以结合 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。
六、重试:先判断失败是否值得重试
重试不是“捕获异常后再执行一次”。一次正确的重试决策至少需要回答:
- 失败发生在请求发送前、发送中还是服务端处理后?
- 服务端是否已经执行了这次写入?
- 这次操作是否幂等?
- 错误是临时资源不足,还是永久数据错误?
- 重试是否会放大流量和故障?
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 ms、cap = 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 可能恢复,也可能继续失败。
- 若使用非幂等操作,重复发送可能产生重复数据。
正确流程是:
- 保留原始操作与业务 ID 的对应关系。
- 检查每个 item 的状态。
- 只提取可重试的失败项。
- 丢弃永久失败项并记录原因。
- 对失败项重新组成更小的 Bulk。
- 达到上限后进入死信队列或人工处理。
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()));
}
}
}
每一步的因果关系是:
- 先检查索引是否存在,避免重复创建。
- 索引不存在时创建显式 Mapping。
- Bulk 使用业务 ID,重复执行时不会生成新的随机文档 ID。
- 设置
wait_for,让示例后续查询等待 refresh。 - 检查 Bulk 的逐项错误,而不是只看请求是否返回。
- 查询使用
name的全文语义、category的精确语义和price的数值范围语义。
这个程序仍不是完整生产实现,因为它没有处理认证、TLS、连接池参数、Bulk 分批、失败项重试、死信、指标和优雅停机。这些属于生命周期和故障处理的一部分,而不是客户端 API 调用本身能够自动完成的功能。
十、常见失败表现与诊断路径
10.1 “写入成功但查询不到”
优先检查:
- 查询是否打到了正确的索引或别名。
- 查询字段是否写成了
name.keyword或name。 term是否错误地用于text字段。- 是否只是 refresh 尚未发生。
_source中是否真的包含预期值。- 是否有路由参数导致查询没有覆盖目标分片。
如果只是短暂不可见,应检查 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 完整学习路线:从 Java 25 语言与 JVM 到 Spring、微服务和生产交付
- 上一篇:Java Redis 工程:Lettuce、连接、序列化、缓存、锁和 Streams
- 下一篇:Java MongoDB 工程:文档建模、驱动、事务、索引和变更流
官方资料
本文依据 Java、Spring 与相关项目官方文档重新梳理;正文、示例与生产清单由 WR BLOG 编写。

评论
0 条讨论