Elasticsearch聚合查询

Eave 2026.09.01

一、创建API客户端

/**
 * 创建API客户端
 *
 * @return
 */
public static ElasticsearchClient client()
{
    // 1. 配置认证信息
    // final CredentialsProvider credentialsProvider = new BasicCredentialsProvider();
    // credentialsProvider.setCredentials(AuthScope.ANY, new UsernamePasswordCredentials("elastic", "your_password"));

    // 1. 创建低级别的 RestClient
    // 注意:HttpHost 的构造参数是 (hostname, port, scheme)
    // RestClient restClient = RestClient
    //         .builder(new HttpHost("localhost", 9200, "http"))
    //         .setHttpClientConfigCallback(httpClientBuilder -> httpClientBuilder.setDefaultCredentialsProvider(credentialsProvider)
    // ).build();

    // 1. 创建低级别的 RestClient
    // 注意:HttpHost 的构造参数是 (hostname, port, scheme)
    RestClient restClient = RestClient.builder(new HttpHost("localhost", 9200, "http")).build();

    // 2. 使用 Jackson 映射器创建 Transport 层
    ElasticsearchTransport transport = new RestClientTransport(restClient, new JacksonJsonpMapper());

    // 3. 创建 API 客户端
    return new ElasticsearchClient(transport);
}

二、列出所有索引

/**
 * 列出所有索引
 *
 * @param client
 * @throws Exception
 */
public static void allIndex(ElasticsearchClient client) throws Exception
{
    GetIndexResponse response = client.indices().get(new GetIndexRequest.Builder().index("*").build());

    response.indices().forEach((indexName, indexState) -> {
        logger.info("index: {}, state: {}", indexName, indexState);
    });

    IndicesResponse catResponse = client.cat().indices();
    logger.info(JsonUtil.toJson(catResponse.indices()));
}

三、创建索引

/**
 * 创建索引
 *
 * @param client
 * @param indexName
 * @throws Exception
 */
public static void createIndex(ElasticsearchClient client, String indexName) throws Exception
{
    // 1. 检查索引是否已存在
    BooleanResponse exists = client.indices().exists(ExistsRequest.of(e -> e.index(indexName)));
    logger.info("exists: {}", exists.value());

    if(exists.value())
    {
        return;
    }

    // 2. 创建索引
    String settingsJson = FileReader.read("D:\\Users\\user\\Desktop\\" + indexName + ".json");
    CreateIndexResponse response = client.indices().create(c -> c.index(indexName).withJson(new ByteArrayInputStream(settingsJson.getBytes(StandardCharsets.UTF_8))));
    logger.info("acknowledged: {}", response.acknowledged());
}

四、删除索引

/**
 * 删除索引
 *
 * @param client
 * @param indexName
 * @throws Exception
 */
public static void deleteIndex(ElasticsearchClient client, String indexName) throws Exception
{
    // 1. 检查索引是否存在
    BooleanResponse exists = client.indices().exists(ExistsRequest.of(e -> e.index(indexName)));
    logger.info("exists: {}", exists.value());

    if(!exists.value())
    {
        return;
    }

    // 2. 删除索引
    DeleteIndexResponse response = client.indices().delete(d -> d.index(indexName));
    logger.info("acknowledged: {}", response.acknowledged());
}

五、获取索引别名

/**
 * 获取索引别名
 *
 * @param client
 * @throws Exception
 */
public static void getAlias(ElasticsearchClient client) throws Exception
{
    // 查看某个索引的所有别名
    GetAliasResponse response1 = client.indices().getAlias(g -> g.index("user_v1"));
    logger.info("索引 user_v1 的别名: {}", response1.aliases().keySet());

    // 查看别名指向哪些索引
    GetAliasResponse response2 = client.indices().getAlias(g -> g.name("user"));
    logger.info("别名 user 指向: {}", response2.aliases().keySet());
}

六、添加别名

/**
 * 添加别名
 *
 * @param client
 * @throws Exception
 */
public static void addAlias(ElasticsearchClient client) throws Exception
{
    client.indices().updateAliases(u -> u
            .actions(a -> a
                    .add(add -> add
                            .index("user_v2")
                            .alias("user")
                    )
            )
    );
}

七、检查别名是否存在

/**
 * 检查别名是否存在
 *
 * @param client
 * @throws Exception
 */
public static void existsAlias(ElasticsearchClient client) throws Exception
{
    BooleanResponse response = client.indices().existsAlias(e -> e.name("user"));
    logger.info("exists: {}", response.value());
}

