You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何使用OpenSearch Java Client指定search_pipeline?

使用OpenSearch Java Client指定Search Pipeline参数

核心实现方式

在OpenSearch Java Client中,你可以直接通过SearchRequest的searchPipeline()方法指定要调用的搜索流水线名称。

同步客户端示例

import org.opensearch.client.opensearch.OpenSearchClient;
import org.opensearch.client.opensearch.core.SearchRequest;
import org.opensearch.client.opensearch.core.SearchResponse;
import org.opensearch.client.opensearch.core.search.Query;

public class SearchPipelineDemo {
    public static void main(String[] args) throws Exception {
        // 假设已完成OpenSearchClient实例初始化
        OpenSearchClient client = initOpenSearchClient();

        // 构建查询条件
        Query query = new Query.Builder()
                .match(m -> m.field("title").query("技术文档"))
                .build();

        // 构建搜索请求并指定目标流水线
        SearchRequest request = new SearchRequest.Builder()
                .index("docs_index") // 替换为你的目标索引名
                .query(query)
                .searchPipeline("my_pipeline") // 这里设置要调用的流水线名称
                .build();

        // 执行搜索并处理结果
        SearchResponse<Object> response = client.search(request, Object.class);
        System.out.println("匹配结果总数: " + response.hits().total().value());
    }

    // 客户端初始化方法(根据你的集群环境调整实现)
    private static OpenSearchClient initOpenSearchClient() {
        // 此处省略具体初始化逻辑,可基于RestClientBuilder完成配置
        return null;
    }
}

异步客户端示例

如果使用异步客户端,写法逻辑一致,仅调用异步搜索方法:

import org.opensearch.client.opensearch.OpenSearchAsyncClient;
import org.opensearch.client.opensearch.core.SearchRequest;
import org.opensearch.client.opensearch.core.SearchResponse;

import java.util.concurrent.CompletableFuture;

public class AsyncSearchPipelineDemo {
    public static void main(String[] args) throws Exception {
        OpenSearchAsyncClient asyncClient = initAsyncOpenSearchClient();

        Query query = new Query.Builder()
                .match(m -> m.field("title").query("技术文档"))
                .build();

        SearchRequest request = new SearchRequest.Builder()
                .index("docs_index")
                .query(query)
                .searchPipeline("my_pipeline")
                .build();

        // 异步执行搜索并处理回调
        CompletableFuture<SearchResponse<Object>> future = asyncClient.search(request, Object.class);
        future.thenAccept(response -> {
            System.out.println("异步搜索结果总数: " + response.hits().total().value());
        }).join();
    }

    // 异步客户端初始化方法
    private static OpenSearchAsyncClient initAsyncOpenSearchClient() {
        // 此处省略具体初始化逻辑
        return null;
    }
}

注意事项

  • 确保指定的my_pipeline已在OpenSearch集群中创建并处于可用状态,否则会返回流水线不存在的错误。
  • 若需执行多个流水线,可传递多个名称给searchPipeline()方法,例如.searchPipeline("pipeline_a", "pipeline_b"),OpenSearch会按传入顺序依次执行。

内容的提问来源于stack exchange,提问作者AlexElin

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.14 02:22:36