全部产品
Search
文档中心

云原生数据库 PolarDB:使用客户端连接PolarSearch

更新时间:Apr 20, 2026

PolarSearch完全兼容OpenSearch的官方客户端,您可以使用相应的客户端与PolarSearch进行交互。这使得您能够通过Java或Python等常用编程语言高效地管理索引、操作文档(包括增、删、改、查)以及执行复杂搜索,从而将搜索功能无缝集成到您的应用程序中。

准备工作

  1. 已创建含有PolarSearch节点的集群并设置了节点的管理员账号

  2. 获取连接地址:在集群的数据库节点区域,将鼠标悬浮在搜索节点,根据您的业务环境,获取PolarSearch节点的私网或公网地址。

连接集群

OpenSearch Java client

OpenSearch Java客户端允许您通过Java方法和数据结构与OpenSearch集群交互,而不是使用HTTP方法和原始JSON。比如,您可以使用对象向集群提交请求,通过客户端内置的方法来创建索引、向文档写入数据,或完成其他操作。有关该客户端完整的API文档和更多示例,请参阅javadoc

1. 选择传输层并添加依赖

OpenSearch Java Client需要搭配一个传输层框架来处理HTTP请求。

  • PolarSearch 1.x版本:Java Client仅能使用RestClient Transport。

  • PolarSearch 3.x版本:Java Client则可以使用Apache HttpClient 5 Transport或者RestClient Transport。

PolarSearch 1.x版本

Maven示例:pom.xml文件中添加以下依赖:

<!-- OpenSearch Java Client核心库 -->
<dependency>
    <groupId>org.opensearch.client</groupId>
    <artifactId>opensearch-java</artifactId>
    <version>1.0.0</version>
</dependency>
<!-- RestClient Transport 传输层 -->
<dependency>
    <groupId>org.opensearch.client</groupId>
    <artifactId>opensearch-rest-client</artifactId>
    <version>1.3.20</version>
</dependency>

PolarSearch 3.x版本

Apache HttpClient5 Transport Maven示例:pom.xml文件中添加以下依赖:

<!-- OpenSearch Java Client核心库 -->
<dependency>
    <groupId>org.opensearch.client</groupId>
    <artifactId>opensearch-java</artifactId>
    <version>3.3.0</version>
</dependency>
<!-- Apache HttpClient5 传输层 -->
<dependency>
    <groupId>org.apache.httpcomponents.client5</groupId>
    <artifactId>httpclient5</artifactId>
    <version>5.4.3</version>
</dependency>

2. 设置数据类

创建一个测试数据类,以供后续测试PolarSearch功能使用。

static class IndexData {
    private String title;
    private String text;

    public IndexData() {}

    public IndexData(String title, String text) {
        this.title = title;
        this.text = text;
    }

    public String getTitle() {
        return title;
    }

    public void setTitle(String title) {
        this.title = title;
    }

    public String getText() {
        return text;
    }

    public void setText(String text) {
        this.text = text;
    }

    @Override
    public String toString() {
        return String.format("IndexData{title='%s', text='%s'}", title, text);
    }
}

3. 初始化客户端

根据您选择的传输层,使用PolarSearch的连接信息初始化客户端。在实际使用中,您可以根据以下代码配置将PolarSearch的相关信息设置为环境变量,或直接将其赋值为参数的默认值,以便于测试使用。

PolarSearch 1.x版本

以下示例展示了如何使用RestClient Transport初始化一个客户端,当前示例代码禁用了SSL。

public static OpenSearchTransport createTransport() throws Exception {
    var env = System.getenv();
    var hostname = env.getOrDefault("HOST", "<polarsearch_host>");
    var port = Integer.parseInt(env.getOrDefault("PORT", "<polarsearch_port>"));
    var scheme = env.getOrDefault("SCHEME", "http");
    var user = env.getOrDefault("USERNAME", "<polarsearch_username>");
    var pass = env.getOrDefault("PASSWORD", "<polarsearch_password>");

    final HttpHost host = new HttpHost(hostname, port, scheme);
    final BasicCredentialsProvider credentialsProvider = new BasicCredentialsProvider();
    credentialsProvider.setCredentials(AuthScope.ANY,
            new UsernamePasswordCredentials(user, pass));

    final SSLContext sslContext = SSLContextBuilder.create()
            .loadTrustMaterial(null, (chains, authType) -> true)
            .build();

    RestClientBuilder builder = RestClient.builder(host)
            .setHttpClientConfigCallback(httpClientBuilder ->
                    httpClientBuilder
                            .setDefaultCredentialsProvider(credentialsProvider)
                            .setSSLContext(sslContext));

    final RestClient restClient = builder.build();
    return new RestClientTransport(restClient, new JacksonJsonpMapper());
}

transport = createTransport();
var client = new OpenSearchClient(transport);

PolarSearch 3.x版本

以下示例展示了如何使用Apache HttpClient5 Transport初始化一个客户端,当前示例代码禁用了SSL。

public static OpenSearchTransport createTransport() throws NoSuchAlgorithmException, KeyStoreException, KeyManagementException {
    var env = System.getenv();
    var hostname = env.getOrDefault("HOST", "<polarsearch_host>");
    var port = Integer.parseInt(env.getOrDefault("PORT", "<polarsearch_port>"));
    var scheme = env.getOrDefault("SCHEME", "http");
    var user = env.getOrDefault("USERNAME", "<polarsearch_username>");
    var pass = env.getOrDefault("PASSWORD", "<polarsearch_password>");

    final HttpHost host = new HttpHost(scheme, hostname, port);
    final BasicCredentialsProvider credentialsProvider = new BasicCredentialsProvider();
    credentialsProvider.setCredentials(new AuthScope(host), new UsernamePasswordCredentials(user, pass.toCharArray()));

    final SSLContext sslContext = SSLContextBuilder
            .create()
            .loadTrustMaterial(null, (chains, authType) -> true)
            .build();

    final ApacheHttpClient5TransportBuilder builder = ApacheHttpClient5TransportBuilder.builder(host);
    builder.setHttpClientConfigCallback(httpClientBuilder -> {
        final TlsStrategy tlsStrategy = ClientTlsStrategyBuilder.create()
                .setSslContext(sslContext)
                .setTlsDetailsFactory(new Factory<SSLEngine, TlsDetails>() {
                    @Override
                    public TlsDetails create(final SSLEngine sslEngine) {
                        return new TlsDetails(sslEngine.getSession(), sslEngine.getApplicationProtocol());
                    }
                })
                .build();

        final PoolingAsyncClientConnectionManager connectionManager = PoolingAsyncClientConnectionManagerBuilder
                .create()
                .setTlsStrategy(tlsStrategy)
                .build();

        return httpClientBuilder
                .setDefaultCredentialsProvider(credentialsProvider)
                .setConnectionManager(connectionManager);
    });

    final OpenSearchTransport transport = builder.build();
    return transport;
}

transport = createTransport();
var client = new OpenSearchClient(transport);

4. 执行基本操作

以下代码片段展示了如何执行常见的索引和文档操作。

创建索引

// 1. Creating an index
final var index = "my-index";
if (!client.indices().exists(r -> r.index(index)).value()) {
    CreateIndexRequest createIndexRequest = new CreateIndexRequest.Builder().index(index)
            .build();
    client.indices().create(createIndexRequest);
}