八、删除别名

/**
 * 删除别名
 *
 * @param client
 * @throws Exception
 */
public static void removeAlias(ElasticsearchClient client) throws Exception
{
    client.indices().updateAliases(u -> u
            .actions(a -> a
                    .remove(remove -> remove
                            .index("user_v1")
                            .alias("user")
                    )
            )
    );
}

九、切换别名

/**
 * 切换别名
 *
 * @param client
 * @throws Exception
 */
public void switchAlias(ElasticsearchClient client) throws Exception
{
    // 将别名 "user" 从 user_v1 切换到 user_v2
    client.indices().updateAliases(u -> u
            .actions(a -> a
                    .remove(remove -> remove
                            .index("user_v1")
                            .alias("user")
                    )
            )
            .actions(a -> a
                    .add(add -> add
                            .index("user_v2")
                            .alias("user")
                    )
            )
    );
}

十、添加数据(完全替换)

/**
 * 添加数据(完全替换)
 *
 * @param client
 * @param indexName
 * @throws Exception
 */
public static void putObject(ElasticsearchClient client, String indexName) throws Exception
{
    // 1. 准备数据(可以用 Map 或 POJO)
    IdAndName document = new IdAndName(UuidUtil.gen24Id(), "张三测");

    // 2. 构建 Index 请求
    IndexRequest<IdAndName> request = IndexRequest.of(i -> i
            .index(indexName)          // 索引名
            .id(document.getId())      // 文档 ID(如果不指定,ES 会自动生成)
            .document(document)        // 要写入的数据
    );

    // 3. 执行请求
    IndexResponse response = client.index(request);

    // 4. 处理响应
    logger.info("索引名: {}", response.index());
    logger.info("文档ID: {}", response.id());
    logger.info("版本号: {}", response.version());
    logger.info("结果: {}", response.result().jsonValue()); // "created" 或 "updated"
}

十一、批量插入(Bulk)

/**
 * 批量插入(Bulk)
 *
 * @param client
 * @param indexName
 * @throws Exception
 */
public static void bulkInsert(ElasticsearchClient client, String indexName) throws Exception
{
    // 准备要更新的数据
    List<IdAndName> inserts = List.of(
            new IdAndName("P001", "price1"),
            new IdAndName("P002", "price2"),
            new IdAndName("P003", "price3")
    );

    BulkResponse response = client.bulk(b -> {
        for(IdAndName data : inserts)
        {
            b.operations(op -> op
                    .index(idx -> idx
                            .index(indexName)
                            .id(data.getId())
                            .document(data)
                    )
            );
        }
        return b;
    });

    if(response.errors())
    {
        logger.info("批量更新有错误!");
        response.items().forEach(item -> logger.warn("ID: {}, 错误: {}", item.id(), item.error()));
    }
    else
    {
        logger.info("批量更新成功,共 {} 条", response.items().size());
    }
}

十二、更新数据(部分更新)

/**
 * 更新数据(部分更新)
 *
 * @param client
 * @param indexName
 * @throws Exception
 */
public static void updateObject(ElasticsearchClient client, String indexName) throws Exception
{
    // 1. 准备数据(可以用 Map 或 POJO)
    IdAndName document = new IdAndName(UuidUtil.gen24Id(), "张三测");

    // 2. 构建 Update 请求
    UpdateRequest<IdAndName, IdAndName> request = UpdateRequest.of(i -> i
            .index(indexName)          // 索引名
            .id(document.getId())      // 文档 ID(如果不指定,ES 会自动生成)
            .doc(document)             // 要写入的数据
            .docAsUpsert(false)        // 不存在时不插入
    );

    // 3. 执行请求
    UpdateResponse<IdAndName> response = client.update(request, IdAndName.class);  // 返回更新后的完整文档

    // 4. 处理响应
    logger.info("索引名: {}", response.index());
    logger.info("文档ID: {}", response.id());
    logger.info("版本号: {}", response.version());
    logger.info("结果: {}", response.result().jsonValue()); // "created" / "updated" / "noop" / "deleted"
}

十三、批量更新(Bulk)

/**
 * 批量更新(Bulk)
 *
 * @param client
 * @param indexName
 * @throws Exception
 */
