import co.elastic.clients.elasticsearch._types.FieldValue;
import co.elastic.clients.elasticsearch._types.SortOptions;
import co.elastic.clients.elasticsearch._types.query_dsl.Query;
import co.elastic.clients.elasticsearch._types.query_dsl.RangeQuery;
import co.elastic.clients.elasticsearch.core.SearchRequest;
import co.elastic.clients.elasticsearch.core.search.Hit;
import co.elastic.clients.json.JsonData;
import com.gy.hsxt.elasticsearch.constant.SearchResultConstant;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import co.elastic.clients.elasticsearch.ElasticsearchClient;
import co.elastic.clients.elasticsearch.core.SearchResponse;
import co.elastic.clients.util.ObjectBuilder;
import org.springframework.stereotype.Service;

import java.math.BigDecimal;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.Function;
/**
 * @author Administrator
 * @Description:
 * @date 2025/10/13 9:26
 */
@Slf4j
@Service
public class ElasticsearchGenericSearchService {
    // 注入ES客户端(通用服务无需绑定具体索引,客户端全局复用)
    @Autowired
    private  ElasticsearchClient esClient;
    // 分页最大值
    private static final int PAGE_SIZE_MAX= 1000;

    private static final int PAGE_SIZE_MIN= BigDecimal.ONE.intValue();

    public <T> Map<String,Object> pageMapSearch(
            String indexName,
            int pageNo,
            int pageSize,
            List<SortOptions> sortOptionsList,
            Function<Query.Builder, ObjectBuilder<Query>> queryBuilder,
            Class<T> dataClass
    ) throws Exception {
        Map<String,Object> dataMap = new HashMap<>();
        // 1. 参数校验(通用逻辑,所有类型复用)
        validateParams(indexName, pageNo, pageSize, dataClass);

        // 2. 分页核心逻辑(Search After通用实现,与具体数据类型解耦)
        int targetPage = pageNo;
        int currentPage = 1;
        long total = 0;
        List<T> targetData = new ArrayList<>();
        AtomicReference<List<FieldValue>> searchAfterRef = new AtomicReference<>(null);
        boolean hasMore = true;
        while (hasMore && currentPage <= targetPage) {
            // 3. 构建通用查询请求(索引、排序、查询条件均由外部传入)
            SearchRequest request = buildGenericSearchRequest(
                    indexName, pageSize, searchAfterRef.get(), sortOptionsList, queryBuilder
            );

            // 4. 执行查询(泛型适配不同数据类型)
            SearchResponse<T> response = esClient.search(request, dataClass);
            log.info("查询结果:{}条 sql:{}",response.hits().total().value(),request.toString());
            // 5. 提取数据(通用逻辑)
            List<T> currentData = extractData(response);

            // 6. 目标页判断(通用逻辑)
            if (currentPage == targetPage) {
                targetData = currentData;
                total = response.hits().total().value();
            }

            // 7. 更新分页状态(通用逻辑)
            hasMore = updatePaginationState(response, pageSize, searchAfterRef);

            currentPage++;
        }
        dataMap.put(SearchResultConstant.DATA, targetData);
        // 获取总数
        dataMap.put(SearchResultConstant.TOTAL, total);

        return dataMap;
    }

    /**
     * 通用分页查询(Search After模式,支持任意索引、任意数据类型)
     * @param indexName 目标索引名(如"order_index"/"product_index")
     * @param pageNo 目标页码(从1开始)
     * @param pageSize 每页条数
     * @param sortOptionsList 排序条件构造器(外部传入,适配不同数据类型的排序字段)
     * @param queryBuilder 查询条件构造器(外部传入,适配不同业务的过滤逻辑)
     * @param dataClass 数据类型Class(如OrderIndex.class/ProductIndex.class)
     * @return 目标页数据列表(非null)
     * @param <T> 数据类型泛型
     * @throws Exception 查询异常(上层处理)
     */
    public <T> List<T> pageSearch(
            String indexName,
            int pageNo,
            int pageSize,
            List<SortOptions> sortOptionsList,
            Function<Query.Builder, ObjectBuilder<Query>> queryBuilder,
            Class<T> dataClass
    ) throws Exception {
        Map<String,Object> dataMap =  pageMapSearch(
                indexName,
                pageNo,
                pageSize,
                sortOptionsList,
                queryBuilder,
                dataClass
        );
        return (List<T>) dataMap.get(SearchResultConstant.DATA);
    }




    // ------------------------------ 通用工具方法(参数校验、请求构建等)------------------------------
    private <T> void validateParams(String indexName, int pageNo, int pageSize, Class<T> dataClass) {
        if (indexName == null || indexName.trim().isEmpty()) {
            throw new IllegalArgumentException("索引名indexName不能为空");
        }
        if (pageNo < PAGE_SIZE_MIN) {
            throw new IllegalArgumentException("页码pageNo必须≥1,当前值:" + pageNo);
        }
        if (pageSize < PAGE_SIZE_MIN || pageSize > PAGE_SIZE_MAX) {
            throw new IllegalArgumentException("页大小pageSize必须在1-1000之间,当前值:" + pageSize);
        }
        if (dataClass == null) {
            throw new IllegalArgumentException("数据类型Class不能为空");
        }
    }

