本文目录导读:

我来为您详细讲解Java对接Elasticsearch的完整案例,包括主流方式(Spring Data Elasticsearch 和 RestHighLevelClient)。
环境准备
1 Maven依赖配置
<!-- Spring Data Elasticsearch -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-elasticsearch</artifactId>
<version>2.7.0</version>
</dependency>
<!-- 或使用RestHighLevelClient -->
<dependency>
<groupId>org.elasticsearch.client</groupId>
<artifactId>elasticsearch-rest-high-level-client</artifactId>
<version>7.17.0</version>
</dependency>
<!-- Elasticsearch版本保持一致 -->
<dependency>
<groupId>org.elasticsearch</groupId>
<artifactId>elasticsearch</artifactId>
<version>7.17.0</version>
</dependency>
2 配置文件 (application.yml)
spring:
elasticsearch:
uris: http://localhost:9200
connection-timeout: 10s
socket-timeout: 30s
# 或者使用Spring Data Elasticsearch
data:
elasticsearch:
cluster-name: elasticsearch
cluster-nodes: localhost:9300
repositories:
enabled: true
方式一:Spring Data Elasticsearch(推荐)
1 实体类映射
package com.example.esdemo.entity;
import lombok.Data;
import org.springframework.data.annotation.Id;
import org.springframework.data.elasticsearch.annotations.Document;
import org.springframework.data.elasticsearch.annotations.Field;
import org.springframework.data.elasticsearch.annotations.FieldType;
import java.util.Date;
@Data
@Document(indexName = "user_index", shards = 3, replicas = 1)
public class UserDocument {
@Id
private Long id;
@Field(type = FieldType.Text, analyzer = "ik_max_word")
private String name;
@Field(type = FieldType.Integer)
private Integer age;
@Field(type = FieldType.Keyword)
private String email;
@Field(type = FieldType.Text, analyzer = "ik_max_word")
private String description;
@Field(type = FieldType.Date, format = DateFormat.date_time)
private Date createTime;
}
2 Repository接口
package com.example.esdemo.repository;
import com.example.esdemo.entity.UserDocument;
import org.springframework.data.elasticsearch.repository.ElasticsearchRepository;
import org.springframework.stereotype.Repository;
import java.util.List;
@Repository
public interface UserRepository extends ElasticsearchRepository<UserDocument, Long> {
// 根据名称模糊查询
List<UserDocument> findByNameLike(String name);
// 根据年龄范围查询
List<UserDocument> findByAgeBetween(Integer minAge, Integer maxAge);
// 复合查询
List<UserDocument> findByNameAndAge(String name, Integer age);
}
3 Service层实现
package com.example.esdemo.service;
import com.example.esdemo.entity.UserDocument;
import com.example.esdemo.repository.UserRepository;
import org.elasticsearch.index.query.QueryBuilders;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.domain.Page;
import org.springframework.data.domain.PageRequest;
import org.springframework.data.domain.Sort;
import org.springframework.data.elasticsearch.core.ElasticsearchRestTemplate;
import org.springframework.data.elasticsearch.core.SearchHits;
import org.springframework.data.elasticsearch.core.query.NativeSearchQueryBuilder;
import org.springframework.stereotype.Service;
import java.util.List;
import java.util.Optional;
@Service
public class UserService {
@Autowired
private UserRepository userRepository;
@Autowired
private ElasticsearchRestTemplate elasticsearchRestTemplate;
// 保存文档
public UserDocument save(UserDocument user) {
return userRepository.save(user);
}
// 批量保存
public Iterable<UserDocument> saveAll(List<UserDocument> users) {
return userRepository.saveAll(users);
}
// 根据ID查询
public Optional<UserDocument> findById(Long id) {
return userRepository.findById(id);
}
// 根据名称模糊查询
public List<UserDocument> findByName(String name) {
return userRepository.findByNameLike(name);
}
// 复杂查询示例
public Page<UserDocument> searchWithPagination(String keyword, int page, int size) {
NativeSearchQueryBuilder queryBuilder = new NativeSearchQueryBuilder();
// 构建多字段查询
queryBuilder.withQuery(QueryBuilders.multiMatchQuery(keyword, "name", "description"));
// 分页
queryBuilder.withPageable(PageRequest.of(page, size));
// 排序
queryBuilder.withSort(Sort.by(Sort.Direction.DESC, "createTime"));
// 高亮设置
// ...
// 执行查询
SearchHits<UserDocument> searchHits = elasticsearchRestTemplate.search(
queryBuilder.build(), UserDocument.class);
return searchHits.getSearchHits().stream()
.map(hit -> hit.getContent())
.collect(Collectors.toList());
}
}
方式二:RestHighLevelClient(更灵活)
1 客户端配置
package com.example.esdemo.config;
import org.apache.http.HttpHost;
import org.elasticsearch.client.RestClient;
import org.elasticsearch.client.RestHighLevelClient;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@Configuration
public class ElasticsearchConfig {
@Bean
public RestHighLevelClient restHighLevelClient() {
return new RestHighLevelClient(
RestClient.builder(
new HttpHost("localhost", 9200, "http")
// 可以添加多个节点
// new HttpHost("localhost", 9201, "http")
)
);
}
}
2 ES操作工具类
package com.example.esdemo.util;
import com.alibaba.fastjson.JSON;
import org.elasticsearch.action.admin.indices.delete.DeleteIndexRequest;
import org.elasticsearch.action.bulk.BulkRequest;
import org.elasticsearch.action.bulk.BulkResponse;
import org.elasticsearch.action.delete.DeleteRequest;
import org.elasticsearch.action.delete.DeleteResponse;
import org.elasticsearch.action.get.GetRequest;
import org.elasticsearch.action.get.GetResponse;
import org.elasticsearch.action.index.IndexRequest;
import org.elasticsearch.action.index.IndexResponse;
import org.elasticsearch.action.search.SearchRequest;
import org.elasticsearch.action.search.SearchResponse;
import org.elasticsearch.action.update.UpdateRequest;
import org.elasticsearch.action.update.UpdateResponse;
import org.elasticsearch.client.RequestOptions;
import org.elasticsearch.client.RestHighLevelClient;
import org.elasticsearch.client.indices.CreateIndexRequest;
import org.elasticsearch.client.indices.GetIndexRequest;
import org.elasticsearch.common.xcontent.XContentType;
import org.elasticsearch.index.query.QueryBuilders;
import org.elasticsearch.search.SearchHit;
import org.elasticsearch.search.builder.SearchSourceBuilder;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
@Component
public class ElasticsearchUtil {
@Autowired
private RestHighLevelClient client;
// 创建索引
public boolean createIndex(String index) throws Exception {
CreateIndexRequest request = new CreateIndexRequest(index);
client.indices().create(request, RequestOptions.DEFAULT);
return true;
}
// 判断索引是否存在
public boolean existsIndex(String index) throws Exception {
GetIndexRequest request = new GetIndexRequest(index);
return client.indices().exists(request, RequestOptions.DEFAULT);
}
// 删除索引
public boolean deleteIndex(String index) throws Exception {
DeleteIndexRequest request = new DeleteIndexRequest(index);
client.indices().delete(request, RequestOptions.DEFAULT);
return true;
}
// 添加文档
public String addDocument(String index, String id, Object object) throws Exception {
IndexRequest request = new IndexRequest(index);
request.id(id);
request.source(JSON.toJSONString(object), XContentType.JSON);
IndexResponse response = client.index(request, RequestOptions.DEFAULT);
return response.getId();
}
// 批量添加文档
public boolean batchAddDocument(String index, List<Map<String, Object>> documents) throws Exception {
BulkRequest request = new BulkRequest();
for (Map<String, Object> doc : documents) {
request.add(new IndexRequest(index)
.source(JSON.toJSONString(doc), XContentType.JSON));
}
BulkResponse response = client.bulk(request, RequestOptions.DEFAULT);
return !response.hasFailures();
}
// 查询文档
public Map<String, Object> getDocument(String index, String id) throws Exception {
GetRequest request = new GetRequest(index, id);
GetResponse response = client.get(request, RequestOptions.DEFAULT);
return response.getSource();
}
// 更新文档
public boolean updateDocument(String index, String id, Map<String, Object> document) throws Exception {
UpdateRequest request = new UpdateRequest(index, id);
request.doc(JSON.toJSONString(document), XContentType.JSON);
UpdateResponse response = client.update(request, RequestOptions.DEFAULT);
return response.getResult().getClass().equals(UpdateResponse.Result.UPDATED);
}
// 删除文档
public boolean deleteDocument(String index, String id) throws Exception {
DeleteRequest request = new DeleteRequest(index, id);
DeleteResponse response = client.delete(request, RequestOptions.DEFAULT);
return response.getResult().getClass().equals(DeleteResponse.Result.DELETED);
}
// 搜索文档
public List<Map<String, Object>> searchDocument(String index, String field, String keyword) throws Exception {
SearchRequest request = new SearchRequest(index);
SearchSourceBuilder sourceBuilder = new SearchSourceBuilder();
// 构建查询
sourceBuilder.query(QueryBuilders.matchQuery(field, keyword));
sourceBuilder.from(0);
sourceBuilder.size(10);
request.source(sourceBuilder);
SearchResponse response = client.search(request, RequestOptions.DEFAULT);
List<Map<String, Object>> resultList = new ArrayList<>();
for (SearchHit hit : response.getHits().getHits()) {
resultList.add(hit.getSourceAsMap());
}
return resultList;
}
// 复杂搜索(聚合查询等)
public SearchResponse complexSearch(SearchRequest searchRequest) throws Exception {
return client.search(searchRequest, RequestOptions.DEFAULT);
}
}
3 控制器示例
package com.example.esdemo.controller;
import com.example.esdemo.entity.UserDocument;
import com.example.esdemo.service.UserService;
import com.example.esdemo.util.ElasticsearchUtil;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.*;
import java.util.List;
import java.util.Map;
@RestController
@RequestMapping("/es")
public class EsController {
@Autowired
private UserService userService;
@Autowired
private ElasticsearchUtil elasticsearchUtil;
// 保存用户
@PostMapping("/user")
public UserDocument saveUser(@RequestBody UserDocument user) {
return userService.save(user);
}
// 查询用户
@GetMapping("/user/{id}")
public UserDocument getUser(@PathVariable Long id) {
return userService.findById(id).orElse(null);
}
// 搜索用户
@GetMapping("/user/search")
public List<Map<String, Object>> searchUser(
@RequestParam String keyword,
@RequestParam(defaultValue = "0") int page,
@RequestParam(defaultValue = "10") int size) throws Exception {
return elasticsearchUtil.searchDocument("user_index", "name", keyword);
}
// 批量导入
@PostMapping("/user/batch")
public boolean batchImport(@RequestBody List<Map<String, Object>> users) throws Exception {
return elasticsearchUtil.batchAddDocument("user_index", users);
}
}
高级用法示例
1 聚合查询
public Map<String, Long> aggregateByField(String index, String field) throws Exception {
SearchRequest searchRequest = new SearchRequest(index);
SearchSourceBuilder searchSourceBuilder = new SearchSourceBuilder();
// 创建聚合
AggregationBuilder aggregation = AggregationBuilders
.terms("by_" + field)
.field(field)
.size(10);
searchSourceBuilder.aggregation(aggregation);
searchRequest.source(searchSourceBuilder);
SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT);
// 解析聚合结果
Terms terms = response.getAggregations().get("by_" + field);
Map<String, Long> result = new HashMap<>();
for (Terms.Bucket bucket : terms.getBuckets()) {
result.put(bucket.getKeyAsString(), bucket.getDocCount());
}
return result;
}
2 地理空间查询
public List<Map<String, Object>> geoSearch(String index,
double lat, double lon,
double distance) throws Exception {
SearchRequest searchRequest = new SearchRequest(index);
SearchSourceBuilder searchSourceBuilder = new SearchSourceBuilder();
// 地理距离查询
GeoDistanceQueryBuilder geoQuery = QueryBuilders
.geoDistanceQuery("location")
.point(lat, lon)
.distance(distance, DistanceUnit.KILOMETERS);
searchSourceBuilder.query(geoQuery);
searchRequest.source(searchSourceBuilder);
SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT);
List<Map<String, Object>> result = new ArrayList<>();
for (SearchHit hit : response.getHits().getHits()) {
result.add(hit.getSourceAsMap());
}
return result;
}
性能优化建议
1 连接池配置
@Bean
public RestHighLevelClient restHighLevelClient() {
RestClientBuilder builder = RestClient.builder(
new HttpHost("localhost", 9200, "http")
);
// 设置连接池
builder.setRequestConfigCallback(requestConfigBuilder ->
requestConfigBuilder
.setConnectTimeout(5000)
.setSocketTimeout(60000)
.setConnectionRequestTimeout(5000)
);
// 设置节点选择器
builder.setNodeSelector(NodeSelector.SKIP_DEDICATED_MASTERS);
return new RestHighLevelClient(builder);
}
2 批量操作优化
// 使用bulkProcessor进行批量操作
BulkProcessor bulkProcessor = BulkProcessor.builder(
(request, bulkListener) -> client.bulkAsync(request, RequestOptions.DEFAULT, bulkListener),
new BulkProcessor.Listener() {
@Override
public void beforeBulk(long executionId, BulkRequest request) {
// 批量执行前
}
@Override
public void afterBulk(long executionId, BulkRequest request, BulkResponse response) {
// 批量执行后
}
@Override
public void afterBulk(long executionId, BulkRequest request, Throwable failure) {
// 批量执行失败
log.error("Bulk operation failed", failure);
}
})
.setBulkActions(1000) // 每1000条批量执行一次
.setBulkSize(new ByteSizeValue(5, ByteSizeUnit.MB)) // 5MB大小
.setFlushInterval(TimeValue.timeValueSeconds(5)) // 每5秒刷新
.setConcurrentRequests(1)
.build();
常见问题解决
- 版本兼容性:确保客户端版本与ES服务器版本一致
- 连接超时:合理设置超时时间
- 索引不存在:操作前检查索引是否存在
- 数据类型不匹配:确保Java类型与ES mapping类型对应
这个完整的案例覆盖了Java对接Elasticsearch的主要方式和使用场景,您可以根据具体需求选择合适的方式。