写入文档

// 2. Indexing data
IndexData indexData = new IndexData("first_name", "Bruce");
IndexRequest<IndexData> indexRequest = new IndexRequest.Builder<IndexData>().index(index).id("1").document(indexData).build();
client.index(indexRequest);

搜索文档

// 3. Searching for documents
SearchResponse<IndexData> searchResponse = client.search(s -> s.index(index), IndexData.class);
for (int i = 0; i< searchResponse.hits().hits().size(); i++) {
    LOGGER.info(searchResponse.hits().hits().get(i).source());
}        

删除文档

// 4. Deleting a document
client.delete(b -> b.index(index).id("1"));

删除索引

// 5. Deleting an index
DeleteIndexRequest deleteIndexRequest = new DeleteIndexRequest.Builder().index(index).build();
DeleteIndexResponse deleteIndexResponse = client.indices().delete(deleteIndexRequest);

OpenSearch Java high-level REST client

建议您使用OpenSearch Java client,因为Java high-level REST client已在OpenSearch中被废弃,未来版本将不再支持该客户端。

1. 添加依赖

pom.xml文件中添加以下依赖:

PolarSearch 1.x版本

<!-- OpenSearch Java rest high level client核心库 -->
<dependency>
    <groupId>org.opensearch.client</groupId>
    <artifactId>opensearch-rest-high-level-client</artifactId>
    <version>1.3.20</version>
</dependency>

PolarSearch 3.x版本

<!-- OpenSearch Java rest high level client核心库 -->
<dependency>
    <groupId>org.opensearch.client</groupId>
    <artifactId>opensearch-rest-high-level-client</artifactId>
    <version>3.3.2</version>
</dependency>

2. 初始化客户端

在实际使用中,您可以根据以下代码配置将PolarSearch的相关信息设置为环境变量,或直接将其赋值为参数的默认值,以便于测试使用。

PolarSearch 1.x版本

以下示例展示了如何初始化一个1.x版本的客户端,当前示例代码禁用了SSL。

public static RestHighLevelClient createClient() throws Exception {
    var env = System.getenv();
    var hostname = env.getOrDefault("HOST", "<polarsearch_host>");
    var port = Integer.parseInt(env.getOrDefault("PORT", "<polarsearch_port>"));
    var scheme = env.getOrDefault("SCHEME", "http");
    var user = env.getOrDefault("USERNAME", "<polarsearch_username>");
    var pass = env.getOrDefault("PASSWORD", "<polarsearch_password>");

    final HttpHost host = new HttpHost(hostname, port, scheme);
    final BasicCredentialsProvider credentialsProvider = new BasicCredentialsProvider();
    credentialsProvider.setCredentials(AuthScope.ANY,
            new UsernamePasswordCredentials(user, pass));

    final SSLContext sslContext = SSLContextBuilder.create()
            .loadTrustMaterial(null, (chains, authType) -> true)
            .build();

    RestClientBuilder builder = RestClient.builder(host)
            .setHttpClientConfigCallback(httpClientBuilder ->
                    httpClientBuilder
                            .setDefaultCredentialsProvider(credentialsProvider)
                            .setSSLContext(sslContext));

    return new RestHighLevelClient(builder);
}

RestHighLevelClient client = createClient();

PolarSearch 3.x版本

以下示例展示了如何初始化一个3.x版本的客户端,当前示例代码禁用了SSL。

public static RestHighLevelClient createClient() throws Exception {
    var env = System.getenv();
    var hostname = env.getOrDefault("HOST", "<polarsearch_host>");
    var port = Integer.parseInt(env.getOrDefault("PORT", "<polarsearch_port>"));
    var scheme = env.getOrDefault("SCHEME", "http");
    var user = env.getOrDefault("USERNAME", "<polarsearch_username>");
    var pass = env.getOrDefault("PASSWORD", "<polarsearch_password>");

    final HttpHost host = new HttpHost(scheme, hostname, port);
    final BasicCredentialsProvider credentialsProvider = new BasicCredentialsProvider();
    credentialsProvider.setCredentials(new AuthScope(host),
            new UsernamePasswordCredentials(user, pass.toCharArray()));

    final SSLContext sslContext = SSLContextBuilder.create()
            .loadTrustMaterial(null, (chains, authType) -> true)
            .build();

    final TlsStrategy tlsStrategy = ClientTlsStrategyBuilder.create()
            .setSslContext(sslContext)
            .build();

    final PoolingAsyncClientConnectionManager connectionManager = PoolingAsyncClientConnectionManagerBuilder
            .create()
            .setTlsStrategy(tlsStrategy)
            .build();

    RestClientBuilder builder = RestClient.builder(host)
            .setHttpClientConfigCallback(httpClientBuilder ->
                    httpClientBuilder
                            .setDefaultCredentialsProvider(credentialsProvider)
                            .setConnectionManager(connectionManager));

    return new RestHighLevelClient(builder);
}

RestHighLevelClient client = createClient();

3. 执行基本操作

以下代码片段展示了如何执行常见的索引和文档操作。

创建索引

// 1. Creating an index
final var index = "my-index";
boolean exists = client.indices().exists(new GetIndexRequest(index), RequestOptions.DEFAULT);
if (!exists) {
    CreateIndexRequest createIndexRequest = new CreateIndexRequest(index);
    client.indices().create(createIndexRequest, RequestOptions.DEFAULT);
}

写入文档

// 2. Indexing data
IndexData indexData = new IndexData("first_name", "Bruce");
Map<String, Object> document = new HashMap<>();
document.put("title", indexData.getTitle());
document.put("text", indexData.getText());
IndexRequest indexRequest = new IndexRequest(index).id("1").source(document);
client.index(indexRequest, RequestOptions.DEFAULT);

搜索文档

// 3. Searching for documents
SearchRequest searchRequest = new SearchRequest(index);
searchRequest.source(new SearchSourceBuilder().query(QueryBuilders.matchAllQuery()));
SearchResponse searchResponse = client.search(searchRequest, RequestOptions.DEFAULT);
for (SearchHit hit : searchResponse.getHits().getHits()) {
    Map<String, Object> sourceMap = hit.getSourceAsMap();
    IndexData data = new IndexData((String) sourceMap.get("title"), (String) sourceMap.get("text"));
    LOGGER.info(data);
}

删除文档

// 4. Deleting a document
client.delete(new DeleteRequest(index, "1"), RequestOptions.DEFAULT);

删除索引

// 5. Deleting an index
DeleteIndexRequest deleteIndexRequest = new DeleteIndexRequest(index);
AcknowledgedResponse deleteIndexResponse = client.indices().delete(deleteIndexRequest, RequestOptions.DEFAULT);

ElasticSearch Java High Level REST Client

PolarSearch 1.x版本100%兼容ElasticSearch Java High Level REST Client的7.0.0 ~ 7.13.4的所有版本,若使用的客户端版本为上述兼容版本,则可以直接修改连接地址,而无需改造任何业务代码。

1. 添加依赖

pom.xml文件中添加以下依赖:

<!-- ElasticSearch Java rest high level client核心库 -->
<dependency>
    <groupId>org.elasticsearch.client</groupId>
    <artifactId>elasticsearch-rest-high-level-client</artifactId>
    <version>7.13.4</version>
</dependency>

2. 初始化客户端

