第五阶段 42 · 生产实战解析(查询构建 / 取数 / mapping 的最佳实践落地) 42 · 生产实战解析查询构建 / 取数 / mapping 的最佳实践落地阶段第五阶段 / 进阶与实战目标把前面学的概念串成一套可直接用于生产的实现骨架——动态查询构建、稳定的全量取数、mapping 管理。示例是通用最佳实践不绑定任何具体工程拿到任何 Spring Boot 官方co.elastic.clients的项目都能落地。1. 整体链路一条「DB → ES → 查询」的销售数据链路生产上通常这样分层建索引: IndexAdminService.createIndexIfAbsent → 读 mapping 定义JSON/代码 写数据: DB 批量取数keyset 分页 → BulkIngester 攒批 → 写 ES 查数据: SearchCriteria入参 → QueryBuilder 拼 bool → SearchService → PIT search_after 取数对照前面的文档mapping第 2/33 篇、bulk第 31 篇、bool/term第 11/13 篇、深分页第 40 篇。下面按「查询构建 → 取数 → mapping」逐层给出生产骨架。2. 查询构建把入参翻译成 Query生产里查询入口一般接收一个 DTOSearchCriteria支持两种模式结构化字段查询前端传字段条件后端逐字段判空拼bool对照第 05、13 篇。原始 DSL 直通前端已经拼好一段 Kibana DSL后端原样透传不逐条翻译。publicQuerybuildBaseQuery(SearchCriteriacriteria){// 前端直接给了原始 DSL就走直通否则按字段拼 boolreturncriteria.hasRawDsl()?wrapRawDsl(criteria.getRawDsl()):buildFieldQuery(criteria);}2.1 字段动态拼条件对照第 13 篇 bool核心思路和 MyBatis 的if一样逐个入参判空非空才 add 子句精确/范围条件放filter不打分、可缓存、更快。privateQuerybuildFieldQuery(SearchCriteriac){if(cnull||!c.hasCondition()){// 无条件 SELECT *用 match_all别返回 nullreturnQuery.of(q-q.matchAll(m-m));}BoolQuery.BuilderboolnewBoolQuery.Builder();if(CollectionUtils.isNotEmpty(c.getRecordIdList())){bool.filter(f-f.terms(t-t.field(FieldNames.RECORD_ID).terms(tv-tv.value(toFieldValues(c.getRecordIdList())))));}if(CollectionUtils.isNotEmpty(c.getInvoiceNumberList())){bool.filter(f-f.terms(t-t.field(FieldNames.INVOICE_NUMBER).terms(tv-tv.value(toFieldValues(c.getInvoiceNumberList())))));}// material / region / 日期区间 …… 同理为空则跳过returnQuery.of(q-q.bool(bool.build()));}privateListFieldValuetoFieldValues(ListStringvalues){returnvalues.stream().map(FieldValue::of).collect(Collectors.toList());}要点每个「多值」条件用terms相当于 SQLIN (...)作用在keyword字段上做精确匹配。「条件为空就跳过」 MyBatisif的思路全空则退化成match_all。字段名集中到常量类如FieldNames避免散落的魔法字符串导致拼写漂移。新增查询字段时FieldNames加常量 SearchCriteria加入参 这里加一段if。2.2 原始 DSL 直通wrapper前端复杂查询Kibana 里调好的 DSL不必逐条翻译成 Java用WrapperQuery直接透传它接收一段 Base64 编码的 JSON DSL。privateQuerywrapRawDsl(StringrawJson){StringencodedBase64.getEncoder().encodeToString(rawJson.getBytes(StandardCharsets.UTF_8));returnQuery.of(q-q.wrapper(w-w.query(encoded)));}⚠️ 直通很方便但等于把「查询构造」暴露给上游。生产上要校验/白名单传入的 DSL避免注入代价高的查询深聚合、script、超大terms拖垮集群。3. 取数稳定的全量遍历PIT search_after导出/同步这类「把匹配到的数据全捞出来」的场景生产首选 PIT search_after对照第 40 篇而不是旧的 scroll——scroll 有状态、占段资源官方已不推荐用于新场景。publicvoidscanAll(Stringindex,Queryquery,ConsumerListOrderDocbatchConsumer)throwsIOException{// 1) 开一个 Point In Time拿到一致性快照StringpitIdclient.openPointInTime(p-p.index(index).keepAlive(t-t.time(2m))).id();try{ListFieldValuesearchAfternull;while(true){ListFieldValueaftersearchAfter;SearchResponseOrderDocrespclient.search(s-{s.size(1000).pit(p-p.id(pitId).keepAlive(t-t.time(2m)))// 排序必须带唯一 tiebreakerrecord_id否则翻页会漏/重.sort(so-so.field(f-f.field(record_id).order(SortOrder.Asc)));if(after!null){s.searchAfter(after);}returns;},OrderDoc.class);ListHitOrderDochitsresp.hits().hits();if(hits.isEmpty()){break;}batchConsumer.accept(hits.stream().map(Hit::source).filter(Objects::nonNull).collect(Collectors.toList()));// 用最后一条的排序值作为下一页游标searchAfterhits.get(hits.size()-1).sort();}}finally{// 2) PIT 用完必须释放否则占用段资源client.closePointInTime(c-c.id(pitId));}}要点对照第 40 篇一致性PIT 是快照遍历期间数据变化不会导致漏/重。稳定排序search_after依赖唯一排序键业务字段不唯一时务必追加record_id兜底。释放资源finally里closePointInTime异常也要清理。_source 裁剪第 19 篇只导出需要的列size取几百到一千即可别贪大。如果只是「有限量预览」比如导出前 1000 条可以用简单的from/size但不要用它做真正的全量深翻页fromsize上限默认 10000见第 40 篇。4. 写入BulkIngester 攒批DB → ES 的批量写生产上推荐官方BulkIngester自动攒批、并发、背压比手工拼BulkRequest更省心对照第 31 篇try(BulkIngesterVoidingesterBulkIngester.of(b-b.client(client).maxOperations(2000)// 攒够 2000 条自动 flush.maxSize(5*1024*1024)// 或攒够 5MB.flushInterval(2,TimeUnit.SECONDS)// 或每 2s flush 一次.listener(newBulkListener(){// 回调里检查每条结果OverridepublicvoidafterBulk(longid,BulkRequestreq,ListVoidctx,BulkResponseresp){if(resp.errors()){resp.items().stream().filter(i-i.error()!null).forEach(i-log.error(bulk failed id{}, reason{},i.id(),i.error().reason()));}}OverridepublicvoidafterBulk(longid,BulkRequestreq,ListVoidctx,Throwablefailure){log.error(bulk request failed,failure);}}))){for(OrderDocdoc:docs){ingester.add(op-op.index(i-i.index(index).id(doc.getRecordId())// 显式 _id → 幂等重跑不产生重复.document(doc)));}}// try-with-resources 关闭时自动 flush 剩余数据要点幂等显式_id重跑同一批不会产生重复文档。错误可见bulk整体 200 不代表每条成功必须在 listener 里查resp.errors()。导入期临时调优大批导入可临时refresh_interval-1、number_of_replicas0完成后恢复。5. mapping 管理显式定义 别名生产索引必须显式 mapping别靠动态映射猜类型业务字段基本是keyword精确匹配 可聚合日期字段用date并指定format对照第 2/33 篇。mapping 可以放在资源文件里如mappings/order_idx.json集中维护建索引时读取{settings:{number_of_shards:3,number_of_replicas:1},mappings:{properties:{record_id:{type:keyword},invoice_number:{type:keyword},material:{type:keyword},net_amount:{type:double},invoice_dt:{type:date,format:yyyy-MM-dd}}}}新增字段进 ES就是在这份 mapping 的properties里追加一段customer_group:{type:keyword}对已上线索引用PUT /index/_mapping追加字段只能加、不能改类型需要改类型时走 reindex 别名切换第 33 篇。应用只认别名如order底层索引版本切换对上层透明、零停机。6. 想改点什么就看这里需求改哪里参考新增可查询字段FieldNames常量 SearchCriteria入参 buildFieldQuery加if第 05/13 篇新增 ES 存储字段mappingproperties追加 写 ES 的取数 SQL/DTO第 33 篇换全量导出方案scanAll用 PIT search_after别用 scroll第 40 篇批量写吞吐不够调BulkIngester的maxOperations/maxSize、导入期关 refresh第 31 篇复杂前端查询走原始 DSL WrapperQuery直通记得校验白名单§2.2改索引类型/分片新建索引 reindex 别名原子切换第 33 篇