    private <T> SearchRequest buildGenericSearchRequest(
            String indexName,
            int pageSize,
            List<FieldValue> searchAfter,
            List<SortOptions> sortOptionsList,
            Function<Query.Builder, ObjectBuilder<Query>> queryBuilder
    ) {
        return SearchRequest.of(s -> {
            s.index(indexName)
                    .query(queryBuilder) // 外部传入查询条件(如订单查商户ID、商品查分类)
                    .size(pageSize)
                    .from(0);
            if(sortOptionsList != null && !sortOptionsList.isEmpty()){
                s.sort(sortOptionsList);   // 外部传入排序条件
            }
            if (searchAfter != null && !searchAfter.isEmpty()) {
                s.searchAfter(searchAfter);
            }
            return s;
        });
    }



    private <T> List<T> extractData(SearchResponse<T> response) {
        if (response.hits() == null || response.hits().hits() == null) {
            return new ArrayList<>();
        }
        return response.hits().hits().stream()
                .map(hit -> hit.source())
                .filter(java.util.Objects::nonNull)
                .toList();
    }


    private <T> boolean updatePaginationState(
            SearchResponse<T> response,
            int pageSize,
            AtomicReference<List<FieldValue>> searchAfterRef
    ) {
        List<Hit<T>> hits = response.hits().hits();
        if (hits == null || hits.size() < pageSize) {
            return false; // 数据不足,无下一页
        }
        // 更新searchAfter为最后一条数据的排序值
        List<FieldValue> lastSort = hits.get(hits.size() - 1).sort();
        searchAfterRef.set(lastSort);
        return true;
    }



    /**
     * 添加时间范围查询条件到布尔查询构建器中
     * 支持开始时间和结束时间的范围查询
     *
     * @param fieldName 查询参数字段名
     * @param startTime 查询开始时间
     * @param endTime 查询结束时间
     * @see RangeQuery
     */
    public RangeQuery addTimeRangeQuery(String fieldName, Long startTime, Long endTime) {
        // 时间范围查询优化版
        if (startTime != null || endTime != null) {
            RangeQuery.Builder rangeBuilder = new RangeQuery.Builder();

            // 设置字段名
            rangeBuilder.field(fieldName);

            // 添加开始时间条件(大于等于)
            if (startTime != null) {
                rangeBuilder.gte(JsonData.of(startTime));
            }

            // 添加结束时间条件(小于等于)
            if (endTime != null) {
                rangeBuilder.lte(JsonData.of(endTime));
            }
            return rangeBuilder.build();
        }
        return null;
    }
}

调用示例:

@Slf4j
@Service
public class SearchOrderServiceImpl  {
  /**
     * 索引名称常量
     */
    private static final String INDEX_NAME = "order_index";


    @Autowired
    private ElasticsearchGenericSearchService elasticsearchGenericSearchService;


     public List<SearchOrderRespVO> pageQueryOrderListByEntCustId(AbstractOrderQuery params) {

        if(params.getCurPage() <= BigDecimal.ZERO.intValue())
            params.setCurPage(Constants.DEFAULT_PAGE_NUM);
        // 1. 定义排序字段(必须与 searchAfter 的值列表对应)
        List<SortOptions> sortOptions = new ArrayList<>();
        // 排序字段1:create_time(长整型)
        // 排序条件:先按 createDate 升序,再按 id.keyword 升序(订单专属排序)
        sortOptions.add(SortOptions.of(s -> s
                .field(f -> f.field("createDate").order(SortOrder.Desc))
        ));
        // 排序字段2:_id(字符串,确保唯一性)
        sortOptions.add(SortOptions.of(s -> s
                .field(f -> f.field("id.keyword").order(SortOrder.Desc))
        ));

        // 查询条件:商户ID=MERCHANT_001(订单专属过滤)
        Function<Query.Builder, ObjectBuilder<Query>> orderQuery = q ->{
            q.bool(b -> b
                    .must(m -> m.term(t -> t.field("userId").value(params.getUserId())))
                    .must(m -> m.term(t -> t.field("mobile").value(params.getMobile())))
                    .must(m -> m.range(elasticsearchGenericSearchService.addTimeRangeQuery("createDate", params.getStartTime(), params.getEndTime())))
            );
            return q;
        };


        List<SearchOrderRespVO> vos = new ArrayList<>();
        try {
            List<OrderIndex> list = elasticsearchGenericSearchService.pageSearch(
                    INDEX_NAME, params.getCurPage(), 10, // 第1页,每页10条
                    sortOptions, orderQuery,
                    OrderIndex.class // 订单数据类型
            );

            // 遍历结果
            list.forEach(order -> {
               System.out.println("count======="+order );
            });

        }catch (ElasticsearchException  e) {
            log.error("调用ES查询订单列表失败,参数:{}",params, e);
        } catch (Exception e) {
            log.error("查询订单列表失败,参数:{}",params, e);
        }

        return vos;
    }

}

pom.xml 依赖配置(JDK21)

<dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-data-elasticsearch</artifactId>
            <version>${spring-boot.version}</version>
        </dependency>
        <!-- Kafka -->
        <dependency>
            <groupId>org.springframework.kafka</groupId>
            <artifactId>spring-kafka</artifactId>
        </dependency>
        <dependency>
            <groupId>org.elasticsearch</groupId>
            <artifactId>elasticsearch</artifactId>
            <version>8.15.1</version>
        </dependency>


        <dependency>
            <groupId>org.elasticsearch.client</groupId>
            <artifactId>elasticsearch-rest-high-level-client</artifactId>
            <version>7.17.26</version>
        </dependency>

Logo

中国智能体开发者社区,聚焦智能体与大模型开发,提供前沿资讯、实用工具链、开源项目及行业案例。通过技术沙龙、开发者大赛等活动,促进经验交流与协作,助力开发者快速构建创新智能应用。

更多推荐