在实际使用中,您可以根据以下代码配置将PolarSearch的相关信息设置为环境变量,或直接将其赋值为参数的默认值,以便于测试使用。

以下示例展示了如何初始化一个1.x版本的客户端,当前示例代码禁用了SSL。

public static RestHighLevelClient createClient() throws Exception {
    var env = System.getenv();
    var hostname = env.getOrDefault("HOST", "<polarsearch_host>");
    var port = Integer.parseInt(env.getOrDefault("PORT", "<polarsearch_port>"));
    var scheme = env.getOrDefault("SCHEME", "http");
    var user = env.getOrDefault("USERNAME", "<polarsearch_username>");
    var pass = env.getOrDefault("PASSWORD", "<polarsearch_password>");

    final HttpHost host = new HttpHost(hostname, port, scheme);
    final BasicCredentialsProvider credentialsProvider = new BasicCredentialsProvider();
    credentialsProvider.setCredentials(AuthScope.ANY,
            new UsernamePasswordCredentials(user, pass));

    final SSLContext sslContext = SSLContextBuilder.create()
            .loadTrustMaterial(null, (chains, authType) -> true)
            .build();

    RestClientBuilder builder = RestClient.builder(host)
            .setHttpClientConfigCallback(httpClientBuilder ->
                    httpClientBuilder
                            .setDefaultCredentialsProvider(credentialsProvider)
                            .setSSLContext(sslContext));

    return new RestHighLevelClient(builder);
}

RestHighLevelClient client = createClient();

3. 执行基本操作

以下代码片段展示了如何执行常见的索引和文档操作。

创建索引

// 1. Creating an index
final var index = "my-index";
boolean exists = client.indices().exists(new GetIndexRequest(index), RequestOptions.DEFAULT);
if (!exists) {
    CreateIndexRequest createIndexRequest = new CreateIndexRequest(index);
    client.indices().create(createIndexRequest, RequestOptions.DEFAULT);
}

写入文档

// 2. Indexing data
IndexData indexData1 = new IndexData("first_name", "Bruce");
Map<String, Object> document1 = new HashMap<>();
document1.put("title", indexData1.getTitle());
document1.put("text", indexData1.getText());
client.index(new IndexRequest(index).id("1").source(document1), RequestOptions.DEFAULT);

搜索文档

// 3. Searching for documents
SearchRequest searchRequest = new SearchRequest(index);
searchRequest.source(new SearchSourceBuilder().query(QueryBuilders.matchAllQuery()));
SearchResponse searchResponse = client.search(searchRequest, RequestOptions.DEFAULT);
for (SearchHit hit : searchResponse.getHits().getHits()) {
    Map<String, Object> sourceMap = hit.getSourceAsMap();
    IndexData data = new IndexData((String) sourceMap.get("title"), (String) sourceMap.get("text"));
    LOGGER.info(data);
}

删除文档

// 4. Deleting a document
client.delete(new DeleteRequest(index, "1"), RequestOptions.DEFAULT);

删除索引

// 5. Deleting an index
DeleteIndexRequest deleteIndexRequest = new DeleteIndexRequest(index);
AcknowledgedResponse deleteIndexResponse = client.indices().delete(deleteIndexRequest, RequestOptions.DEFAULT);

Low-level Python client

建议您使用Low-level Python client,因为High-level Python client在OpenSearch 2.1.0版本后已被废弃。

1.配置环境

  1. 根据实际业务需求,进入指定的项目目录。此处以/home/PolarSearchTestPython为例。

    mkdir /home/PolarSearchTestPython
    cd /home/PolarSearchTestPython
  2. /home/PolarSearchTestPython目录中,创建虚拟环境(venv)隔离项目依赖,避免全局污染。

    python3 -m venv myenv
  3. 激活虚拟环境

    source myenv/bin/activate
  4. 安装所需的Python依赖库。

    pip3 install opensearch-py

2. 连接PolarSearch

在您的Python代码中导入OpenSearch类,并使用PolarSearch的连接信息创建客户端实例。在实际使用中,您可以根据以下代码配置将PolarSearch的相关信息设置为环境变量,或直接将其赋值为参数的默认值,以便于测试使用。

from opensearchpy import OpenSearch

host = os.getenv("HOST", default="<polarsearch_host>")
port = int(os.getenv("PORT", <polarsearch_port>))
auth = (os.getenv("USERNAME", "<polarsearch_username>"), os.getenv("PASSWORD", "<polarsearch_password>"))

client = OpenSearch(
    hosts=[{"host": host, "port": port}],
    http_auth=auth,
    use_ssl=False,
    verify_certs=False,
    ssl_show_warn=False,
)

3. 操作示例

以下代码片段展示了如何执行常见的索引和文档操作。

创建索引

使用client.indices.create()方法创建一个新索引。

index_name = 'python-test-index'
index_body = {
  'settings': {
    'index': {
      'number_of_shards': 4
    }
  }
}

response = client.indices.create(index=index_name, body=index_body)

写入文档

使用client.index()方法向指定索引中添加一个文档。

document = {
  'title': 'Moneyball',
  'director': 'Bennett Miller',
  'year': '2011'
}

response = client.index(
    index = 'python-test-index',
    body = document,
    id = '1',
    refresh = True
)

搜索文档

使用client.search()方法根据查询条件搜索文档。

q = 'miller'
query = {
  'size': 5,
  'query': {
    'multi_match': {
      'query': q,
      'fields': ['title^2', 'director']
    }
  }
}

response = client.search(
    body = query,
    index = 'python-test-index'
)

删除文档

使用client.delete()方法删除指定ID的文档。

response = client.delete(
    index = 'python-test-index',
    id = '1'
)

删除索引

使用client.indices.delete()方法删除整个索引。

response = client.indices.delete(
    index = 'python-test-index'
)

完整示例代码

下面是一个完整的示例,演示了从创建索引到最终删除索引的全过程。

OpenSearch Java client

示例代码

PolarSearch 1.x版本

依赖配置(pom.xml)

<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
         xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
    <modelVersion>4.0.0</modelVersion>

    <groupId>org.example</groupId>
    <artifactId>OpenSearchJavaClientSample</artifactId>
    <version>1.0-SNAPSHOT</version>

    <properties>
        <maven.compiler.source>21</maven.compiler.source>
        <maven.compiler.target>21</maven.compiler.target>
        <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
    </properties>
    <dependencies>
        <dependency>
            <groupId>org.opensearch.client</groupId>
            <artifactId>opensearch-java</artifactId>
            <version>1.0.0</version>
        </dependency>
        <dependency>
            <groupId>org.opensearch.client</groupId>
            <artifactId>opensearch-rest-client</artifactId>
            <version>1.3.20</version>
        </dependency>
        <dependency>
            <groupId>org.apache.logging.log4j</groupId>
            <artifactId>log4j-core</artifactId>
            <version>2.21.0</version>
        </dependency>
    </dependencies>
</project>

示例程序(OpenSearchClientExample.java)

package samples;

import java.io.IOException;

import javax.net.ssl.SSLContext;

