diff --git a/docs/en/changes/changes.md b/docs/en/changes/changes.md index b860e4e4b2..3411f62abb 100644 --- a/docs/en/changes/changes.md +++ b/docs/en/changes/changes.md @@ -5,6 +5,7 @@ #### OAP Server +* ElasticSearchClient: Add `deleteById` API. #### UI diff --git a/oap-server/server-library/library-client/src/main/java/org/apache/skywalking/oap/server/library/client/elasticsearch/ElasticSearchClient.java b/oap-server/server-library/library-client/src/main/java/org/apache/skywalking/oap/server/library/client/elasticsearch/ElasticSearchClient.java index eea5004e0d..a30fd72ccf 100644 --- a/oap-server/server-library/library-client/src/main/java/org/apache/skywalking/oap/server/library/client/elasticsearch/ElasticSearchClient.java +++ b/oap-server/server-library/library-client/src/main/java/org/apache/skywalking/oap/server/library/client/elasticsearch/ElasticSearchClient.java @@ -346,6 +346,12 @@ public class ElasticSearchClient implements Client, HealthCheckable { es.get().documents().update(wrapper.getRequest(), params); } + public void deleteById(String indexName, String id) { + indexName = indexNameConverter.apply(indexName); + Map params = ImmutableMap.of("refresh", "true"); + es.get().documents().deleteById(indexName, TYPE, id, params); + } + public IndexRequestWrapper prepareInsert(String indexName, String id, Map source) { return prepareInsert(indexName, id, Optional.empty(), source); diff --git a/oap-server/server-library/library-client/src/test/java/org/apache/skywalking/library/elasticsearch/bulk/ElasticSearchIT.java b/oap-server/server-library/library-client/src/test/java/org/apache/skywalking/library/elasticsearch/bulk/ElasticSearchIT.java index a2ebcc79f9..cdaca0a2b4 100644 --- a/oap-server/server-library/library-client/src/test/java/org/apache/skywalking/library/elasticsearch/bulk/ElasticSearchIT.java +++ b/oap-server/server-library/library-client/src/test/java/org/apache/skywalking/library/elasticsearch/bulk/ElasticSearchIT.java @@ -169,7 +169,8 @@ public class ElasticSearchIT { .next() .getSource() .get("message")); - + client.deleteById(indexName, id); + Assertions.assertFalse(client.existDoc(indexName, id)); client.shutdown(); server.stop(); } diff --git a/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/client/DocumentClient.java b/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/client/DocumentClient.java index b6d242e3c0..33d6acec33 100644 --- a/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/client/DocumentClient.java +++ b/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/client/DocumentClient.java @@ -157,4 +157,26 @@ public final class DocumentClient { }); future.join(); } + + @SneakyThrows + public void deleteById(String index, String type, String id, Map params) { + final CompletableFuture future = version.thenCompose( + v -> client.execute(v.requestFactory().document().deleteById(index, type, id, params)) + .aggregate().thenAccept(response -> { + final HttpStatus status = response.status(); + if (status != HttpStatus.OK) { + throw new RuntimeException(response.contentUtf8()); + } + })); + future.whenComplete((result, exception) -> { + if (exception != null) { + log.error("Failed to delete doc by id {} in index {}, params: {}", id, index, params, exception); + return; + } + if (log.isDebugEnabled()) { + log.debug("Succeeded delete doc by id {} in index {}, params: {}", id, params, index); + } + }); + future.join(); + } } diff --git a/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/requests/factory/DocumentFactory.java b/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/requests/factory/DocumentFactory.java index 946918fd5c..d1207792fe 100644 --- a/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/requests/factory/DocumentFactory.java +++ b/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/requests/factory/DocumentFactory.java @@ -60,4 +60,9 @@ public interface DocumentFactory { */ HttpRequest delete(String index, String type, Query query, Map params); + + /** + * Returns a request to delete documents matching the given {@code id} in {@code index}. + */ + HttpRequest deleteById(String index, String type, String id, Map params); } diff --git a/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/requests/factory/v6/V6DocumentFactory.java b/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/requests/factory/v6/V6DocumentFactory.java index 71fe55d7bc..63986e1764 100644 --- a/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/requests/factory/v6/V6DocumentFactory.java +++ b/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/requests/factory/v6/V6DocumentFactory.java @@ -198,4 +198,23 @@ final class V6DocumentFactory implements DocumentFactory { .content(MediaType.JSON, content) .build(); } + + @Override + public HttpRequest deleteById(final String index, final String type, final String id, Map params) { + checkArgument(!isNullOrEmpty(index), "index cannot be null or empty"); + checkArgument(!isNullOrEmpty(type), "type cannot be null or empty"); + checkArgument(!isNullOrEmpty(id), "id cannot be null or empty"); + + final HttpRequestBuilder builder = HttpRequest.builder(); + if (params != null) { + params.forEach(builder::queryParam); + } + + return builder + .delete("/{index}/{type}/{id}") + .pathParam("index", index) + .pathParam("type", type) + .pathParam("id", id) + .build(); + } } diff --git a/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/requests/factory/v7plus/V7DocumentFactory.java b/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/requests/factory/v7plus/V7DocumentFactory.java index 1f0abaef56..cd8bf763e5 100644 --- a/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/requests/factory/v7plus/V7DocumentFactory.java +++ b/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/requests/factory/v7plus/V7DocumentFactory.java @@ -191,4 +191,22 @@ final class V7DocumentFactory implements DocumentFactory { .content(MediaType.JSON, content) .build(); } + + @Override + public HttpRequest deleteById(final String index, final String type, final String id, Map params) { + checkArgument(!isNullOrEmpty(index), "index cannot be null or empty"); + checkArgument(!isNullOrEmpty(type), "type cannot be null or empty"); + checkArgument(!isNullOrEmpty(id), "id cannot be null or empty"); + + final HttpRequestBuilder builder = HttpRequest.builder(); + if (params != null) { + params.forEach(builder::queryParam); + } + + return builder + .delete("/{index}/_doc/{id}") + .pathParam("index", index) + .pathParam("id", id) + .build(); + } } diff --git a/oap-server/server-library/library-elasticsearch-client/src/test/java/org/apache/skywalking/library/elasticsearch/ElasticSearchIT.java b/oap-server/server-library/library-elasticsearch-client/src/test/java/org/apache/skywalking/library/elasticsearch/ElasticSearchIT.java index ef9a94c840..bc1130832d 100644 --- a/oap-server/server-library/library-elasticsearch-client/src/test/java/org/apache/skywalking/library/elasticsearch/ElasticSearchIT.java +++ b/oap-server/server-library/library-elasticsearch-client/src/test/java/org/apache/skywalking/library/elasticsearch/ElasticSearchIT.java @@ -389,4 +389,37 @@ public class ElasticSearchIT { server.close(); } + + @ParameterizedTest(name = "version: {0}") + @MethodSource("es") + public void testDocDeleteById(final String ignored, + final ElasticsearchContainer server) { + server.start(); + + final ElasticSearch client = + ElasticSearch.builder() + .endpoints(server.getHttpHostAddress()) + .build(); + client.connect(); + + final String index = "test-index-delete"; + assertTrue(client.index().create(index, null, null)); + + final ImmutableMap doc = ImmutableMap.of("key", "val"); + final String idWithSpace = "an id"; // UI management templates' IDs contains spaces + final String type = "type"; + + client.documents().index( + IndexRequest.builder() + .index(index) + .type(type) + .id(idWithSpace) + .doc(doc) + .build(), null); + + assertTrue(client.documents().exists(index, type, idWithSpace)); + client.documents().deleteById(index, type, idWithSpace, ImmutableMap.of("refresh", "true")); + assertFalse(client.documents().exists(index, type, idWithSpace)); + server.close(); + } }