52 · 并发控制与乐观锁(并发读写下 ES 怎么保证一致)
阶段:第六阶段 / 进阶专题
ES:_seq_no/_primary_term、version、retry_on_conflict、NRT | PostgreSQL:MVCC、行锁、事务、乐观锁
1. 概念:ES 没有事务和行锁,用“乐观并发控制”
从 PG 转过来最大的心智差异:
| 维度 | PostgreSQL | Elasticsearch |
|---|---|---|
| 多语句事务 | ✅ BEGIN/COMMIT | ❌ 没有跨文档事务 |
| 原子性粒度 | 行 + 事务 | 单篇文档(doc 级原子) |
| 并发控制 | MVCC + 行锁(悲观为主) | 乐观并发控制(OCC),冲突让应用重试 |
| 读写互斥 | 读一般不阻塞写(MVCC) | 读写互不加锁(段不可变) |
| 可见性 | 提交即可见(读己所写) | 近实时,写完约 1s 后可搜 |
核心:ES 不加锁等待,而是“乐观地写,冲突了报错让你重试”。
2. 写并发:_seq_no+_primary_term乐观锁
每篇文档都有两个并发控制元数据:
_seq_no:该分片上写操作的序号(每次写自增)。_primary_term:主分片的“任期”,主分片切换时递增。
流程(“比对再写”):
# 1) 先读,拿到当前 _seq_no / _primary_term GET orders_idx/_doc/order-1 # => "_seq_no": 12, "_primary_term": 3 # 2) 带着这两个值去写:只有没被别人改过才成功 PUT orders_idx/_doc/order-1?if_seq_no=12&if_primary_term=3 { "order_no": "SO-1", "status": "PAID" }若期间别人改过(_seq_no变了),返回409version_conflict_engine_exception,
本次写失败——由应用决定重读后重试,而不是阻塞等待。
旧写法
?version=N&version_type=internal已被_seq_no/_primary_term取代,新代码用后者。
3._update自动重试 & 批量冲突策略
3.1 单文档:retry_on_conflict
_update(读-改-写一体)可让 ES 自动重读重试,省去手写循环:
POST orders_idx/_update/order-1?retry_on_conflict=3 { "script": { "source": "ctx._source.view_count += 1" } }冲突时最多自动重试 3 次,适合计数器这类高并发自增。
3.2 批量:conflicts=proceed
update_by_query/delete_by_query遇到版本冲突默认中止;
想“跳过冲突项继续跑完”用conflicts=proceed:
POST orders_idx/_update_by_query?conflicts=proceed { "script": {...}, "query": {...} }4. external version(用外部系统的版本号)
数据权威在别处(如 PG 的updated_at/递增版本)时,用外部版本,让 ES 只接受“更新的版本”:
PUT orders_idx/_doc/order-1?version=170000&version_type=external { ... }只有传入version大于当前值才写入——天然实现“旧数据不覆盖新数据”,
非常适合 DB → ES 同步(乱序消息也不会用旧值盖新值)。
5. 读并发:读不阻塞写、近实时可见
- 段不可变:Lucene 写入生成新段,旧段继续服务查询;读写互不加锁(类似 MVCC)。
- 近实时(NRT):写入先进内存 buffer + translog,默认每
refresh_interval(~1s)
才 refresh 成可搜索段。所以“写完立刻查不到”是正常的(第 08 篇)。 - 持久性:translog 保证宕机不丢已确认的写;flush 把段落盘。
- 遍历期间一致:一次
search看当次段快照;长时间分页/导出用PIT(第 18/40 篇)
锁定快照,避免翻页途中被并发写“串动”。
6. Spring Boot 实现
@ComponentpublicclassDoc52Concurrency{@AutowiredprivateElasticsearchClientelasticsearchClient;/** 乐观锁更新:带 if_seq_no/if_primary_term,冲突则由调用方重试 */publicbooleanupdateWithOcc(Stringindex,Stringid,Map<String,Object>doc,longseqNo,longprimaryTerm)throwsIOException{try{elasticsearchClient.index(i->i.index(index).id(id).ifSeqNo(seqNo).ifPrimaryTerm(primaryTerm).document(doc));returntrue;}catch(ElasticsearchExceptione){if(e.status()==409){// version_conflict_engine_exceptionlog.warn("并发冲突,需重读重试 id={}",id);returnfalse;// 交给上层重读最新 _seq_no 后重试}throwe;}}/** 读时拿到 _seq_no / _primary_term,供后续乐观锁写入使用 */publicGetResult<Map>getForUpdate(Stringindex,Stringid)throwsIOException{GetResponse<Map>resp=elasticsearchClient.get(g->g.index(index).id(id),Map.class);// resp.seqNo()、resp.primaryTerm()、resp.source()returnnewGetResult<>(resp.source(),resp.seqNo(),resp.primaryTerm());}publicrecordGetResult<T>(Tsource,LongseqNo,LongprimaryTerm){}}import:
co.elastic.clients.elasticsearch._types.ElasticsearchException、co.elastic.clients.elasticsearch.core.GetResponse。计数器类高并发优先用_update+retry_on_conflict。
7. 坑与最佳实践
- 没有跨文档事务:需要多文档“要么全成”时,在应用层用幂等 + 补偿,别指望 ES 事务。
- 默认 last-write-wins:不带版本号并发写,后到的覆盖前者;要防覆盖必须用 OCC 或 external version。
- 409 是正常信号不是故障:说明有并发,重读最新版本再写即可(有限次退避重试)。
- 计数器/累加用
_update+retry_on_conflict,别自己读-改-写留下竞态窗口。 - DB→ES 同步用 external version:以源库版本为准,天然防乱序旧值覆盖新值。
- 读己所写要注意 NRT:刚写完就要查到,测试可
refresh=true(生产别滥用,伤性能)。 - 长导出用 PIT:并发写不影响你这次遍历的一致视图(第 40 篇)。
相关
- 近实时可见性:
08-刷新与可见性(refresh与NRT).md - 批量更新/删除与冲突:
32-update_by_query-delete_by_query.md - 一致性快照遍历(PIT):
40-深分页与性能.md