import org.apache.http.HttpHost;
import org.apache.http.auth.AuthScope;
import org.apache.http.auth.UsernamePasswordCredentials;
import org.apache.http.impl.client.BasicCredentialsProvider;
import org.apache.http.ssl.SSLContextBuilder;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import org.opensearch.client.RestClient;
import org.opensearch.client.RestClientBuilder;
import org.opensearch.client.json.jackson.JacksonJsonpMapper;
import org.opensearch.client.opensearch.OpenSearchClient;
import org.opensearch.client.opensearch.core.IndexRequest;
import org.opensearch.client.opensearch.core.SearchResponse;
import org.opensearch.client.opensearch.indices.CreateIndexRequest;
import org.opensearch.client.opensearch.indices.DeleteIndexRequest;
import org.opensearch.client.opensearch.indices.DeleteIndexResponse;
import org.opensearch.client.transport.OpenSearchTransport;
import org.opensearch.client.transport.rest_client.RestClientTransport;


public class OpenSearchClientExample {
    static class IndexData {
        private String title;
        private String text;
    
        public IndexData() {}
    
        public IndexData(String title, String text) {
            this.title = title;
            this.text = text;
        }
    
        public String getTitle() {
            return title;
        }
    
        public void setTitle(String title) {
            this.title = title;
        }
    
        public String getText() {
            return text;
        }
    
        public void setText(String text) {
            this.text = text;
        }
    
        @Override
        public String toString() {
            return String.format("IndexData{title='%s', text='%s'}", title, text);
        }
    }
    
    private static final Logger LOGGER = LogManager.getLogger(OpenSearchClientExample.class);

    public static OpenSearchTransport createTransport() throws Exception {
        var env = System.getenv();
        var hostname = env.getOrDefault("HOST", "<polarsearch_host>");
        var port = Integer.parseInt(env.getOrDefault("PORT", "<polarsearch_port>"));
        var scheme = env.getOrDefault("SCHEME", "http");
        var user = env.getOrDefault("USERNAME", "<polarsearch_username>");
        var pass = env.getOrDefault("PASSWORD", "<polarsearch_password>");

        final HttpHost host = new HttpHost(hostname, port, scheme);
        final BasicCredentialsProvider credentialsProvider = new BasicCredentialsProvider();
        credentialsProvider.setCredentials(AuthScope.ANY,
                new UsernamePasswordCredentials(user, pass));

        final SSLContext sslContext = SSLContextBuilder.create()
                .loadTrustMaterial(null, (chains, authType) -> true)
                .build();

        RestClientBuilder builder = RestClient.builder(host)
                .setHttpClientConfigCallback(httpClientBuilder ->
                        httpClientBuilder
                                .setDefaultCredentialsProvider(credentialsProvider)
                                .setSSLContext(sslContext));

        final RestClient restClient = builder.build();
        return new RestClientTransport(restClient, new JacksonJsonpMapper());
    }

    public static void main(String[] args) {
        OpenSearchTransport transport = null;
        try {
            LOGGER.info("start test...");

            transport = createTransport();
            var client = new OpenSearchClient(transport);
            LOGGER.info("client create");
            final var index = "my-index";
            // 1. Creating an index
            LOGGER.info("1. Creating an index {} start", index);
            if (!client.indices().exists(r -> r.index(index)).value()) {
                CreateIndexRequest createIndexRequest = new CreateIndexRequest.Builder().index(index)
                        .build();
                client.indices().create(createIndexRequest);
            }
            LOGGER.info("1. Creating an index {} finished", index);
    
            // 2. Indexing data
            LOGGER.info("2. Indexing data start");
            IndexData indexData = new IndexData("first_name", "Bruce");
            IndexRequest<IndexData> indexRequest = new IndexRequest.Builder<IndexData>().index(index).id("1").document(indexData).build();
            client.index(indexRequest);
            LOGGER.info("2. Indexing data finished");
    
            Thread.sleep(1500);
    
            // 3. Searching for documents
            LOGGER.info("3. Searching for documents start");
            SearchResponse<IndexData> searchResponse = client.search(s -> s.index(index), IndexData.class);
            for (int i = 0; i< searchResponse.hits().hits().size(); i++) {
                LOGGER.info(searchResponse.hits().hits().get(i).source());
            }        
            LOGGER.info("3. Searching for documents finished");
            
            // 4. Deleting a document
            LOGGER.info("4. Deleting a document start");
            client.delete(b -> b.index(index).id("1"));
            LOGGER.info("4. Deleting a document finished");
    
            // 5. Deleting an index
            LOGGER.info("5. Deleting an index {} start", index);
            DeleteIndexRequest deleteIndexRequest = new DeleteIndexRequest.Builder().index(index).build();
            DeleteIndexResponse deleteIndexResponse = client.indices().delete(deleteIndexRequest);
            LOGGER.info("5. Deleting an index {} finished", index);

            LOGGER.info("end test...");
        } catch (Exception e) {
            LOGGER.error(e.toString());
        } finally {
            if (transport != null) {
                try {
                    transport.close();
                } catch (Exception e) {
                    LOGGER.error("Failed to close transport", e);
                }
            }
        }
    }
}

PolarSearch 3.x版本

依赖配置(pom.xml)

<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
         xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
    <modelVersion>4.0.0</modelVersion>

    <groupId>org.example</groupId>
    <artifactId>OpenSearchJavaClientSample</artifactId>
    <version>1.0-SNAPSHOT</version>

    <properties>
        <maven.compiler.source>21</maven.compiler.source>
        <maven.compiler.target>21</maven.compiler.target>
        <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
    </properties>
    <dependencies>
        <dependency>
            <groupId>org.opensearch.client</groupId>
            <artifactId>opensearch-java</artifactId>
            <version>3.3.0</version>
        </dependency>
        <dependency>
            <groupId>org.apache.httpcomponents.client5</groupId>
            <artifactId>httpclient5</artifactId>
            <version>5.4.3</version>
        </dependency>
        <dependency>
            <groupId>org.apache.logging.log4j</groupId>
            <artifactId>log4j-core</artifactId>
            <version>2.21.0</version>
        </dependency>
    </dependencies>
</project>

示例程序(OpenSearchClientExample.java)

package samples;

import java.security.KeyManagementException;
import java.security.KeyStoreException;
import java.security.NoSuchAlgorithmException;

import javax.net.ssl.SSLContext;
import javax.net.ssl.SSLEngine;

import org.apache.hc.client5.http.auth.AuthScope;
import org.apache.hc.client5.http.auth.UsernamePasswordCredentials;
import org.apache.hc.client5.http.impl.auth.BasicCredentialsProvider;
import org.apache.hc.client5.http.impl.nio.PoolingAsyncClientConnectionManager;
import org.apache.hc.client5.http.impl.nio.PoolingAsyncClientConnectionManagerBuilder;
import org.apache.hc.client5.http.ssl.ClientTlsStrategyBuilder;
import org.apache.hc.core5.function.Factory;
import org.apache.hc.core5.http.HttpHost;
import org.apache.hc.core5.http.nio.ssl.TlsStrategy;
import org.apache.hc.core5.reactor.ssl.TlsDetails;
import org.apache.hc.core5.ssl.SSLContextBuilder;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import org.opensearch.client.opensearch.OpenSearchClient;
import org.opensearch.client.opensearch.core.IndexRequest;
import org.opensearch.client.opensearch.core.SearchResponse;
import org.opensearch.client.opensearch.indices.CreateIndexRequest;
import org.opensearch.client.opensearch.indices.DeleteIndexRequest;
import org.opensearch.client.opensearch.indices.DeleteIndexResponse;
import org.opensearch.client.transport.OpenSearchTransport;
import org.opensearch.client.transport.httpclient5.ApacheHttpClient5TransportBuilder;


