第六阶段 52 · 并发控制与乐观锁(并发读写下 ES 怎么保证一致) 52 · 并发控制与乐观锁并发读写下 ES 怎么保证一致阶段第六阶段 / 进阶专题ES_seq_no/_primary_term、version、retry_on_conflict、NRT | PostgreSQLMVCC、行锁、事务、乐观锁1. 概念ES 没有事务和行锁用“乐观并发控制”从 PG 转过来最大的心智差异维度PostgreSQLElasticsearch多语句事务✅ 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_no12if_primary_term3 { order_no: SO-1, status: PAID }若期间别人改过_seq_no变了返回409version_conflict_engine_exception本次写失败——由应用决定重读后重试而不是阻塞等待。旧写法?versionNversion_typeinternal已被_seq_no/_primary_term取代新代码用后者。3._update自动重试 批量冲突策略3.1 单文档retry_on_conflict_update读-改-写一体可让 ES 自动重读重试省去手写循环POST orders_idx/_update/order-1?retry_on_conflict3 { script: { source: ctx._source.view_count 1 } }冲突时最多自动重试 3 次适合计数器这类高并发自增。3.2 批量conflictsproceedupdate_by_query/delete_by_query遇到版本冲突默认中止想“跳过冲突项继续跑完”用conflictsproceedPOST orders_idx/_update_by_query?conflictsproceed { script: {...}, query: {...} }4. external version用外部系统的版本号数据权威在别处如 PG 的updated_at/递增版本时用外部版本让 ES 只接受“更新的版本”PUT orders_idx/_doc/order-1?version170000version_typeexternal { ... }只有传入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,MapString,Objectdoc,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供后续乐观锁写入使用 */publicGetResultMapgetForUpdate(Stringindex,Stringid)throwsIOException{GetResponseMaprespelasticsearchClient.get(g-g.index(index).id(id),Map.class);// resp.seqNo()、resp.primaryTerm()、resp.source()returnnewGetResult(resp.source(),resp.seqNo(),resp.primaryTerm());}publicrecordGetResultT(Tsource,LongseqNo,LongprimaryTerm){}}importco.elastic.clients.elasticsearch._types.ElasticsearchException、co.elastic.clients.elasticsearch.core.GetResponse。计数器类高并发优先用_updateretry_on_conflict。7. 坑与最佳实践没有跨文档事务需要多文档“要么全成”时在应用层用幂等 补偿别指望 ES 事务。默认 last-write-wins不带版本号并发写后到的覆盖前者要防覆盖必须用 OCC 或 external version。409 是正常信号不是故障说明有并发重读最新版本再写即可有限次退避重试。计数器/累加用_updateretry_on_conflict别自己读-改-写留下竞态窗口。DB→ES 同步用 external version以源库版本为准天然防乱序旧值覆盖新值。读己所写要注意 NRT刚写完就要查到测试可refreshtrue生产别滥用伤性能。长导出用 PIT并发写不影响你这次遍历的一致视图第 40 篇。相关近实时可见性08-刷新与可见性refresh与NRT.md批量更新/删除与冲突32-update_by_query-delete_by_query.md一致性快照遍历PIT40-深分页与性能.md