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());
}