public class OpenSearchClientExample {
    static class IndexData {
        private String title;
        private String text;
    
        public IndexData() {}
    
        public IndexData(String title, String text) {
            this.title = title;
            this.text = text;
        }
    
        public String getTitle() {
            return title;
        }
    
        public void setTitle(String title) {
            this.title = title;
        }
    
        public String getText() {
            return text;
        }
    
        public void setText(String text) {
            this.text = text;
        }
    
        @Override
        public String toString() {
            return String.format("IndexData{title='%s', text='%s'}", title, text);
        }
    }
    
    private static final Logger LOGGER = LogManager.getLogger(OpenSearchClientExample.class);

    public static OpenSearchTransport createTransport() throws NoSuchAlgorithmException, KeyStoreException, KeyManagementException {
        var env = System.getenv();
        var hostname = env.getOrDefault("HOST", "<polarsearch_host>");
        var port = Integer.parseInt(env.getOrDefault("PORT", "<polarsearch_port>"));
        var scheme = env.getOrDefault("SCHEME", "http");
        var user = env.getOrDefault("USERNAME", "<polarsearch_username>");
        var pass = env.getOrDefault("PASSWORD", "<polarsearch_password>");

        final HttpHost host = new HttpHost(scheme, hostname, port);
        final BasicCredentialsProvider credentialsProvider = new BasicCredentialsProvider();
        credentialsProvider.setCredentials(new AuthScope(host), new UsernamePasswordCredentials(user, pass.toCharArray()));

        final SSLContext sslContext = SSLContextBuilder
                .create()
                .loadTrustMaterial(null, (chains, authType) -> true)
                .build();

        final ApacheHttpClient5TransportBuilder builder = ApacheHttpClient5TransportBuilder.builder(host);
        builder.setHttpClientConfigCallback(httpClientBuilder -> {
            final TlsStrategy tlsStrategy = ClientTlsStrategyBuilder.create()
                    .setSslContext(sslContext)
                    .setTlsDetailsFactory(new Factory<SSLEngine, TlsDetails>() {
                        @Override
                        public TlsDetails create(final SSLEngine sslEngine) {
                            return new TlsDetails(sslEngine.getSession(), sslEngine.getApplicationProtocol());
                        }
                    })
                    .build();

            final PoolingAsyncClientConnectionManager connectionManager = PoolingAsyncClientConnectionManagerBuilder
                    .create()
                    .setTlsStrategy(tlsStrategy)
                    .build();

            return httpClientBuilder
                    .setDefaultCredentialsProvider(credentialsProvider)
                    .setConnectionManager(connectionManager);
        });

        final OpenSearchTransport transport = builder.build();
        return transport;
    }

    public static void main(String[] args) {
        OpenSearchTransport transport = null;
        try {
            LOGGER.info("start test...");

            transport = createTransport();
            var client = new OpenSearchClient(transport);
            LOGGER.info("client create");
            final var index = "my-index";
            // 1. Creating an index
            LOGGER.info("1. Creating an index {} start", index);
            if (!client.indices().exists(r -> r.index(index)).value()) {
                CreateIndexRequest createIndexRequest = new CreateIndexRequest.Builder().index(index)
                        .build();
                client.indices().create(createIndexRequest);
            }
            LOGGER.info("1. Creating an index {} finished", index);
    
            // 2. Indexing data
            LOGGER.info("2. Indexing data start");
            IndexData indexData = new IndexData("first_name", "Bruce");
            IndexRequest<IndexData> indexRequest = new IndexRequest.Builder<IndexData>().index(index).id("1").document(indexData).build();
            client.index(indexRequest);
            LOGGER.info("2. Indexing data finished");
    
            Thread.sleep(1500);
    
            // 3. Searching for documents
            LOGGER.info("3. Searching for documents start");
            SearchResponse<IndexData> searchResponse = client.search(s -> s.index(index), IndexData.class);
            for (int i = 0; i< searchResponse.hits().hits().size(); i++) {
                LOGGER.info(searchResponse.hits().hits().get(i).source());
            }        
            LOGGER.info("3. Searching for documents finished");
            
            // 4. Deleting a document
            LOGGER.info("4. Deleting a document start");
            client.delete(b -> b.index(index).id("1"));
            LOGGER.info("4. Deleting a document finished");
    
            // 5. Deleting an index
            LOGGER.info("5. Deleting an index {} start", index);
            DeleteIndexRequest deleteIndexRequest = new DeleteIndexRequest.Builder().index(index).build();
            DeleteIndexResponse deleteIndexResponse = client.indices().delete(deleteIndexRequest);
            LOGGER.info("5. Deleting an index {} finished", index);

            LOGGER.info("end test...");
        } catch (Exception e) {
            LOGGER.error(e.toString());
        } finally {
            if (transport != null) {
                try {
                    transport.close();
                } catch (Exception e) {
                    LOGGER.error("Failed to close transport", e);
                }
            }
        }
    }
}

运行方式

mvn clean compile exec:java -Dexec.mainClass=samples.OpenSearchClientExample

OpenSearch Java high-level REST client

示例代码

PolarSearch 1.x版本

依赖配置(pom.xml)

<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
         xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
    <modelVersion>4.0.0</modelVersion>

    <groupId>org.example</groupId>
    <artifactId>OpenSearchJavaClientSample</artifactId>
    <version>1.0-SNAPSHOT</version>

    <properties>
        <maven.compiler.source>21</maven.compiler.source>
        <maven.compiler.target>21</maven.compiler.target>
        <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
    </properties>
    <dependencies>
        <dependency>
            <groupId>org.opensearch.client</groupId>
            <artifactId>opensearch-rest-high-level-client</artifactId>
            <version>1.3.20</version>
        </dependency>
        <dependency>
            <groupId>org.apache.logging.log4j</groupId>
            <artifactId>log4j-core</artifactId>
            <version>2.21.0</version>
        </dependency>
    </dependencies>
</project>

示例程序(OpenSearchClientExample.java)

package samples;

import java.util.HashMap;
import java.util.Map;

import javax.net.ssl.SSLContext;

import org.apache.http.HttpHost;
import org.apache.http.auth.AuthScope;
import org.apache.http.auth.UsernamePasswordCredentials;
import org.apache.http.impl.client.BasicCredentialsProvider;
import org.apache.http.ssl.SSLContextBuilder;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import org.opensearch.action.delete.DeleteRequest;
import org.opensearch.action.index.IndexRequest;
import org.opensearch.action.admin.indices.delete.DeleteIndexRequest;
import org.opensearch.action.search.SearchRequest;
import org.opensearch.action.search.SearchResponse;
import org.opensearch.action.support.master.AcknowledgedResponse;
import org.opensearch.client.RequestOptions;
import org.opensearch.client.RestClient;
import org.opensearch.client.RestClientBuilder;
import org.opensearch.client.RestHighLevelClient;
import org.opensearch.client.indices.CreateIndexRequest;
import org.opensearch.client.indices.GetIndexRequest;
import org.opensearch.index.query.QueryBuilders;
import org.opensearch.search.SearchHit;
import org.opensearch.search.builder.SearchSourceBuilder;