public static void bulkUpdate(ElasticsearchClient client, String indexName) throws Exception
{
    // 准备要更新的数据
    List<IdAndName> updates = List.of(
            new IdAndName("P001", "price1"),
            new IdAndName("P002", "price2"),
            new IdAndName("P003", "price3")
    );

    BulkResponse response = client.bulk(b -> {
        for(IdAndName data : updates)
        {
            Map<String, Object> doc = new HashMap<>();
            doc.put("name", data.getName());

            b.operations(op -> op
                    .update(up -> up
                            .index(indexName)
                            .id(data.getId())
                            .action(ua -> ua
                                    .doc(doc)
                            )
                    )
            );
        }
        return b;
    });

    if(response.errors())
    {
        logger.info("批量更新有错误");
        response.items().forEach(item -> logger.warn("ID: {}, 错误: {}", item.id(), item.error()));
    }
    else
    {
        logger.info("批量更新成功,共 {} 条", response.items().size());
    }
}

十四、模糊查询&高亮

/**
 * 模糊查询&高亮
 *
 * @param client
 * @param indexName
 * @param fuzzyPattern
 * @throws Exception
 */
public static void searchByWildcard(ElasticsearchClient client, String indexName, String fuzzyPattern) throws Exception
{
    // fuzzyPattern 示例:
    // 查前缀: "123e4567*"
    // 查中间: "*e89b*"
    // 查后缀: "*174000"
    // 查完整: "123e4567-e89b-12d3-a456-426614174000"
    SearchResponse<IdAndName> response1 = client.search(s -> s
            .index(indexName) // 可以直接使用别名
            .size(100)
            .query(q -> q
                    .wildcard(w -> w
                            .field("id")         // 字段名
                            .value(fuzzyPattern) // 带 * 的模糊表达式
                    )
            ), IdAndName.class                   // 你的 POJO 类
    );

    SearchResponse<IdAndName> response = client.search(s -> s
            .index(indexName)
            .size(100)
            .query(q -> q
                    .bool(b -> b
                            .must(must -> must
                                    .wildcard(t -> t
                                            .field("id")
                                            .value(fuzzyPattern)
                                    )
                            )
                            .must(must -> must
                                    .match(m -> m
                                            .field("name")
                                            .query("测")
                                            .fuzziness("AUTO")        // 自动模糊度
                                            .operator(Operator.And)   // 所有词必须匹配
                                    )
                            )
                    )
            )
            .highlight(h -> h
                    .preTags("<em>")
                    .postTags("</em>")
                    .fields(NamedValue.of("name", HighlightField.of(f -> f.numberOfFragments(5).fragmentSize(80))))
            ), IdAndName.class
    );

    for(Hit<IdAndName> hit : response.hits().hits())
    {
        logger.info("Found: {}", JsonUtil.toJson(hit.source()));

        // 获取高亮
        hit.highlight().forEach((key, value) -> logger.info("{} 高亮: {}", key, JsonUtil.toJson(value)));
    }
}

十五、按自定义价格区间聚合

/**
 * 按自定义价格区间聚合
 *
 * @param client
 * @throws Exception
 */
public static void priceRangeAggregation(ElasticsearchClient client) throws Exception
{
    SearchResponse<Void> response = client.search(s -> s
                    .index("products")
                    .size(0)  // 不返回文档,只返回聚合结果
                    .aggregations("price_ranges", a -> a
                            .range(r -> r
                                    .field("price")
                                    .ranges(
                                            AggregationRange.of(rb -> rb.key("0-100").to(100.0)),
                                            AggregationRange.of(rb -> rb.key("100-500").from(100.0).to(500.0)),
                                            AggregationRange.of(rb -> rb.key("500-1000").from(500.0).to(1000.0)),
                                            AggregationRange.of(rb -> rb.key("1000+").from(1000.0))
                                    )
                            )
                    ), Void.class
    );

    // 解析结果
    RangeAggregate priceRanges = response.aggregations().get("price_ranges").range();
    for(RangeBucket bucket : priceRanges.buckets().array())
    {
        logger.info("区间: {}, 数量: {}", bucket.key(), bucket.docCount());
    }
}

十六、按固定间隔聚合

/**
 * 按固定间隔聚合
 *
 * @param client
 * @throws Exception
 */
public static void priceHistogramAggregation(ElasticsearchClient client) throws Exception
{
    SearchResponse<Void> response = client.search(s -> s
                    .index("products")
                    .size(0)
                    .aggregations("price_histogram", a -> a
                            .histogram(h -> h
                                    .field("price")
                                    .interval(100.0)      // 每 100 元一段
                                    .minDocCount(1)       // 至少 1 个文档才返回
                            )
                    ), Void.class
    );

    HistogramAggregate histogram = response.aggregations().get("price_histogram").histogram();
    for(HistogramBucket bucket : histogram.buckets().array())
    {
        logger.info("价格区间: {} - {}, 数量: {}", bucket.key(), bucket.key() + 100, bucket.docCount());
    }
}

