ElasticsearchNotes

第 05 章:文档 CRUD:写入、更新与并发控制

zjc 于 2026-01-05 发布

这是《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 不是原地修改某个字段,而是重新索引整个文档。

设计约束:

  1. 大文档频繁更新代价高;
  2. nested 对象越多更新成本越高;
  3. 高频计数和状态变更应先聚合;
  4. 并发更新会出现版本冲突;
  5. 原始字段最好能在同步层完整重建。

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

优点:

  1. 同一用户文档落在同一分片;
  2. 按用户查询时只命中部分分片;
  3. 可优化租户隔离。

风险:

  1. 路由值分布不均导致数据倾斜;
  2. 忘记 routing 读取不到文档;
  3. 大租户可能打爆单个分片;
  4. 查询灵活性下降。

多租户系统更常见的是按租户建索引或使用权限过滤,而不是过度依赖自定义路由。

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"
    }
  }
}

生产建议:

  1. 时间序列数据优先整索引删除;
  2. _delete_by_query 会产生标记删除和合并压力;
  3. 低峰执行;
  4. 控制批次和并发;
  5. 任务失败要记录断点。

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": "轻薄笔记本"
}

价值:

  1. 排查同步链路;
  2. 判断数据新旧;
  3. 支持重放和幂等;
  4. 支持软删除过滤;
  5. 重建索引时校验一致性。

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 本章小结

5.14 思考题

  1. 为什么 Elasticsearch 无法像 MySQL 一样原地更新一个字段?
  2. 上游 MySQL 有 updated_at 和自增版本,你会如何设计 ES 幂等同步?
  3. refresh=truerefresh=wait_for 有什么区别?
  4. 高频订单状态变更是否适合每次 _update?如何优化?
  5. 为什么 bulk HTTP 200 不代表全部写入成功?