public class OpenSearchClientExample {

    static class IndexData {
        private String title;
        private String text;

        public IndexData() {}

        public IndexData(String title, String text) {
            this.title = title;
            this.text = text;
        }

        public String getTitle() {
            return title;
        }

        public void setTitle(String title) {
            this.title = title;
        }

        public String getText() {
            return text;
        }

        public void setText(String text) {
            this.text = text;
        }

        @Override
        public String toString() {
            return String.format("IndexData{title='%s', text='%s'}", title, text);
        }
    }

    private static final Logger LOGGER = LogManager.getLogger(OpenSearchClientExample.class);

    public static RestHighLevelClient createClient() throws Exception {
        var env = System.getenv();
        var hostname = env.getOrDefault("HOST", "<polarsearch_host>");
        var port = Integer.parseInt(env.getOrDefault("PORT", "<polarsearch_port>"));
        var scheme = env.getOrDefault("SCHEME", "http");
        var user = env.getOrDefault("USERNAME", "<polarsearch_username>");
        var pass = env.getOrDefault("PASSWORD", "<polarsearch_password>");

        final HttpHost host = new HttpHost(hostname, port, scheme);
        final BasicCredentialsProvider credentialsProvider = new BasicCredentialsProvider();
        credentialsProvider.setCredentials(AuthScope.ANY,
                new UsernamePasswordCredentials(user, pass));

        final SSLContext sslContext = SSLContextBuilder.create()
                .loadTrustMaterial(null, (chains, authType) -> true)
                .build();

        RestClientBuilder builder = RestClient.builder(host)
                .setHttpClientConfigCallback(httpClientBuilder ->
                        httpClientBuilder
                                .setDefaultCredentialsProvider(credentialsProvider)
                                .setSSLContext(sslContext));

        return new RestHighLevelClient(builder);
    }

    public static void main(String[] args) {
        try (RestHighLevelClient client = createClient()) {
            LOGGER.info("start test...");
            LOGGER.info("client create");
            final String index = "my-index";

            // 1. Creating an index
            LOGGER.info("1. Creating an index {} start", index);
            boolean exists = client.indices().exists(new GetIndexRequest(index), RequestOptions.DEFAULT);
            if (!exists) {
                CreateIndexRequest createIndexRequest = new CreateIndexRequest(index);
                client.indices().create(createIndexRequest, RequestOptions.DEFAULT);
            }
            LOGGER.info("1. Creating an index {} finished", index);

            // 2. Indexing data
            LOGGER.info("2. Indexing data start");
            IndexData indexData = new IndexData("first_name", "Bruce");
            Map<String, Object> document = new HashMap<>();
            document.put("title", indexData.getTitle());
            document.put("text", indexData.getText());
            IndexRequest indexRequest = new IndexRequest(index).id("1").source(document);
            client.index(indexRequest, RequestOptions.DEFAULT);
            LOGGER.info("2. Indexing data finished");

            Thread.sleep(1500);

            // 3. Searching for documents
            LOGGER.info("3. Searching for documents start");
            SearchRequest searchRequest = new SearchRequest(index);
            searchRequest.source(new SearchSourceBuilder().query(QueryBuilders.matchAllQuery()));
            SearchResponse searchResponse = client.search(searchRequest, RequestOptions.DEFAULT);
            for (SearchHit hit : searchResponse.getHits().getHits()) {
                Map<String, Object> sourceMap = hit.getSourceAsMap();
                IndexData data = new IndexData((String) sourceMap.get("title"), (String) sourceMap.get("text"));
                LOGGER.info(data);
            }
            LOGGER.info("3. Searching for documents finished");

            // 4. Deleting a document
            LOGGER.info("4. Deleting a document start");
            client.delete(new DeleteRequest(index, "1"), RequestOptions.DEFAULT);
            LOGGER.info("4. Deleting a document finished");

            // 5. Deleting an index
            LOGGER.info("5. Deleting an index {} start", index);
            DeleteIndexRequest deleteIndexRequest = new DeleteIndexRequest(index);
            AcknowledgedResponse deleteIndexResponse = client.indices().delete(deleteIndexRequest, RequestOptions.DEFAULT);
            LOGGER.info("5. Deleting an index {} finished", index);

            LOGGER.info("end test...");
        } catch (Exception e) {
            LOGGER.error(e.toString(), e);
        }
    }
}

PolarSearch 3.x版本

依赖配置(pom.xml)

<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
         xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
    <modelVersion>4.0.0</modelVersion>

    <groupId>org.example</groupId>
    <artifactId>OpenSearchJavaClientSample</artifactId>
    <version>1.0-SNAPSHOT</version>

    <properties>
        <maven.compiler.source>21</maven.compiler.source>
        <maven.compiler.target>21</maven.compiler.target>
        <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
    </properties>
    <dependencies>
        <dependency>
            <groupId>org.opensearch.client</groupId>
            <artifactId>opensearch-rest-high-level-client</artifactId>
            <version>3.3.2</version>
        </dependency>
        <dependency>
            <groupId>org.apache.logging.log4j</groupId>
            <artifactId>log4j-core</artifactId>
            <version>2.21.0</version>
        </dependency>
    </dependencies>
</project>

示例程序(OpenSearchClientExample.java)

package samples;

import java.util.HashMap;
import java.util.Map;

import javax.net.ssl.SSLContext;

import org.apache.hc.client5.http.auth.AuthScope;
import org.apache.hc.client5.http.auth.UsernamePasswordCredentials;
import org.apache.hc.client5.http.impl.auth.BasicCredentialsProvider;
import org.apache.hc.client5.http.impl.nio.PoolingAsyncClientConnectionManager;
import org.apache.hc.client5.http.impl.nio.PoolingAsyncClientConnectionManagerBuilder;
import org.apache.hc.client5.http.ssl.ClientTlsStrategyBuilder;
import org.apache.hc.core5.http.HttpHost;
import org.apache.hc.core5.http.nio.ssl.TlsStrategy;
import org.apache.hc.core5.ssl.SSLContextBuilder;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import org.opensearch.action.delete.DeleteRequest;
import org.opensearch.action.index.IndexRequest;
import org.opensearch.action.search.SearchRequest;
import org.opensearch.action.search.SearchResponse;
import org.opensearch.action.support.clustermanager.AcknowledgedResponse;
import org.opensearch.client.RequestOptions;
import org.opensearch.client.RestClient;
import org.opensearch.client.RestClientBuilder;
import org.opensearch.client.RestHighLevelClient;
import org.opensearch.client.indices.CreateIndexRequest;
import org.opensearch.client.indices.GetIndexRequest;
import org.opensearch.index.query.QueryBuilders;
import org.opensearch.search.SearchHit;
import org.opensearch.search.builder.SearchSourceBuilder;
import org.opensearch.action.admin.indices.delete.DeleteIndexRequest;


public class OpenSearchClientExample {

    static class IndexData {
        private String title;
        private String text;

        public IndexData() {}

        public IndexData(String title, String text) {
            this.title = title;
            this.text = text;
        }

        public String getTitle() {
            return title;
        }

        public void setTitle(String title) {
            this.title = title;
        }

        public String getText() {
            return text;
        }

        public void setText(String text) {
            this.text = text;
        }