十七、按分类聚合+价格子聚合

/**
 * 按分类聚合+价格子聚合
 *
 * @param client
 * @throws Exception
 */
public static void categoryWithPriceAggregation(ElasticsearchClient client) throws Exception
{
    SearchResponse<Void> response = client.search(s -> s
                    .index("products")
                    .size(0)
                    .aggregations("by_category", a -> a
                            .terms(t -> t
                                    .field("category")
                                    .size(10)
                            )
                            .aggregations("avg_price", a2 -> a2
                                    .avg(avg -> avg.field("price"))
                            )
                            .aggregations("max_price", a2 -> a2
                                    .max(max -> max.field("price"))
                            )
                            .aggregations("min_price", a2 -> a2
                                    .min(min -> min.field("price"))
                            )
                    ), Void.class
    );

    StringTermsAggregate categoryAgg = response.aggregations().get("by_category").sterms();
    for(StringTermsBucket bucket : categoryAgg.buckets().array())
    {
        String category = bucket.key().stringValue();
        long count = bucket.docCount();
        double avgPrice = bucket.aggregations().get("avg_price").avg().value();
        double maxPrice = bucket.aggregations().get("max_price").max().value();
        double minPrice = bucket.aggregations().get("min_price").min().value();

        logger.info("分类: {}, 数量: {}, 均价: {}, 最高: {}, 最低: {}", category, count, avgPrice, maxPrice, minPrice);
    }
}

十八、每个分类下价格最高的前3个商品

/**
 * 每个分类下价格最高的前3个商品
 *
 * @param client
 * @throws Exception
 */
public static void topHitsAggregation(ElasticsearchClient client) throws Exception
{
    SearchResponse<Product> response = client.search(s -> s
                    .index("products")
                    .size(0)
                    .aggregations("by_category", a -> a
                            .terms(t -> t
                                    .field("category")
                                    .size(5)
                            )
                            .aggregations("top_products", a2 -> a2
                                    .topHits(th -> th
                                            .size(3)
                                            .sort(so -> so
                                                    .field(f -> f
                                                            .field("price")
                                                            .order(SortOrder.Desc)
                                                    )
                                            )
                                    )
                            )
                    ), Product.class
    );

    StringTermsAggregate categoryAgg = response.aggregations().get("by_category").sterms();
    for(StringTermsBucket bucket : categoryAgg.buckets().array())
    {
        String category = bucket.key().stringValue();
        logger.info("分类: {}", category);

        TopHitsAggregate topHits = bucket.aggregations().get("top_products").topHits();
        for(Hit<JsonData> hit : topHits.hits().hits())
        {
            // Jackson默认不支持java.time.LocalDateTime hit.source().to(Product.class);
            JsonData source = hit.source();
            Product product = JsonUtil.parse(source.toString(), Product.class);
            logger.info(" - {}, 价格: {}", product.getName(), product.getPrice());
        }
    }
}

十九、每个分类下价格最高的前3个商品(带查询条件)

/**
 * 每个分类下价格最高的前3个商品(带查询条件)
 *
 * @param client
 * @throws Exception
 */
