这是《Elasticsearch 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 第 4 章已经见过基础 API。本章系统讲解文档写入语义、自动 ID 与手动 ID、创建、替换、局部更新、批量保存、乐观并发、路由、刷新和失败处理。
5.1 文档写入 API 对比
| API | 行为 | 适用 |
|---|---|---|
POST /{index}/_doc |
自动生成 ID | 日志、事件、无需业务主键 |
PUT /{index}/_doc/{id} |
不存在创建,存在整体替换 | 数据库同步、幂等写入 |
PUT /{index}/_create/{id} |
只创建,存在则报错 | 防止覆盖 |
POST /{index}/_update/{id} |
合并局部字段或执行脚本 | 少量字段更新 |
POST /_bulk |
批量执行 index/create/update/delete | 高吞吐写入 |
5.2 自动 ID 与业务 ID
自动 ID:
POST /events/_doc
{
"level": "INFO",
"message": "user login",
"timestamp": "2026-08-25T10:00:00Z"
}
优点是写入路径简单,适合只追加的日志和事件;缺点是不直观、无法天然幂等。
业务 ID:
PUT /products/_doc/10001
{
"product_id": 10001,
"title": "轻薄笔记本电脑",
"price": 6999,
"status": "ON_SALE"
}
第一次返回 result=created、_version=1;再次写入同 ID 返回 result=updated、_version=2。这里的 updated 表示替换了旧文档,即使 JSON 完全相同也会生成新版本。
商品、订单、用户等数据库实体同步通常使用业务 ID;日志、埋点、审计事件通常使用自动 ID。
5.3 创建语义
只允许创建:
PUT /products/_create/10002
{
"product_id": 10002,
"title": "游戏手机",
"price": 4999
}
已存在时返回 409 Conflict。等价写法:
PUT /products/_doc/10002?op_type=create
{
"product_id": 10002,
"title": "游戏手机",
"price": 4999
}
适合初始化、迁移和首次导入。
5.4 局部更新
使用 doc 覆盖部分字段:
POST /products/_update/10001
{
"doc": {
"price": 6899,
"tags": ["轻薄", "新品"]
}
}
删除字段:
POST /products/_update/10001
{
"script": {
"source": "ctx._source.remove('tags')"
}
}
增量更新:
POST /products/_update/10001
{
"script": {
"source": "ctx._source.sales += params.count",
"params": { "count": 1 }
}
}
数组去重添加:
POST /products/_update/10001
{
"script": {
"source": """
if (ctx._source.tags == null) {
ctx._source.tags = new ArrayList();
}
if (!ctx._source.tags.contains(params.tag)) {
ctx._source.tags.add(params.tag);
}
""",
"params": { "tag": "轻薄" }
}
}
doc_as_upsert:
POST /products/_update/10003
{
"doc": {
"product_id": 10003,
"title": "机械键盘",
"price": 499,
"status": "ON_SALE"
},
"doc_as_upsert": true
}
如果文档不存在则插入,存在则合并更新。
5.5 更新的内部流程
Elasticsearch 文档不可变,局部更新流程近似为:
定位主分片
-> 读取现有 _source
-> 合并修改
-> 标记旧文档删除
-> 写入新文档
-> 写 translog
-> 同步副本
-> 返回客户端
因此 _update 不是原地修改某个字段,而是重新索引整个文档。
设计约束:
- 大文档频繁更新代价高;
- nested 对象越多更新成本越高;
- 高频计数和状态变更应先聚合;
- 并发更新会出现版本冲突;
- 原始字段最好能在同步层完整重建。
5.6 乐观并发控制
ES 使用 _seq_no 和 _primary_term 实现乐观并发。
先读取文档并记录:
_seq_no = 42
_primary_term = 1
带条件更新:
POST /products/_update/10001?if_seq_no=42&if_primary_term=1
{
"doc": {
"price": 6799
}
}
如果中间已有其他更新,本次请求返回 409 Conflict。
简单增量更新可以自动重试:
POST /products/_update/10001?retry_on_conflict=3
{
"script": {
"source": "ctx._source.sales += params.count",
"params": { "count": 1 }
}
}
如果上游数据库有版本号,可使用外部版本:
PUT /products/_doc/10001?version=10&version_type=external
{
"product_id": 10001,
"title": "轻薄笔记本电脑",
"price": 6799
}
只有新版本号大于当前版本号才写入,适合数据库到 ES 的单向同步。
5.7 自定义路由
默认路由值是 _id:
shard = hash(routing) % number_of_primary_shards
自定义路由:
PUT /orders/_doc/O10001?routing=U8888
{
"order_id": "O10001",
"user_id": "U8888",
"amount": 199
}
读取也必须携带:
GET /orders/_doc/O10001?routing=U8888
优点:
- 同一用户文档落在同一分片;
- 按用户查询时只命中部分分片;
- 可优化租户隔离。
风险:
- 路由值分布不均导致数据倾斜;
- 忘记 routing 读取不到文档;
- 大租户可能打爆单个分片;
- 查询灵活性下降。
多租户系统更常见的是按租户建索引或使用权限过滤,而不是过度依赖自定义路由。
5.8 刷新策略
写入时可以指定 refresh:
PUT /products/_doc/10004?refresh=true
{
"title": "无线鼠标"
}
可选值:
| 值 | 行为 | 适用 |
|---|---|---|
true |
请求后立即 refresh | 测试、低频关键数据 |
wait_for |
等待下一次 refresh | 需要可见性,又不想强制刷新 |
false |
默认,不触发 | 大多数批量写入 |
示例:
PUT /products/_doc/10005?refresh=wait_for
{
"title": "显示器"
}
高吞吐批量写入不要每条都 refresh=true,否则会产生大量小 segment,导致性能下降。
5.9 Bulk 与失败处理
批量写入:
POST /_bulk
{"index": {"_index": "products", "_id": "10001"}}
{"product_id": 10001, "title": "轻薄笔记本", "price": 6999, "status": "ON_SALE"}
{"index": {"_index": "products", "_id": "10002"}}
{"product_id": 10002, "title": "游戏手机", "price": 4999, "status": "ON_SALE"}
{"index": {"_index": "products", "_id": "10003"}}
{"product_id": 10003, "title": "机械键盘", "price": 499, "status": "OFF_SALE"}
返回中必须检查 errors 和每个 item.status:
{
"took": 30,
"errors": true,
"items": [
{
"index": {
"status": 429,
"error": {
"type": "es_rejected_execution_exception",
"reason": "rejected execution of coordinating operation"
}
}
}
]
}
Java 侧示意:
BulkResponse response = client.bulk(request);
if (response.errors()) {
for (BulkItemResponse item : response.items()) {
if (item.isFailed()) {
log.warn("bulk failed index={} id={} reason={}",
item.index(), item.id(), item.failure().message());
failedBuffer.add(item);
}
}
}
HTTP 200 只表示请求被处理,不代表所有操作成功。
5.10 删除文档
按 ID 删除:
DELETE /products/_doc/10003
按查询删除:
POST /products/_delete_by_query
{
"query": {
"term": {
"status": "OFF_SALE"
}
}
}
生产建议:
- 时间序列数据优先整索引删除;
_delete_by_query会产生标记删除和合并压力;- 低峰执行;
- 控制批次和并发;
- 任务失败要记录断点。
5.11 同步文档元数据实践
推荐携带:
{
"id": 10001,
"version": 128,
"updated_at": "2026-08-25T10:00:00Z",
"synced_at": "2026-08-25T10:00:01Z",
"source": "mysql",
"deleted": false,
"title": "轻薄笔记本"
}
价值:
- 排查同步链路;
- 判断数据新旧;
- 支持重放和幂等;
- 支持软删除过滤;
- 重建索引时校验一致性。
5.12 常见错误
| 错误 | 原因 | 处理 |
|---|---|---|
409 Conflict |
版本冲突 | 重新读取或使用上游版本 |
document_missing_exception |
更新时文档不存在 | 使用 upsert 或先初始化 |
mapper_parsing_exception |
字段类型不匹配 | 修正 Mapping 和数据 |
illegal_argument_exception |
脚本或参数错误 | 检查 Painless 语法 |
es_rejected_execution_exception |
写入队列满 | 降低并发或扩容 |
cluster_block_exception |
磁盘满或索引只读 | 处理磁盘和 block |
5.13 本章小结
- 自动 ID 适合日志和事件,业务 ID 适合实体同步;
_doc/{id}是整体替换,_create是只创建,_update是局部合并;- 更新本质是重新索引整个文档,高频更新应先聚合;
- 使用
_seq_no、_primary_term或外部版本实现并发控制; - 自定义路由可以优化局部查询,但会带来数据倾斜风险;
- bulk 响应必须逐条检查失败项。
5.14 思考题
- 为什么 Elasticsearch 无法像 MySQL 一样原地更新一个字段?
- 上游 MySQL 有
updated_at和自增版本,你会如何设计 ES 幂等同步? refresh=true和refresh=wait_for有什么区别?- 高频订单状态变更是否适合每次
_update?如何优化? - 为什么 bulk HTTP 200 不代表全部写入成功?