        @Override
        public String toString() {
            return String.format("IndexData{title='%s', text='%s'}", title, text);
        }
    }

    private static final Logger LOGGER = LogManager.getLogger(OpenSearchClientExample.class);

    public static RestHighLevelClient createClient() throws Exception {
        var env = System.getenv();
        var hostname = env.getOrDefault("HOST", "<polarsearch_host>");
        var port = Integer.parseInt(env.getOrDefault("PORT", "<polarsearch_port>"));
        var scheme = env.getOrDefault("SCHEME", "http");
        var user = env.getOrDefault("USERNAME", "<polarsearch_username>");
        var pass = env.getOrDefault("PASSWORD", "<polarsearch_password>");

        final HttpHost host = new HttpHost(scheme, hostname, port);
        final BasicCredentialsProvider credentialsProvider = new BasicCredentialsProvider();
        credentialsProvider.setCredentials(new AuthScope(host),
                new UsernamePasswordCredentials(user, pass.toCharArray()));

        final SSLContext sslContext = SSLContextBuilder.create()
                .loadTrustMaterial(null, (chains, authType) -> true)
                .build();

        final TlsStrategy tlsStrategy = ClientTlsStrategyBuilder.create()
                .setSslContext(sslContext)
                .build();

        final PoolingAsyncClientConnectionManager connectionManager = PoolingAsyncClientConnectionManagerBuilder
                .create()
                .setTlsStrategy(tlsStrategy)
                .build();

        RestClientBuilder builder = RestClient.builder(host)
                .setHttpClientConfigCallback(httpClientBuilder ->
                        httpClientBuilder
                                .setDefaultCredentialsProvider(credentialsProvider)
                                .setConnectionManager(connectionManager));

        return new RestHighLevelClient(builder);
    }

    public static void main(String[] args) {
        try (RestHighLevelClient client = createClient()) {
            LOGGER.info("start test...");
            LOGGER.info("client create");
            final String index = "my-index";

            // 1. Creating an index
            LOGGER.info("1. Creating an index {} start", index);
            boolean exists = client.indices().exists(new GetIndexRequest(index), RequestOptions.DEFAULT);
            if (!exists) {
                CreateIndexRequest createIndexRequest = new CreateIndexRequest(index);
                client.indices().create(createIndexRequest, RequestOptions.DEFAULT);
            }
            LOGGER.info("1. Creating an index {} finished", index);

            // 2. Indexing data
            LOGGER.info("2. Indexing data start");
            IndexData indexData = new IndexData("first_name", "Bruce");
            Map<String, Object> document = new HashMap<>();
            document.put("title", indexData.getTitle());
            document.put("text", indexData.getText());
            IndexRequest indexRequest = new IndexRequest(index).id("1").source(document);
            client.index(indexRequest, RequestOptions.DEFAULT);
            LOGGER.info("2. Indexing data finished");

            Thread.sleep(1500);

            // 3. Searching for documents
            LOGGER.info("3. Searching for documents start");
            SearchRequest searchRequest = new SearchRequest(index);
            searchRequest.source(new SearchSourceBuilder().query(QueryBuilders.matchAllQuery()));
            SearchResponse searchResponse = client.search(searchRequest, RequestOptions.DEFAULT);
            for (SearchHit hit : searchResponse.getHits().getHits()) {
                Map<String, Object> sourceMap = hit.getSourceAsMap();
                IndexData data = new IndexData((String) sourceMap.get("title"), (String) sourceMap.get("text"));
                LOGGER.info(data);
            }
            LOGGER.info("3. Searching for documents finished");

            // 4. Deleting a document
            LOGGER.info("4. Deleting a document start");
            client.delete(new DeleteRequest(index, "1"), RequestOptions.DEFAULT);
            LOGGER.info("4. Deleting a document finished");

            // 5. Deleting an index
            LOGGER.info("5. Deleting an index {} start", index);
            DeleteIndexRequest deleteIndexRequest = new DeleteIndexRequest(index);
            AcknowledgedResponse deleteIndexResponse = client.indices().delete(deleteIndexRequest, RequestOptions.DEFAULT);
            LOGGER.info("5. Deleting an index {} finished", index);

            LOGGER.info("end test...");
        } catch (Exception e) {
            LOGGER.error(e.toString(), e);
        }
    }
}

运行方式

mvn clean compile exec:java -Dexec.mainClass=samples.OpenSearchClientExample

ElasticSearch Java High Level REST Client

示例程序

依赖配置(pom.xml)

<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
         xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
    <modelVersion>4.0.0</modelVersion>

    <groupId>org.example</groupId>
    <artifactId>OpenSearchJavaSample</artifactId>
    <version>1.0-SNAPSHOT</version>

    <properties>
        <maven.compiler.source>17</maven.compiler.source>
        <maven.compiler.target>17</maven.compiler.target>
        <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
    </properties>
    <dependencies>
        <dependency>
            <groupId>org.elasticsearch.client</groupId>
            <artifactId>elasticsearch-rest-high-level-client</artifactId>
            <version>7.13.4</version>
        </dependency>
        <dependency>
            <groupId>org.apache.logging.log4j</groupId>
            <artifactId>log4j-core</artifactId>
            <version>2.20.0</version>
        </dependency>
    </dependencies>
</project>

示例程序(OpenSearchClientExample.java)

package samples;

import java.util.HashMap;
import java.util.Map;

import javax.net.ssl.SSLContext;

import org.apache.http.HttpHost;
import org.apache.http.auth.AuthScope;
import org.apache.http.auth.UsernamePasswordCredentials;
import org.apache.http.impl.client.BasicCredentialsProvider;
import org.apache.http.ssl.SSLContextBuilder;
import org.apache.http.util.EntityUtils;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import org.elasticsearch.action.delete.DeleteRequest;
import org.elasticsearch.action.get.GetRequest;
import org.elasticsearch.action.get.GetResponse;
import org.elasticsearch.action.index.IndexRequest;
import org.elasticsearch.action.search.SearchRequest;
import org.elasticsearch.action.search.SearchResponse;
import org.elasticsearch.action.support.master.AcknowledgedResponse;
import org.elasticsearch.client.Request;
import org.elasticsearch.client.RequestOptions;
import org.elasticsearch.client.Response;
import org.elasticsearch.client.RestClient;
import org.elasticsearch.client.RestClientBuilder;
import org.elasticsearch.client.RestHighLevelClient;
import org.elasticsearch.client.indices.CreateIndexRequest;
import org.elasticsearch.client.indices.GetIndexRequest;
import org.elasticsearch.action.admin.indices.delete.DeleteIndexRequest;
import org.elasticsearch.index.query.QueryBuilders;
import org.elasticsearch.search.SearchHit;
import org.elasticsearch.search.builder.SearchSourceBuilder;


public class ElasticSearchClientExample {

    static class IndexData {
        private String title;
        private String text;

        public IndexData() {}

        public IndexData(String title, String text) {
            this.title = title;
            this.text = text;
        }

        public String getTitle() {
            return title;
        }

        public void setTitle(String title) {
            this.title = title;
        }

        public String getText() {
            return text;
        }

        public void setText(String text) {
            this.text = text;
        }

        @Override
        public String toString() {
            return String.format("IndexData{title='%s', text='%s'}", title, text);
        }
    }

    private static final Logger LOGGER = LogManager.getLogger(ElasticSearchClientExample.class);

