From ad9d616297fbb535cd0e7c341ae52b2ed60397a1 Mon Sep 17 00:00:00 2001 From: kezhenxu94 Date: Fri, 24 Sep 2021 22:35:56 +0800 Subject: [PATCH] Add `socketTimeout` back to the new implementation (#7798) --- .../client/elasticsearch/ElasticSearchClient.java | 1 + .../library/elasticsearch/ElasticSearchBuilder.java | 9 +++++++++ .../library/elasticsearch/bulk/BulkProcessor.java | 7 ++++++- 3 files changed, 16 insertions(+), 1 deletion(-) 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 b6801bf48b..03a91161fd 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 @@ -117,6 +117,7 @@ public class ElasticSearchClient implements Client, HealthCheckable { .endpoints(clusterNodes.split(",")) .protocol(protocol) .connectTimeout(connectTimeout) + .socketTimeout(socketTimeout) .numHttpClientThread(numHttpClientThread) .healthyListener(healthy -> { if (healthy) { diff --git a/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/ElasticSearchBuilder.java b/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/ElasticSearchBuilder.java index d154f3277b..06975df0c2 100644 --- a/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/ElasticSearchBuilder.java +++ b/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/ElasticSearchBuilder.java @@ -63,6 +63,8 @@ public final class ElasticSearchBuilder { private Duration connectTimeout = Duration.ofMillis(500); + private Duration socketTimeout = Duration.ofSeconds(30); + private Consumer healthyListener; private int numHttpClientThread; @@ -117,6 +119,12 @@ public final class ElasticSearchBuilder { return this; } + public ElasticSearchBuilder socketTimeout(int socketTimeout) { + checkArgument(socketTimeout > 0, "socketTimeout must be positive"); + this.socketTimeout = Duration.ofMillis(socketTimeout); + return this; + } + public ElasticSearchBuilder healthyListener(Consumer healthyListener) { requireNonNull(healthyListener, "healthyListener"); this.healthyListener = healthyListener; @@ -138,6 +146,7 @@ public final class ElasticSearchBuilder { final ClientFactoryBuilder factoryBuilder = ClientFactory.builder() .connectTimeout(connectTimeout) + .idleTimeout(socketTimeout) .useHttp2Preface(false) .workerGroup(numHttpClientThread > 0 ? numHttpClientThread : NUM_PROC); diff --git a/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/bulk/BulkProcessor.java b/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/bulk/BulkProcessor.java index a8f72ec8e4..93c49aec5c 100644 --- a/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/bulk/BulkProcessor.java +++ b/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/bulk/BulkProcessor.java @@ -84,9 +84,10 @@ public final class BulkProcessor { return this; } + @SneakyThrows private void internalAdd(Object request) { requireNonNull(request, "request"); - requests.add(request); + requests.put(request); flushIfNeeded(); } @@ -120,6 +121,10 @@ public final class BulkProcessor { private CompletableFuture doFlush(final List batch) { log.debug("Executing bulk with {} requests", batch.size()); + if (batch.isEmpty()) { + return CompletableFuture.completedFuture(null); + } + final CompletableFuture future = es.get().version().thenCompose(v -> { try { final RequestFactory rf = v.requestFactory();