O Java High Level REST Client é um cliente Elasticsearch que oferece uma API de alto nível para gerenciar índices e documentos. O LindormSearch é compatível com o Elasticsearch 7.10 e versões anteriores. Portanto, use este cliente para se conectar ao LindormSearch e executar consultas e pesquisas complexas sem alterar o código da aplicação existente.
O Java High Level REST Client tem compatibilidade futura. Por exemplo, a versão 6.7.0 comunica-se com clusters do Elasticsearch 6.7.0 ou posteriores. Use a versão 7.10.0 ou anterior para se conectar ao LindormSearch.
Pré-requisitos
Antes de começar, verifique se você possui:
JDK 1.8 ou posterior instalado
Mecanismo de busca do LindormSearch ativado. Consulte o Guia de ativação.
Endereço IP do cliente adicionado à lista de permissões da instância Lindorm. Consulte Configure listas de permissões.
Adicionar dependências
Em projetos Maven, adicione as seguintes dependências ao arquivo pom.xml:
<dependency>
<groupId>org.elasticsearch.client</groupId>
<artifactId>elasticsearch-rest-high-level-client</artifactId>
<version>7.10.0</version>
</dependency>
<dependency>
<groupId>org.apache.logging.log4j</groupId>
<artifactId>log4j-core</artifactId>
<version>2.20.0</version>
</dependency>
<dependency>
<groupId>org.apache.logging.log4j</groupId>
<artifactId>log4j-api</artifactId>
<version>2.20.0</version>
</dependency>
Conectar-se ao LindormSearch
Use RestClient.builder() para criar um objeto RestHighLevelClient. O cliente usa BasicCredentialsProvider para autenticação.
// Set the Elasticsearch-compatible endpoint and port for LindormSearch.
String search_url = "ld-t4n5668xk31ui****-proxy-search-public.lindorm.rds.aliyuncs.com";
int search_port = 30070;
// Set the username and password. Retrieve them from the Lindorm console:
// navigate to Database Connections > Search Engine tab.
String username = "user";
String password = "test";
final CredentialsProvider credentials_provider = new BasicCredentialsProvider();
credentials_provider.setCredentials(AuthScope.ANY, new UsernamePasswordCredentials(username, password));
RestHighLevelClient highClient = new RestHighLevelClient(
RestClient.builder(new HttpHost(search_url, search_port, "http"))
.setHttpClientConfigCallback(new RestClientBuilder.HttpClientConfigCallback() {
public HttpAsyncClientBuilder customizeHttpClient(HttpAsyncClientBuilder httpClientBuilder) {
return httpClientBuilder.setDefaultCredentialsProvider(credentials_provider);
}
})
);
Parâmetros de conexão
|
Parâmetro |
Descrição |
|
|
Endpoint compatível com Elasticsearch para o mecanismo de busca do LindormSearch. Para obter o endpoint, consulte Endereço compatível com Elasticsearch. Use o endereço da Virtual Private Cloud (VPC) ao se conectar a partir de uma VPC ou o endereço de Internet para conexões via rede pública. |
|
|
Porta para compatibilidade com Elasticsearch no LindormSearch: |
|
|
Nome de usuário do mecanismo de busca. |
|
|
Senha do mecanismo de busca. |
Escolha o tipo de conexão de rede
VPC (recomendado): Se a aplicação for executada em uma instância do Elastic Compute Service (ECS) na mesma VPC da instância Lindorm, conecte-se via VPC para garantir menor latência e maior segurança. Defina
search_urlcomo o endereço VPC do endpoint compatível com Elasticsearch.Rede pública: Se a aplicação estiver fora da Alibaba Cloud, ative primeiro o endpoint público. No console do Lindorm, acesse Database Connections, clique em na aba Search Engine e selecione Enable Public Endpoint no canto superior direito. Em seguida, defina
search_urlcomo o endereço de Internet do endpoint compatível com Elasticsearch.
Crie um índice
Use CreateIndexRequest para criar um índice. O exemplo abaixo cria um índice chamado lindorm_index com 4 shards.
String index_name = "lindorm_index";
// Create a CreateIndexRequest and configure index settings.
CreateIndexRequest createIndexRequest = new CreateIndexRequest(index_name);
Map<String, Object> settingsMap = new HashMap<>();
settingsMap.put("index.number_of_shards", 4);
createIndexRequest.settings(settingsMap);
CreateIndexResponse createIndexResponse = highClient.indices().create(createIndexRequest, COMMON_OPTIONS);
if (createIndexResponse.isAcknowledged()) {
System.out.println("Create index [" + index_name + "] successfully.");
}
Indexar um documento
Use IndexRequest para gravar um único documento. Especifique um ID de documento ou deixe-o em branco para que o sistema gere um automaticamente. Omitir o ID pode melhorar o desempenho de gravação.
// Specify the document ID.
String doc_id = "test";
// Build the document fields. Replace with actual field names and values for your use case.
Map<String, Object> jsonMap = new HashMap<>();
jsonMap.put("field1", "value1");
jsonMap.put("field2", "value2");
IndexRequest indexRequest = new IndexRequest(index_name);
indexRequest.id(doc_id).source(jsonMap);
IndexResponse indexResponse = highClient.index(indexRequest, COMMON_OPTIONS);
System.out.println("Index document with id[" + indexResponse.getId() + "] successfully.");
Indexar documentos em massa
Combine BulkProcessor com bulkAsync() para gravar grandes volumes de documentos de forma eficiente. O BulkProcessor agrupa solicitações individuais de indexação e as envia em lote sempre que qualquer limiar configurado for atingido.
Configuração do BulkProcessor
|
Parâmetro |
Descrição |
Valor de exemplo |
|
|
Número máximo de solicitações em massa simultâneas. Aumente este valor para melhorar o throughput sob cargas intensas de gravação. Padrão: |
|
|
|
Intervalo de tempo após o qual uma solicitação em massa é enviada, independentemente do tamanho ou da quantidade. Funciona como medida de segurança para evitar que os dados fiquem em buffer por muito tempo. |
|
|
|
Quantidade de operações individuais que aciona um flush. Ajuste este valor com base no tamanho médio dos documentos. |
|
|
|
Tamanho total das operações em buffer que dispara um flush. |
|
int bulkTotal = 100000;
AtomicLong failedBulkItemCount = new AtomicLong();
BulkProcessor.Builder builder = BulkProcessor.builder(
(request, bulkListener) -> highClient.bulkAsync(request, COMMON_OPTIONS, bulkListener),
new BulkProcessor.Listener() {
@Override
public void beforeBulk(long executionId, BulkRequest request) {}
@Override
public void afterBulk(long executionId, BulkRequest request, BulkResponse response) {
// Count failed items in the bulk response.
for (BulkItemResponse bulkItemResponse : response) {
if (bulkItemResponse.isFailed()) {
failedBulkItemCount.incrementAndGet();
}
}
}
@Override
public void afterBulk(long executionId, BulkRequest request, Throwable failure) {
// If this callback fires, all requests in the bulk were not executed.
if (null != failure) {
failedBulkItemCount.addAndGet(request.numberOfActions());
}
}
});
// Maximum concurrent bulk requests. Default is 1; increase for higher throughput.
builder.setConcurrentRequests(10);
// Flush thresholds — a bulk request is sent when any of these is met.
builder.setFlushInterval(TimeValue.timeValueSeconds(5)); // every 5 seconds
builder.setBulkActions(5000); // every 5,000 operations
builder.setBulkSize(new ByteSizeValue(5, ByteSizeUnit.MB)); // every 5 MB
BulkProcessor bulkProcessor = builder.build();
Random random = new Random();
for (int i = 0; i < bulkTotal; i++) {
// Replace with actual field names and values for your use case.
Map<String, Object> map = new HashMap<>();
map.put("field1", random.nextInt() + "");
map.put("field2", random.nextInt() + "");
IndexRequest bulkItemRequest = new IndexRequest(index_name);
bulkItemRequest.source(map);
bulkProcessor.add(bulkItemRequest);
}
// Wait up to 120 seconds for all pending operations to complete.
bulkProcessor.awaitClose(120, TimeUnit.SECONDS);
long failure = failedBulkItemCount.get(),
success = bulkTotal - failure;
System.out.println("Bulk using BulkProcessor finished with [" + success + "] requests succeeded, [" + failure + "] requests failed.");
Pesquisar documentos
Envie primeiro uma solicitação de atualização (refresh) para tornar visíveis os dados gravados recentemente e, em seguida, execute as consultas.
Por padrão, uma consulta de pesquisa retorna no máximo 10.000 documentos. Para obter a contagem total exata quando houver mais de 10.000 documentos correspondentes, chame searchSourceBuilder.trackTotalHits(true) antes de executar a consulta.
// Refresh the index to make written data searchable.
RefreshRequest refreshRequest = new RefreshRequest(index_name);
highClient.indices().refresh(refreshRequest, COMMON_OPTIONS);
System.out.println("Refresh on index [" + index_name + "] successfully.");
// Query all documents.
SearchRequest searchRequest = new SearchRequest(index_name);
SearchSourceBuilder searchSourceBuilder = new SearchSourceBuilder();
// Uncomment the next line to get the exact total count when results exceed 10,000.
// searchSourceBuilder.trackTotalHits(true);
QueryBuilder queryMatchAllBuilder = new MatchAllQueryBuilder();
searchSourceBuilder.query(queryMatchAllBuilder);
searchRequest.source(searchSourceBuilder);
SearchResponse searchResponse = highClient.search(searchRequest, COMMON_OPTIONS);
long totalHit = searchResponse.getHits().getTotalHits().value;
System.out.println("Search query match all hits [" + totalHit + "] in total.");
// Query documents by ID.
QueryBuilder queryByIdBuilder = new MatchQueryBuilder("_id", doc_id);
searchSourceBuilder.query(queryByIdBuilder);
searchRequest.source(searchSourceBuilder);
searchResponse = highClient.search(searchRequest, COMMON_OPTIONS);
for (SearchHit searchHit : searchResponse.getHits()) {
System.out.println("Search query by id response [" + searchHit.getSourceAsString() + "]");
}
Exclua documentos e índices
Use DeleteRequest para remover um único documento e DeleteIndexRequest para excluir um índice.
// Delete a single document by ID.
DeleteRequest deleteRequest = new DeleteRequest(index_name);
deleteRequest.id(doc_id);
DeleteResponse deleteResponse = highClient.delete(deleteRequest, COMMON_OPTIONS);
System.out.println("Delete document with id [" + deleteResponse.getId() + "] successfully.");
// Delete the index.
DeleteIndexRequest deleteIndexRequest = new DeleteIndexRequest(index_name);
AcknowledgedResponse deleteIndexResponse = highClient.indices().delete(deleteIndexRequest, COMMON_OPTIONS);
if (deleteIndexResponse.isAcknowledged()) {
System.out.println("Delete index [" + index_name + "] successfully.");
}
highClient.close();
Exemplo completo
import org.apache.http.HttpHost;
import org.apache.http.auth.AuthScope;
import org.apache.http.auth.UsernamePasswordCredentials;
import org.apache.http.client.CredentialsProvider;
import org.apache.http.impl.client.BasicCredentialsProvider;
import org.apache.http.impl.nio.client.HttpAsyncClientBuilder;
import org.elasticsearch.action.admin.indices.delete.DeleteIndexRequest;
import org.elasticsearch.action.admin.indices.refresh.RefreshRequest;
import org.elasticsearch.action.admin.indices.refresh.RefreshResponse;
import org.elasticsearch.action.bulk.BulkItemResponse;
import org.elasticsearch.action.bulk.BulkProcessor;
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.index.IndexRequest;
import org.elasticsearch.action.index.IndexResponse;
import org.elasticsearch.action.search.SearchRequest;
import org.elasticsearch.action.search.SearchResponse;
import org.elasticsearch.action.support.master.AcknowledgedResponse;
import org.elasticsearch.client.HttpAsyncResponseConsumerFactory;
import org.elasticsearch.client.RequestOptions;
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.CreateIndexResponse;
import org.elasticsearch.common.unit.ByteSizeUnit;
import org.elasticsearch.common.unit.ByteSizeValue;
import org.elasticsearch.common.unit.TimeValue;
import org.elasticsearch.index.query.MatchAllQueryBuilder;
import org.elasticsearch.index.query.MatchQueryBuilder;
import org.elasticsearch.index.query.QueryBuilder;
import org.elasticsearch.search.SearchHit;
import org.elasticsearch.search.builder.SearchSourceBuilder;
import java.util.HashMap;
import java.util.Map;
import java.util.Random;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
public class RestHClientTest {
private static final RequestOptions COMMON_OPTIONS;
static {
// Set the maximum response buffer size to 30 MB. The default is 100 MB.
RequestOptions.Builder builder = RequestOptions.DEFAULT.toBuilder();
builder.setHttpAsyncResponseConsumerFactory(
new HttpAsyncResponseConsumerFactory
.HeapBufferedResponseConsumerFactory(30 * 1024 * 1024));
COMMON_OPTIONS = builder.build();
}
public static void main(String[] args) {
// Set the Elasticsearch-compatible endpoint and port for LindormSearch.
String search_url = "ld-t4n5668xk31ui****-proxy-search-public.lindorm.rds.aliyuncs.com";
int search_port = 30070;
// Set the username and password. Retrieve them from the Lindorm console.
String username = "user";
String password = "test";
final CredentialsProvider credentials_provider = new BasicCredentialsProvider();
credentials_provider.setCredentials(AuthScope.ANY, new UsernamePasswordCredentials(username, password));
RestHighLevelClient highClient = new RestHighLevelClient(
RestClient.builder(new HttpHost(search_url, search_port, "http"))
.setHttpClientConfigCallback(new RestClientBuilder.HttpClientConfigCallback() {
public HttpAsyncClientBuilder customizeHttpClient(HttpAsyncClientBuilder httpClientBuilder) {
return httpClientBuilder.setDefaultCredentialsProvider(credentials_provider);
}
})
);
try {
String index_name = "lindorm_index";
// Create an index.
CreateIndexRequest createIndexRequest = new CreateIndexRequest(index_name);
Map<String, Object> settingsMap = new HashMap<>();
settingsMap.put("index.number_of_shards", 4);
createIndexRequest.settings(settingsMap);
CreateIndexResponse createIndexResponse = highClient.indices().create(createIndexRequest, COMMON_OPTIONS);
if (createIndexResponse.isAcknowledged()) {
System.out.println("Create index [" + index_name + "] successfully.");
}
// Index a single document.
// Specify the document ID. If you do not specify the document ID, an ID is automatically generated,
// which can improve write performance.
String doc_id = "test";
Map<String, Object> jsonMap = new HashMap<>();
jsonMap.put("field1", "value1");
jsonMap.put("field2", "value2");
IndexRequest indexRequest = new IndexRequest(index_name);
indexRequest.id(doc_id).source(jsonMap);
IndexResponse indexResponse = highClient.index(indexRequest, COMMON_OPTIONS);
System.out.println("Index document with id[" + indexResponse.getId() + "] successfully.");
// Bulk index documents using BulkProcessor.
int bulkTotal = 100000;
AtomicLong failedBulkItemCount = new AtomicLong();
BulkProcessor.Builder builder = BulkProcessor.builder(
(request, bulkListener) -> highClient.bulkAsync(request, COMMON_OPTIONS, bulkListener),
new BulkProcessor.Listener() {
@Override
public void beforeBulk(long executionId, BulkRequest request) {}
@Override
public void afterBulk(long executionId, BulkRequest request, BulkResponse response) {
for (BulkItemResponse bulkItemResponse : response) {
if (bulkItemResponse.isFailed()) {
failedBulkItemCount.incrementAndGet();
}
}
}
@Override
public void afterBulk(long executionId, BulkRequest request, Throwable failure) {
if (null != failure) {
failedBulkItemCount.addAndGet(request.numberOfActions());
}
}
});
builder.setConcurrentRequests(10);
builder.setFlushInterval(TimeValue.timeValueSeconds(5));
builder.setBulkActions(5000);
builder.setBulkSize(new ByteSizeValue(5, ByteSizeUnit.MB));
BulkProcessor bulkProcessor = builder.build();
Random random = new Random();
for (int i = 0; i < bulkTotal; i++) {
Map<String, Object> map = new HashMap<>();
map.put("field1", random.nextInt() + "");
map.put("field2", random.nextInt() + "");
IndexRequest bulkItemRequest = new IndexRequest(index_name);
bulkItemRequest.source(map);
bulkProcessor.add(bulkItemRequest);
}
bulkProcessor.awaitClose(120, TimeUnit.SECONDS);
long failure = failedBulkItemCount.get(),
success = bulkTotal - failure;
System.out.println("Bulk using BulkProcessor finished with [" + success + "] requests succeeded, [" + failure + "] requests failed.");
// Refresh the index to make written data searchable.
RefreshRequest refreshRequest = new RefreshRequest(index_name);
RefreshResponse refreshResponse = highClient.indices().refresh(refreshRequest, COMMON_OPTIONS);
System.out.println("Refresh on index [" + index_name + "] successfully.");
// Query all documents. By default, at most 10,000 results are returned.
// To get the exact total count, call searchSourceBuilder.trackTotalHits(true).
SearchRequest searchRequest = new SearchRequest(index_name);
SearchSourceBuilder searchSourceBuilder = new SearchSourceBuilder();
QueryBuilder queryMatchAllBuilder = new MatchAllQueryBuilder();
searchSourceBuilder.query(queryMatchAllBuilder);
searchRequest.source(searchSourceBuilder);
SearchResponse searchResponse = highClient.search(searchRequest, COMMON_OPTIONS);
long totalHit = searchResponse.getHits().getTotalHits().value;
System.out.println("Search query match all hits [" + totalHit + "] in total.");
// Query documents by ID.
QueryBuilder queryByIdBuilder = new MatchQueryBuilder("_id", doc_id);
searchSourceBuilder.query(queryByIdBuilder);
searchRequest.source(searchSourceBuilder);
searchResponse = highClient.search(searchRequest, COMMON_OPTIONS);
for (SearchHit searchHit : searchResponse.getHits()) {
System.out.println("Search query by id response [" + searchHit.getSourceAsString() + "]");
}
// Delete a document by ID.
DeleteRequest deleteRequest = new DeleteRequest(index_name);
deleteRequest.id(doc_id);
DeleteResponse deleteResponse = highClient.delete(deleteRequest, COMMON_OPTIONS);
System.out.println("Delete document with id [" + deleteResponse.getId() + "] successfully.");
// Delete the index.
DeleteIndexRequest deleteIndexRequest = new DeleteIndexRequest(index_name);
AcknowledgedResponse deleteIndexResponse = highClient.indices().delete(deleteIndexRequest, COMMON_OPTIONS);
if (deleteIndexResponse.isAcknowledged()) {
System.out.println("Delete index [" + index_name + "] successfully.");
}
highClient.close();
} catch (Exception exception) {
System.out.println("msg " + exception);
}
}
}
A saída esperada é:
Create index [lindorm_index] successfully.
Index document with id[test] successfully.
Bulk using BulkProcessor finished with [100000] requests succeeded, [0] requests failed.
Refresh on index [lindorm_index] successfully.
Search query match all hits [10000] in total.
Search query by id response [{"field1":"value1","field2":"value2"}]
Delete document with id [test] successfully.
Delete index [lindorm_index] successfully.
A consulta match-all retorna 10.000 resultados mesmo com 100.000 documentos gravados. Esse é o limite padrão de resultados. Para obter a contagem total exata, chame searchSourceBuilder.trackTotalHits(true) antes de executar a consulta.