    public static RestHighLevelClient createClient() throws Exception {
        var env = System.getenv();
        var hostname = env.getOrDefault("HOST", "<polarsearch_host>");
        var port = Integer.parseInt(env.getOrDefault("PORT", "<polarsearch_port>"));
        var scheme = env.getOrDefault("SCHEME", "http");
        var user = env.getOrDefault("USERNAME", "<polarsearch_username>");
        var pass = env.getOrDefault("PASSWORD", "<polarsearch_password>");

        final HttpHost host = new HttpHost(hostname, port, scheme);
        final BasicCredentialsProvider credentialsProvider = new BasicCredentialsProvider();
        credentialsProvider.setCredentials(AuthScope.ANY,
                new UsernamePasswordCredentials(user, pass));

        final SSLContext sslContext = SSLContextBuilder.create()
                .loadTrustMaterial(null, (chains, authType) -> true)
                .build();

        RestClientBuilder builder = RestClient.builder(host)
                .setHttpClientConfigCallback(httpClientBuilder ->
                        httpClientBuilder
                                .setDefaultCredentialsProvider(credentialsProvider)
                                .setSSLContext(sslContext));

        return new RestHighLevelClient(builder);
    }

    public static void main(String[] args) {
        try (RestHighLevelClient client = createClient()) {
            LOGGER.info("start test...");
            LOGGER.info("client create");
            final String index = "my-index";

            // 0. Ping cluster
            LOGGER.info("0. Ping cluster start");
            Response response = client.getLowLevelClient().performRequest(new Request("HEAD", "/"));
            int statusCode = response.getStatusLine().getStatusCode();
            LOGGER.info("Ping status code: {}", statusCode);
            if (statusCode == 200) {
                Response infoResponse = client.getLowLevelClient().performRequest(new Request("GET", "/"));
                String responseBody = EntityUtils.toString(infoResponse.getEntity());
                LOGGER.info("Cluster info: {}", responseBody);
            }
            LOGGER.info("0. Ping cluster finished");

            // 1. Creating an index
            LOGGER.info("1. Creating an index {} start", index);
            boolean exists = client.indices().exists(new GetIndexRequest(index), RequestOptions.DEFAULT);
            if (!exists) {
                CreateIndexRequest createIndexRequest = new CreateIndexRequest(index);
                client.indices().create(createIndexRequest, RequestOptions.DEFAULT);
            }
            LOGGER.info("1. Creating an index {} finished", index);

            // 2. Indexing data
            LOGGER.info("2. Indexing data start");
            IndexData indexData1 = new IndexData("first_name", "Bruce");
            Map<String, Object> document1 = new HashMap<>();
            document1.put("title", indexData1.getTitle());
            document1.put("text", indexData1.getText());
            client.index(new IndexRequest(index).id("1").source(document1), RequestOptions.DEFAULT);

            IndexData indexData2 = new IndexData("last_name", "Wayne");
            Map<String, Object> document2 = new HashMap<>();
            document2.put("title", indexData2.getTitle());
            document2.put("text", indexData2.getText());
            client.index(new IndexRequest(index).id("2").source(document2), RequestOptions.DEFAULT);
            LOGGER.info("2. Indexing data finished");

            Thread.sleep(1500);

            // 3. Getting a document by ID
            LOGGER.info("3. Getting a document by ID start");
            GetResponse getResponse = client.get(new GetRequest(index, "1"), RequestOptions.DEFAULT);
            if (getResponse.isExists()) {
                Map<String, Object> sourceMap = getResponse.getSourceAsMap();
                IndexData data = new IndexData((String) sourceMap.get("title"), (String) sourceMap.get("text"));
                LOGGER.info(data);
            }
            LOGGER.info("3. Getting a document by ID finished");

            // 4. Searching for documents
            LOGGER.info("4. Searching for documents start");
            SearchRequest searchRequest = new SearchRequest(index);
            searchRequest.source(new SearchSourceBuilder().query(QueryBuilders.matchAllQuery()));
            SearchResponse searchResponse = client.search(searchRequest, RequestOptions.DEFAULT);
            for (SearchHit hit : searchResponse.getHits().getHits()) {
                Map<String, Object> sourceMap = hit.getSourceAsMap();
                IndexData data = new IndexData((String) sourceMap.get("title"), (String) sourceMap.get("text"));
                LOGGER.info(data);
            }
            LOGGER.info("4. Searching for documents finished");

            // 5. Deleting a document
            LOGGER.info("5. Deleting a document start");
            client.delete(new DeleteRequest(index, "1"), RequestOptions.DEFAULT);
            LOGGER.info("5. Deleting a document finished");

            // 6. Deleting an index
            LOGGER.info("6. Deleting an index {} start", index);
            DeleteIndexRequest deleteIndexRequest = new DeleteIndexRequest(index);
            AcknowledgedResponse deleteIndexResponse = client.indices().delete(deleteIndexRequest, RequestOptions.DEFAULT);
            LOGGER.info("6. Deleting an index {} finished", index);

            LOGGER.info("end test...");
        } catch (Exception e) {
            LOGGER.error(e.toString(), e);
        }
    }
}

Low-level Python client

依赖配置

pip3 install opensearch-py

示例程序

import os
from opensearchpy import OpenSearch

host = os.getenv("HOST", default="<polarsearch_host>")
port = int(os.getenv("PORT", <polarsearch_port>))
auth = (os.getenv("USERNAME", "<polarsearch_username>"), os.getenv("PASSWORD", "<polarsearch_password>"))


client = OpenSearch(
    hosts=[{"host": host, "port": port}],
    http_auth=auth,
    use_ssl=False,
    verify_certs=False,
    ssl_show_warn=False,
)

# Create an index with non-default settings.
index_name = 'python-test-index'
index_body = {
  'settings': {
    'index': {
      'number_of_shards': 4
    }
  }
}

response = client.indices.create(index=index_name, body=index_body)
print('\nCreating index:')
print(response)

# Add a document to the index.
document = {
  'title': 'Moneyball',
  'director': 'Bennett Miller',
  'year': '2011'
}
id = '1'

response = client.index(
    index = index_name,
    body = document,
    id = id,
    refresh = True
)

print('\nAdding document:')
print(response)

# Perform bulk operations

movies = '{ "index" : { "_index" : "my-dsl-index", "_id" : "2" } } \n { "title" : "Interstellar", "director" : "Christopher Nolan", "year" : "2014"} \n { "create" : { "_index" : "my-dsl-index", "_id" : "3" } } \n { "title" : "Star Trek Beyond", "director" : "Justin Lin", "year" : "2015"} \n { "update" : {"_id" : "3", "_index" : "my-dsl-index" } } \n { "doc" : {"year" : "2016"} }'

client.bulk(body=movies)

# Search for the document.
q = 'miller'
query = {
  'size': 5,
  'query': {
    'multi_match': {
      'query': q,
      'fields': ['title^2', 'director']
    }
  }
}

response = client.search(
    body = query,
    index = index_name
)
print('\nSearch results:')
print(response)

# Delete the document.
response = client.delete(
    index = index_name,
    id = id
)

print('\nDeleting document:')
print(response)

# Delete the index.
response = client.indices.delete(
    index = index_name
)

print('\nDeleting index:')
print(response)

相关文档