public static void topHitsAggregationQuery(ElasticsearchClient client) throws Exception
{
    // 构建基础请求
    SearchRequest.Builder builder = new SearchRequest.Builder();
    builder.index("products").size(0);

    // 动态设置 query(示例:按品牌过滤)
    String name = "pro";   // 可能是 null
    String brand = null;   // 可能是 null
    Double minPrice = 5000.0; // 可能是 null

    BoolQuery.Builder boolQuery = new BoolQuery.Builder();
    if(name != null)
    {
        boolQuery.must(q -> q
                .term(t -> t
                        .field("name")
                        .value(name)
                )
        );
    }

    if(brand != null)
    {
        boolQuery.must(q -> q
                .term(t -> t
                        .field("brand")
                        .value(brand)
                )
        );

        // 多字段or匹配
        // boolQuery.must(q -> q.multiMatch(m -> m.fields("name", "brand").query("华为")));
        // 前缀匹配
        // boolQuery.must(q -> q.prefix(p -> p.field("name").value("华为")));
    }

    if(minPrice != null)
    {
        // NumberRangeQuery.Builder numberRangeQuery = new NumberRangeQuery.Builder();
        // numberRangeQuery.field("price").gte(minPrice);

        boolQuery.must(Query.of(q -> q
                .range(r -> r
                        .number(n -> n
                                .field("price")
                                .gte(minPrice)
                        )
                )
        ));
    }

    boolQuery.must(Query.of(q -> q
            .range(r -> r
                    .date(d -> d
                            .field("createdAt")
                            .gte("2026-08-01 00:00:00")
                            .format("yyyy-MM-dd HH:mm:ss")
                    )
            )
    ));

    builder.query(q -> q.bool(boolQuery.build()));
    builder.aggregations("by_category", a -> a
            .terms(t -> t
                    .field("category")
                    .size(5)
            )
            .aggregations("top_products", a2 -> a2
                    .topHits(th -> th
                            .size(3)
                            .sort(so -> so
                                    .field(f -> f
                                            .field("price")
                                            .order(SortOrder.Desc)
                                    )
                            )
                    )
            )
    );

    // 执行查询
    SearchResponse<Product> response = client.search(builder.build(), Product.class);
    StringTermsAggregate categoryAgg = response.aggregations().get("by_category").sterms();
    for(StringTermsBucket bucket : categoryAgg.buckets().array())
    {
        String category = bucket.key().stringValue();
        logger.info("分类: {}", category);

        TopHitsAggregate topHits = bucket.aggregations().get("top_products").topHits();
        for(Hit<JsonData> hit : topHits.hits().hits())
        {
            // Jackson默认不支持java.time.LocalDateTime hit.source().to(Product.class);
            JsonData source = hit.source();
            Product product = JsonUtil.parse(source.toString(), Product.class);
            logger.info(" - {}, 价格: {}", product.getName(), product.getPrice());
        }
    }
}

二十、搜索商品(带商品可售范围)

/**
 * 搜索商品(带商品可售范围)
 *
 * @param client
 * @throws Exception
 */
public static void sellAreasQuery(ElasticsearchClient client) throws Exception
{
    // 区域存储:勾选到哪级存哪级,不用向下向上展开
    // {"id":"1001","name":"华为 Mate 80 Pro","price":6999.00,"category":"手机","brand":"华为","sellAreas":["440000"]}
    // {"id":"1002","name":"小米 17 Pro","price":4499.00,"category":"手机","brand":"小米","sellAreas":["440300"]}
    // {"id":"1003","name":"OPPO Find X9","price":4599.00,"category":"手机","brand":"OPPO","sellAreas":["440305"]}
    // {"id":"1004","name":"vivo X300","price":4599.00,"category":"手机","brand":"vivo","sellAreas":["440305001"]}

    // 搜索条件:
    // 南头街道(440305001)
    //     -> 440305001(全匹配/前缀匹配)
    //     -> 440305(全匹配)
    //     -> 440300(全匹配)
    //     -> 440000(全匹配)
    //
    // 南山区(440305)
    //     -> 440305(前缀匹配)
    //     -> 440300(全匹配)
    //     -> 440000(全匹配)
    //
    // 深圳市(440300)
    //     -> 4403(前缀匹配)
    //     -> 440000(全匹配)
    //
    // 广东省(440000)
    //     -> 44(前缀匹配)
    //
    // 总结:
    //     下级和自己 -> 前缀匹配
    //     上级 -> 全匹配

    Function<String, List<String>> ancestorsFunc = code -> {
        // 街道
        if(code.length() == 9)
        {
            return List.of(code, code.substring(0, 6), code.substring(0, 4) + "00", code.substring(0, 2) + "0000");
        }

        if(code.length() == 6)
        {
            // 省
            if(code.endsWith("0000"))
            {
                return List.of(code.substring(0, 2));
            }

            // 市
            if(code.endsWith("00"))
            {
                return List.of(code.substring(0, 4), code.substring(0, 2) + "0000");
            }

            // 区
            return List.of(code, code.substring(0, 4) + "00", code.substring(0, 2) + "0000");
        }

        return List.of(code);
    };

    // 用户定位:深圳市南山区
    String districtCode = "440305";

    // 1.展开祖先链:区 → 市 → 省
    List<String> ancestors = ancestorsFunc.apply(districtCode);

    // 2.构建 prefix OR 条件
    BoolQuery.Builder bool = new BoolQuery.Builder();
    for(int i = 0; i < ancestors.size(); i++)
    {
        String code = ancestors.get(i);
        Query query = i == 0 ? Query.of(q -> q
                .prefix(p -> p
                        .field("sellAreas")
                        .value(code)
                )
        ) : Query.of(q -> q
                .term(p -> p
                        .field("sellAreas")
                        .value(code)
                )
        );
        bool.should(query);
    }
    bool.minimumShouldMatch("1");

    Query query = Query.of(q -> q.bool(bool.build()));

    // 3.执行
    SearchResponse<Product> response = client.search(s -> s
                    .index("sale_areas")
                    .source(src -> src
                            .filter(f -> f
                                    .includes("name", "price", "brand")  // 只返回这些字段
                            )
                    )
                    .query(query),
            Product.class
    );

    for(Hit<Product> hit : response.hits().hits())
    {
        Product product = hit.source();
        logger.info("{}:{}", product.getName(), product.getSellAreas());
    }
}

