阶段:第一阶段 / 核心概念
目标:把写操作(新增、修改、删除、批量)逐个用例子写出来,和第 06 篇(查询速查)配套。
本篇是独立文档,示例不依赖任何具体项目。
约定:文档实体统一用OrderDoc(定义见第 06 篇 §0.1)。
深入版:单文档写入见第 30 篇、批量见第 31 篇、条件更新/删除见第 32 篇。
0. 统一约定
0.1 写操作执行器
@Slf4j@ComponentpublicclassOrderWriter{@AutowiredprivateElasticsearchClientclient;privatestaticfinalStringINDEX="orders_idx";// 下面各例的方法都放在这个类里}0.2 写 API 与 SQL 对照
| ES API | 语义 | SQL 对照 |
|---|---|---|
index | 有则整篇覆盖、无则新建 | INSERT ... ON CONFLICT DO UPDATE(整行替换) |
create | 仅当不存在时创建,存在报错 | INSERT(主键冲突报错) |
update | 局部更新指定字段 / 脚本更新 | UPDATE ... SET col = ? |
delete | 按_id删除 | DELETE ... WHERE id = ? |
updateByQuery | 按条件批量改 | UPDATE ... WHERE ... |
deleteByQuery | 按条件批量删 | DELETE ... WHERE ... |
bulk | 多条写打成一个请求 | 批量INSERT/COPY |
1. 新增 / 覆盖(index与create)
1.1index—— upsert 整篇(最常用)
有则整篇覆盖、无则新建。显式指定_id保证幂等。
/** 相当于 INSERT ... ON CONFLICT (id) DO UPDATE(整行替换) */publicvoidindexDoc(OrderDocdoc)throwsIOException{client.index(i->i.index(INDEX).id(doc.getId())// 用业务主键做 _id.document(doc));// 直接传实体}1.2create—— 仅新建(已存在则报错)
/** 相当于 INSERT,主键冲突会抛异常 */publicvoidcreateDoc(OrderDocdoc)throwsIOException{client.create(c->c.index(INDEX).id(doc.getId()).document(doc));}
index与create的区别:index会覆盖已有文档,create遇到相同_id直接报错。
2. 修改(update)
2.1 局部更新指定字段(UPDATE ... SET)
只改传入的字段,其它字段不动。用一个「部分对象」承载要改的字段。
/** 相当于 UPDATE orders SET status = ?, amount = ? WHERE id = ? */publicvoidupdateFields(Stringid,Map<String,Object>partial)throwsIOException{client.update(u->u.index(INDEX).id(id).doc(partial),// 只更新给出的字段OrderDoc.class);// 返回类型}// 调用:只改 status 和 amountupdateFields("1001",Map.of("status","SHIPPED","amount",newBigDecimal("1299.00")));想用强类型也可以:
.doc(new OrderDoc())只 set 要改的字段——但对象里为 null 的字段
也会被当成「不更新」还是「设为 null」取决于序列化配置,稳妥起见局部更新用Map更可控。
2.2 upsert:不存在就插入,存在就更新
publicvoidupsert(OrderDocdoc,Map<String,Object>partial)throwsIOException{client.update(u->u.index(INDEX).id(doc.getId()).doc(partial)// 存在 → 局部更新这些字段.upsert(doc),// 不存在 → 用整个 doc 新建OrderDoc.class);}2.3 脚本更新(原子自增等)
需要「基于旧值计算」时用脚本,避免读-改-写的并发问题。
/** 相当于 UPDATE orders SET amount = amount + 100 WHERE id = ? */publicvoidincrAmount(Stringid,doubledelta)throwsIOException{client.update(u->u.index(INDEX).id(id).script(sc->sc.inline(in->in.source("ctx._source.amount += params.d").params("d",JsonData.of(delta)))),OrderDoc.class);}3. 删除(delete)
/** 相当于 DELETE FROM orders WHERE id = ? */publicvoiddeleteDoc(Stringid)throwsIOException{client.delete(d->d.index(INDEX).id(id));}4. 条件批量改 / 删(updateByQuery/deleteByQuery)
按条件一次处理一批,不用先查出来再逐条改。
4.1deleteByQuery—— 条件删除
/** 相当于 DELETE FROM orders WHERE status = 'CANCELLED' */publicvoiddeleteCancelled()throwsIOException{client.deleteByQuery(d->d.index(INDEX).query(q->q.term(t->t.field("status").value("CANCELLED"))));}4.2updateByQuery—— 条件更新(脚本)
/** 相当于 UPDATE orders SET status = 'ARCHIVED' WHERE region = 'AP' */publicvoidarchiveApOrders()throwsIOException{client.updateByQuery(u->u.index(INDEX).query(q->q.term(t->t.field("region").value("AP"))).script(sc->sc.inline(in->in.source("ctx._source.status = 'ARCHIVED'"))));}5. 批量处理(bulk)
把多条写操作打成一个请求发出去,是高吞吐写入的标准做法。
5.1 批量新增 / 覆盖
/** 相当于批量 INSERT ... ON CONFLICT */publicvoidbulkIndex(List<OrderDoc>docs)throwsIOException{List<BulkOperation>ops=newArrayList<>(docs.size());for(OrderDocdoc:docs){ops.add(BulkOperation.of(b->b.index(i->i.index(INDEX).id(doc.getId()).document(doc))));}BulkResponseresp=client.bulk(b->b.operations(ops));// 关键:bulk 整体 200 不代表每条都成功,必须逐条检查if(resp.errors()){resp.items().stream().filter(it->it.error()!=null).forEach(it->log.error("bulk 失败 id={}, reason={}",it.id(),it.error().reason()));}}5.2 一个 bulk 里混合 增 / 改 / 删
publicvoidbulkMixed(List<OrderDoc>toIndex,Map<String,Map<String,Object>>toUpdate,List<String>toDelete)throwsIOException{List<BulkOperation>ops=newArrayList<>();// 新增/覆盖for(OrderDocdoc:toIndex){ops.add(BulkOperation.of(b->b.index(i->i.index(INDEX).id(doc.getId()).document(doc))));}// 局部更新toUpdate.forEach((id,partial)->ops.add(BulkOperation.of(b->b.update(u->u.index(INDEX).id(id).action(a->a.doc(partial))))));// 删除for(Stringid:toDelete){ops.add(BulkOperation.of(b->b.delete(d->d.index(INDEX).id(id))));}BulkResponseresp=client.bulk(b->b.operations(ops));if(resp.errors()){resp.items().stream().filter(it->it.error()!=null).forEach(it->log.error("bulk 失败 id={}, reason={}",it.id(),it.error().reason()));}}相关 import:
co.elastic.clients.elasticsearch.core.bulk.BulkOperation、co.elastic.clients.elasticsearch.core.BulkResponse。
6. 速查总表
| 需求 | 方法 | SQL 类比 | 本篇 | 深入 |
|---|---|---|---|---|
| 新增/覆盖单条 | client.index(...) | INSERT ... ON CONFLICT | §1.1 | 第 30 篇 |
| 仅新建单条 | client.create(...) | INSERT(冲突报错) | §1.2 | 第 30 篇 |
| 改单条部分字段 | client.update(...).doc(...) | UPDATE ... SET | §2.1 | 第 30 篇 |
| upsert | client.update(...).upsert(...) | INSERT ... ON CONFLICT | §2.2 | 第 30 篇 |
| 脚本更新 | client.update(...).script(...) | SET x = x + ? | §2.3 | 第 30 篇 |
| 删单条 | client.delete(...) | DELETE WHERE id=? | §3 | 第 30 篇 |
| 条件删 | client.deleteByQuery(...) | DELETE WHERE ... | §4.1 | 第 32 篇 |
| 条件改 | client.updateByQuery(...) | UPDATE WHERE ... | §4.2 | 第 32 篇 |
| 批量 | client.bulk(...) | 批量INSERT/COPY | §5 | 第 31 篇 |
7. 坑与最佳实践
- 显式指定
_id:用业务主键做_id,重跑不产生重复文档(幂等)。不指定会随机生成 id。 index是整篇替换:漏传字段会丢失;只改部分字段用update。- bulk 必须查
resp.errors():整体 200 不代表每条成功,要逐条检查items()。 - 批大小适中:一批几 MB / 几千条为宜,太大易超时/OOM,太小失去批量意义。
- 写入近实时(NRT):默认 1s 后可搜到(
refresh_interval),不是立即可见。 - 大批导入临时调优:可临时设
refresh_interval: -1、number_of_replicas: 0,完成后恢复。 - 并发更新用乐观锁:
if_seq_no/if_primary_term防丢更新;自增类改动优先用脚本更新。 updateByQuery/deleteByQuery是重操作:大范围执行会扫描很多文档,注意配slices并行与错峰。
第一阶段小结
到这里第一阶段(00–07)完成,你应该能:
- 用 SQL 心智理解 index / document / mapping(01、02)
- 分清
text/keyword,理解分词与_score(02、03) - 独立配置并使用官方
ElasticsearchClient(04) - 看懂查询相关类并写复杂
bool动态拼接(05) - 熟练用各类查询(06)与各类写操作(07),结果/文档统一用实体类
下一篇
进入第二阶段查询能力:10-match-全文匹配.md。
写操作的深入讲解(幂等、乐观锁、bulk 调优、条件更新)见第四阶段 30–33 篇。