From cdc4dd4e0b72dd565610197103e61148dfd206ce Mon Sep 17 00:00:00 2001 From: Keith Massey Date: Tue, 6 Feb 2024 08:27:55 -0600 Subject: [PATCH 1/5] Releasing request memory from BulkRequestBuilder --- .../action/bulk/BulkRequestBuilder.java | 12 ++++++++++-- .../action/bulk/BulkRequestBuilderTests.java | 13 +++++++++++++ 2 files changed, 23 insertions(+), 2 deletions(-) diff --git a/server/src/main/java/org/elasticsearch/action/bulk/BulkRequestBuilder.java b/server/src/main/java/org/elasticsearch/action/bulk/BulkRequestBuilder.java index 24d6fad554935..e5f3202e08846 100644 --- a/server/src/main/java/org/elasticsearch/action/bulk/BulkRequestBuilder.java +++ b/server/src/main/java/org/elasticsearch/action/bulk/BulkRequestBuilder.java @@ -28,6 +28,7 @@ import java.io.IOException; import java.util.ArrayList; +import java.util.Iterator; import java.util.List; /** @@ -51,6 +52,7 @@ public class BulkRequestBuilder extends ActionRequestLazyBuilder, ? extends DocWriteResponse> requestBuilder : requestBuilders) { - DocWriteRequest childRequest = requestBuilder.request(); + for (Iterator, ? extends DocWriteResponse>> requestsIter = requestBuilders + .iterator(); requestsIter.hasNext();) { + DocWriteRequest childRequest = requestsIter.next().request(); request.add(childRequest); + requestsIter.remove(); // The inner request builder can now be garbage collected } for (DocWriteRequest childRequest : requests) { request.add(childRequest); diff --git a/server/src/test/java/org/elasticsearch/action/bulk/BulkRequestBuilderTests.java b/server/src/test/java/org/elasticsearch/action/bulk/BulkRequestBuilderTests.java index 27b1104163d67..2a5916f053b75 100644 --- a/server/src/test/java/org/elasticsearch/action/bulk/BulkRequestBuilderTests.java +++ b/server/src/test/java/org/elasticsearch/action/bulk/BulkRequestBuilderTests.java @@ -12,6 +12,8 @@ import org.elasticsearch.action.index.IndexRequestBuilder; import org.elasticsearch.test.ESTestCase; +import static org.hamcrest.Matchers.equalTo; + public class BulkRequestBuilderTests extends ESTestCase { public void testValidation() { @@ -20,4 +22,15 @@ public void testValidation() { bulkRequestBuilder.add(new IndexRequest()); expectThrows(IllegalStateException.class, bulkRequestBuilder::request); } + + public void testRequestTwice() { + BulkRequestBuilder bulkRequestBuilder = new BulkRequestBuilder(null, null); + bulkRequestBuilder.add(new IndexRequestBuilder(null, randomAlphaOfLength(10))); + bulkRequestBuilder.add(new IndexRequestBuilder(null, randomAlphaOfLength(10))); + bulkRequestBuilder.add(new IndexRequestBuilder(null, randomAlphaOfLength(10))); + BulkRequest bulkRequest = bulkRequestBuilder.request(); + assertNotNull(bulkRequest); + assertThat(bulkRequest.numberOfActions(), equalTo(3)); + expectThrows(IllegalStateException.class, bulkRequestBuilder::request); + } } From 361387533d0ba7af1631a4188b9fd6f0cb2229a7 Mon Sep 17 00:00:00 2001 From: Keith Massey Date: Tue, 6 Feb 2024 08:39:11 -0600 Subject: [PATCH 2/5] Adding a couple of test assertions --- .../org/elasticsearch/action/bulk/BulkRequestBuilderTests.java | 3 +++ 1 file changed, 3 insertions(+) diff --git a/server/src/test/java/org/elasticsearch/action/bulk/BulkRequestBuilderTests.java b/server/src/test/java/org/elasticsearch/action/bulk/BulkRequestBuilderTests.java index 2a5916f053b75..77930e98cf7ab 100644 --- a/server/src/test/java/org/elasticsearch/action/bulk/BulkRequestBuilderTests.java +++ b/server/src/test/java/org/elasticsearch/action/bulk/BulkRequestBuilderTests.java @@ -28,9 +28,12 @@ public void testRequestTwice() { bulkRequestBuilder.add(new IndexRequestBuilder(null, randomAlphaOfLength(10))); bulkRequestBuilder.add(new IndexRequestBuilder(null, randomAlphaOfLength(10))); bulkRequestBuilder.add(new IndexRequestBuilder(null, randomAlphaOfLength(10))); + assertThat(bulkRequestBuilder.numberOfActions(), equalTo(3)); BulkRequest bulkRequest = bulkRequestBuilder.request(); assertNotNull(bulkRequest); assertThat(bulkRequest.numberOfActions(), equalTo(3)); + // Make sure that the bulk request builder is no longer holding onto the child request builders: + assertThat(bulkRequestBuilder.numberOfActions(), equalTo(0)); expectThrows(IllegalStateException.class, bulkRequestBuilder::request); } } From 80e895287c75285059269e8869cc3fa18360bd01 Mon Sep 17 00:00:00 2001 From: Keith Massey Date: Tue, 6 Feb 2024 09:17:42 -0600 Subject: [PATCH 3/5] Adding a runtime assertion --- .../java/org/elasticsearch/action/bulk/BulkRequestBuilder.java | 1 + .../org/elasticsearch/action/bulk/BulkRequestBuilderTests.java | 2 +- 2 files changed, 2 insertions(+), 1 deletion(-) diff --git a/server/src/main/java/org/elasticsearch/action/bulk/BulkRequestBuilder.java b/server/src/main/java/org/elasticsearch/action/bulk/BulkRequestBuilder.java index e5f3202e08846..5485b70ffee79 100644 --- a/server/src/main/java/org/elasticsearch/action/bulk/BulkRequestBuilder.java +++ b/server/src/main/java/org/elasticsearch/action/bulk/BulkRequestBuilder.java @@ -201,6 +201,7 @@ public BulkRequestBuilder setRefreshPolicy(String refreshPolicy) { @Override public BulkRequest request() { + assert requestPreviouslyCalled == false : "Cannot call request() multiple times on the same BulkRequestBuilder object"; if (requestPreviouslyCalled) { throw new IllegalStateException("Cannot call request() multiple times on the same BulkRequestBuilder object"); } diff --git a/server/src/test/java/org/elasticsearch/action/bulk/BulkRequestBuilderTests.java b/server/src/test/java/org/elasticsearch/action/bulk/BulkRequestBuilderTests.java index 77930e98cf7ab..9cfe033f2788c 100644 --- a/server/src/test/java/org/elasticsearch/action/bulk/BulkRequestBuilderTests.java +++ b/server/src/test/java/org/elasticsearch/action/bulk/BulkRequestBuilderTests.java @@ -34,6 +34,6 @@ public void testRequestTwice() { assertThat(bulkRequest.numberOfActions(), equalTo(3)); // Make sure that the bulk request builder is no longer holding onto the child request builders: assertThat(bulkRequestBuilder.numberOfActions(), equalTo(0)); - expectThrows(IllegalStateException.class, bulkRequestBuilder::request); + expectThrows(AssertionError.class, bulkRequestBuilder::request); } } From 40a280c2d401f14f22ed9cac117b31c88ac7d72e Mon Sep 17 00:00:00 2001 From: Keith Massey Date: Wed, 7 Feb 2024 08:18:08 -0600 Subject: [PATCH 4/5] using an ArrayDeque for better runtime performance --- .../action/bulk/BulkRequestBuilder.java | 28 +++++++++++-------- 1 file changed, 16 insertions(+), 12 deletions(-) diff --git a/server/src/main/java/org/elasticsearch/action/bulk/BulkRequestBuilder.java b/server/src/main/java/org/elasticsearch/action/bulk/BulkRequestBuilder.java index 5485b70ffee79..92af9896b30d6 100644 --- a/server/src/main/java/org/elasticsearch/action/bulk/BulkRequestBuilder.java +++ b/server/src/main/java/org/elasticsearch/action/bulk/BulkRequestBuilder.java @@ -27,8 +27,10 @@ import org.elasticsearch.xcontent.XContentType; import java.io.IOException; +import java.util.ArrayDeque; import java.util.ArrayList; -import java.util.Iterator; +import java.util.Collection; +import java.util.Deque; import java.util.List; /** @@ -45,8 +47,8 @@ public class BulkRequestBuilder extends ActionRequestLazyBuilder> requests = new ArrayList<>(); private final List framedData = new ArrayList<>(); - private final List, ? extends DocWriteResponse>> requestBuilders = - new ArrayList<>(); + private final Deque, ? extends DocWriteResponse>> requestBuilders = + new ArrayDeque<>(); private ActiveShardCount waitForActiveShards; private TimeValue timeout; private String globalPipeline; @@ -208,11 +210,13 @@ public BulkRequest request() { requestPreviouslyCalled = true; validate(); BulkRequest request = new BulkRequest(globalIndex); - for (Iterator, ? extends DocWriteResponse>> requestsIter = requestBuilders - .iterator(); requestsIter.hasNext();) { - DocWriteRequest childRequest = requestsIter.next().request(); - request.add(childRequest); - requestsIter.remove(); // The inner request builder can now be garbage collected + /* + * In the following loop we intentionally remove the builders from requestBuilders so that they can be garbage collected. This is + * so that we don't require double the memory of all of the inner requests, which could be really bad for a lage bulk request. + */ + for (ActionRequestLazyBuilder, ? extends DocWriteResponse> builder = requestBuilders + .pollFirst(); builder != null; builder = requestBuilders.pollFirst()) { + request.add(builder.request()); } for (DocWriteRequest childRequest : requests) { request.add(childRequest); @@ -243,17 +247,17 @@ public BulkRequest request() { } private void validate() { - if (countNonEmptyLists(requestBuilders, requests, framedData) > 1) { + if (countNonEmptyCollections(requestBuilders, requests, framedData) > 1) { throw new IllegalStateException( "Must use only request builders, requests, or byte arrays within a single bulk request. Cannot mix and match" ); } } - private int countNonEmptyLists(List... lists) { + private int countNonEmptyCollections(Collection... collections) { int sum = 0; - for (List list : lists) { - if (list.isEmpty() == false) { + for (Collection collection : collections) { + if (collection.isEmpty() == false) { sum++; } } From 289152ddbc176f09edd3d1aa3ce0c43ba210c260 Mon Sep 17 00:00:00 2001 From: Keith Massey Date: Wed, 7 Feb 2024 09:05:34 -0600 Subject: [PATCH 5/5] simplifying for loop --- .../java/org/elasticsearch/action/bulk/BulkRequestBuilder.java | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/server/src/main/java/org/elasticsearch/action/bulk/BulkRequestBuilder.java b/server/src/main/java/org/elasticsearch/action/bulk/BulkRequestBuilder.java index 92af9896b30d6..2e2938b63334e 100644 --- a/server/src/main/java/org/elasticsearch/action/bulk/BulkRequestBuilder.java +++ b/server/src/main/java/org/elasticsearch/action/bulk/BulkRequestBuilder.java @@ -214,8 +214,7 @@ public BulkRequest request() { * In the following loop we intentionally remove the builders from requestBuilders so that they can be garbage collected. This is * so that we don't require double the memory of all of the inner requests, which could be really bad for a lage bulk request. */ - for (ActionRequestLazyBuilder, ? extends DocWriteResponse> builder = requestBuilders - .pollFirst(); builder != null; builder = requestBuilders.pollFirst()) { + for (var builder = requestBuilders.pollFirst(); builder != null; builder = requestBuilders.pollFirst()) { request.add(builder.request()); } for (DocWriteRequest childRequest : requests) {