二十一、搜索附近商家 按距离排序

/**
 * 搜索附近商家 按距离排序
 *
 * "location": { "type": "geo_point" }
 * // 格式1:对象
 * { "location": { "lat": 22.5429, "lon": 114.0596 } }
 * 
 * // 格式2:字符串
 * { "location": "22.5429,114.0596" }
 * 
 * // 格式3:数组(注意:经度在前,纬度在后)
 * { "location": [114.0596, 22.5429] }
 * 
 * // 格式4:GeoHash
 * { "location": "ws10k0dcg" }
 *
 * @param client
 * @throws Exception
 */
public static void searchNearbyWithSort(ElasticsearchClient client) throws Exception
{
    double lat = 22.5429;
    double lon = 114.0596;

    SearchResponse<Shop> response = client.search(s -> s
                    .index("shops")
                    .size(20)
                    .query(q -> q
                            .geoDistance(g -> g
                                    .field("location")
                                    .distance("5km")
                                    .location(loc -> loc
                                            .latlon(latlon -> latlon.lat(lat).lon(lon))
                                    )
                            )
                    )
                    .sort(so -> so
                            .geoDistance(g -> g
                                    .field("location")
                                    .location(loc -> loc
                                            .latlon(latlon -> latlon.lat(lat).lon(lon))
                                    )
                                    .order(SortOrder.Asc)
                                    .unit(DistanceUnit.Kilometers)
                            )
                    ),
            Shop.class
    );

    for(Hit<Shop> hit : response.hits().hits())
    {
        Shop shop = hit.source();
        // 排序值就是距离
        Double distance = hit.sort().get(0).doubleValue();
        logger.info("商家: {}, 距离: {} km", shop.getName(), distance);
    }
}

二十二、距离聚合(统计各距离段的商家)

/**
 * 距离聚合(统计各距离段的商家)
 *
 * @param client
 * @throws Exception
 */
public static void geoDistanceAggregation(ElasticsearchClient client) throws Exception
{
    double lat = 22.5429;
    double lon = 114.0596;

    SearchResponse<Void> response = client.search(s -> s
                    .index("shops")
                    .size(0)
                    .aggregations("shops_by_distance", a -> a
                            .geoDistance(g -> g
                                    .field("location")
                                    .origin(o -> o.latlon(latlon -> latlon.lat(lat).lon(lon)))
                                    .unit(DistanceUnit.Kilometers)
                                    .ranges(
                                            AggregationRange.of(rb -> rb.to(1.0)),
                                            AggregationRange.of(rb -> rb.from(1.0).to(3.0)),
                                            AggregationRange.of(rb -> rb.from(3.0).to(5.0)),
                                            AggregationRange.of(rb -> rb.from(5.0))
                                    )
                            )
                    ),
            Void.class
    );

    var buckets = response.aggregations()
            .get("shops_by_distance")
            .geoDistance()
            .buckets()
            .array();

    logger.info("=== 距离分布 ===");
    for(var bucket : buckets)
    {
        logger.info("  {}: {} 家", bucket.key(), bucket.docCount());
    }
}

二十三、删除文档

/**
 * 删除文档
 *
 * @param client
 * @throws Exception
 */
public static void deleteByCondition(ElasticsearchClient client) throws Exception
{
    // 构建删除条件
    Query.Builder builder = new Query.Builder();

    // 删除 brand = "Apple" 的文档
    builder.term(t -> t
            .field("brand")
            .value("Apple")
    );

    // 构建请求
    DeleteByQueryRequest request = new DeleteByQueryRequest.Builder()
            .index("products")
            .query(builder.build())
            .build();

    // 执行删除
    DeleteByQueryResponse response = client.deleteByQuery(request);
    logger.info("删除文档数: {}", response.deleted());
}