Java Elasticsearch案例如何对接

wen java案例 26

本文目录导读:

Java Elasticsearch案例如何对接

  1. 环境准备
  2. 方式一:Spring Data Elasticsearch(推荐)
  3. 方式二:RestHighLevelClient(更灵活)
  4. 高级用法示例
  5. 性能优化建议
  6. 常见问题解决

我来为您详细讲解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();

常见问题解决

  1. 版本兼容性:确保客户端版本与ES服务器版本一致
  2. 连接超时:合理设置超时时间
  3. 索引不存在:操作前检查索引是否存在
  4. 数据类型不匹配:确保Java类型与ES mapping类型对应

这个完整的案例覆盖了Java对接Elasticsearch的主要方式和使用场景,您可以根据具体需求选择合适的方式。

抱歉,评论功能暂时关闭!