Skip to content

Send Elasticsearch bulk bodies without copying them and filter their responses - #5025

Open
porunov wants to merge 1 commit into
JanusGraph:masterfrom
porunov:es-bulk-request-path
Open

porunov wants to merge 1 commit into
JanusGraph:masterfrom
porunov:es-bulk-request-path

Conversation

@porunov

@porunov porunov commented Oct 5, 2026

Copy link
Copy Markdown
Member

Fixes #5024

An Elasticsearch bulk request held its body in memory up to four times to send it, parsed every field of a response of which the client reads two for each item, and built a map for the action line of each document. Its action also depended on the JVM's default locale.

Changes

  • BulkRequestEntity (new, in RestElasticSearchClient): the body of a bulk request, which writes the items' serialized bytes as they are (writeTo, which the gzip wrapper of index.[X].elasticsearch.compression uses) and reads them through a stream (getContent, which the HTTP client uses otherwise). It knows its length and is repeatable, so the client can send it again to another node; a reattempt builds a new one over the failed items. Its stream behaves as the ByteArrayInputStream of the entity it replaces did: a read fills the buffer across items, available() counts what is left, and it supports mark and reset, which request signers such as those of the AWS SDK use to hash the body before it is sent. It replaces buildBulkRequestInput, which copied the items into a ByteArrayOutputStream and from there into an array, on every attempt. With index.[X].elasticsearch.compression the client still compresses the body into a buffer of its own.
  • Bulk responses: a bulk request asks for filter_path=errors,items.index.status,items.index.error,items.update.status,items.update.error,items.delete.status,items.delete.error: each item's status and error, which the client reads, under the actions JanusGraph sends, and errors, which keeps RestBulkResponse.isErrors() meaningful. Every item keeps its status, so the items line up with the requests they answer; filtering by the error alone leaves out the items which succeeded (checked against Elasticsearch 9.5.4). The actions are named from ElasticSearchMutation.RequestType rather than matched with *, so that the query string holds no asterisk, which a request signer would have to percent-encode the way the cluster does. The response to 1,000 index requests is 25,026 bytes instead of 172,842. RestBulkItemResponse.getResult(), which nothing read, is deprecated, as the response no longer carries it.
  • Action lines: written field by field with a JsonGenerator, instead of as a HashMap in an ImmutableMap serialized by the ObjectMapper, and the action is lower-cased in the root locale. Under a Turkish default locale the action of an index request was ındex, with a dotless i, and Elasticsearch rejected the whole bulk request: Malformed action/metadata line [1], expected field [create], [delete], [index] or [update] but found [ındex].
  • Changelog entry.

Measurements

A harness (not shipped) sends the same bulk of index requests through RestElasticSearchClient.bulkRequest to Elasticsearch 9.5.4 in Docker (one shard, nothing indexed but the source), with master's RestElasticSearchClient and RestBulkResponse against this branch's on the same classpath, in two rounds of alternating order. The time is the median per bulk; the allocation counts every thread, the HTTP client's I/O threads included.

Bulk master this branch
1,000 documents of 1 KB 10.6–13.9 ms, 6.65 MB 8.8–9.8 ms, 2.57 MB
10,000 documents of 100 bytes 39.4–39.6 ms, 30.4 MB 32.4–37.1 ms, 15.9 MB
10,000 documents of 1 KB 84.6–87.0 ms, 78.1 MB 72.1–76.3 ms, 25.2 MB
50,000 documents of 1 KB 428–445 ms, 365.7 MB 385–394 ms, 125.6 MB

The smallest heap in which the bulk of 50,000 documents of 1 KB completes is 320 MB on master in one run and 384 MB in the other, below which it runs out of memory in ByteArrayOutputStream. On this branch it is 144 MB, below which the serialization of the documents themselves runs out.

Measured alone, building the item of a 100-byte document takes 229 ns instead of about 400, and of a 1 KB document about 790 ns instead of 1,440. Parsing a response takes 52 ns per item instead of 200.

Tests

  • RestClientBulkRequestsTest: the entity's body is the items' bytes in order, through writeTo and through its stream read twice, and its length is theirs; the stream gives the same bytes whatever the size of the buffer and the offset it reads into, filling each buffer across items, with available() counting down, and reads the rest again after a reset from a mark at every position, or from the start without one; a bulk request carries the filter and the entity; the action line of each operation, with characters which JSON escapes, the mapping type for Elasticsearch 6, and a Turkish default locale, which yields ındex when the action is lower-cased in the default locale.
  • RestClientRetryTest: the reattempt of a bulk request sends the failed item alone.
  • ElasticsearchIndexTest.testTheFailedItemsOfABulkResponseAreReportedForTheirOwnDocuments (new, Elasticsearch in a container): updates of missing documents between index requests which succeed fail with a 404, and the documents reported are exactly theirs. With a filter of the errors alone it fails, as the successful items are left out and the failures shift onto other documents.
  • ElasticsearchIndexTest.testCompressedRequests sends bulk requests through the gzip wrapper, which writes the entity with writeTo.
  • Locally, Java 11, against Elasticsearch 9.5.4 in a container: ElasticsearchIndexTest (288) and the RestClient*Test classes (77) pass on the final code. BerkeleyElasticsearchTest (92, 1 skipped), ElasticsearchConfigTest (12) and ElasticSearchIndexReattemptTest (7) passed before the last two changes, mark and reset, and the actions named in the filter.

🤖 Generated with Claude Code

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot review overview

🟢 Approval recommended

No unresolved issues remain, and focused regression coverage addresses stream behavior, retries, response alignment, and locale handling.

Review effort: Balanced
Findings: None

What changed in this PR

Reduces memory use and response parsing overhead in JanusGraph’s Elasticsearch bulk client, addressing #5024.

Changes:

  • Sends serialized bulk items through a repeatable entity without assembling another body array.
  • Filters responses to required fields and writes action metadata directly using locale-independent action names.
  • Adds regression coverage and documents the optimization and result-accessor deprecation.
File Description
janusgraph-es/​src/​test/​java/​org/​janusgraph/​diskstorage/​es/​rest/​RestClientRetryTest.java Verifies retries send only failed items.
janusgraph-es/​src/​test/​java/​org/​janusgraph/​diskstorage/​es/​rest/​RestClientBulkRequestsTest.java Tests entity reads, mark/reset, filtering, and action serialization.
janusgraph-es/​src/​test/​java/​org/​janusgraph/​diskstorage/​es/​ElasticsearchIndexTest.java Checks filtered failures retain correct document associations.
janusgraph-es/​src/​main/​java/​org/​janusgraph/​diskstorage/​es/​rest/​RestElasticSearchClient.java Implements copy-free bulk assembly, response filtering, and locale-safe actions.
janusgraph-es/​src/​main/​java/​org/​janusgraph/​diskstorage/​es/​rest/​RestBulkResponse.java Deprecates accessors for the omitted result field.
docs/​changelog.md Documents optimizations, locale correction, and deprecation.

💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.

…responses

A bulk request's body was written into a ByteArrayOutputStream, which
grows by doubling, and then copied into an array, on every attempt, so
a chunk held up to three more copies of its body at once: a bulk of
50,000 documents of 1 KB needed a heap of 320 MB or more. The body is
now an entity which writes, and reads, the items' serialized bytes as
they are, knows its length, can be sent again, and supports mark and
reset for request signers; such a bulk completes in 144 MB.

A bulk request asks Elasticsearch only for each item's status and
error, under the actions JanusGraph sends, and for errors. Every item
keeps its status, so the items still line up with the requests they
answer. The action line of an item is written field by field instead
of as a map through the ObjectMapper, and its action is lower-cased in
the root locale: under a Turkish default locale INDEX became ındex,
which Elasticsearch rejects.

RestBulkItemResponse.getResult() is deprecated, as a filtered response
carries no result.

Fixes JanusGraph#5024

Co-Authored-By: Oleksandr Porunov <alexandr.porunov@gmail.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Signed-off-by: Oleksandr Porunov <alexandr.porunov@gmail.com>
@porunov
porunov force-pushed the es-bulk-request-path branch from b6c84bd to 9fba248 Compare October 5, 2026 12:59
porunov added a commit to porunov/janusgraph that referenced this pull request Oct 5, 2026
toLowerCase() and toUpperCase() without a locale follow the JVM's
default one, in which Turkish and Azerbaijani lower-case I to a dotless
ı and upper-case i to a dotted İ. On such a JVM JanusGraph didn't know
a mapping given in lower case, the type of a JTS geoshape, a time unit
or a backend shorthand in upper case; allowed a key named ID; sent
Elasticsearch geo queries with relations it rejects; named the
Elasticsearch index of a store otherwise than other JVMs; and missed
recovering Solr replicas. Names now convert in the root locale.

Text which JanusGraph lower-cases for text predicates, in memory, for
Elasticsearch queries and for Lucene, is lower-cased one code point at
a time, as Lucene's LowerCaseFilter does in the analyzers of all three
backends: String.toLowerCase differs from it for I in a Turkish locale
and, in any locale, for the dotted capital I and a word-final sigma.

Solr dates are formatted in the root locale. The one conversion left,
the action of an Elasticsearch bulk request, is JanusGraph#5025's.

Fixes JanusGraph#5026

Co-Authored-By: Oleksandr Porunov <alexandr.porunov@gmail.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Signed-off-by: Oleksandr Porunov <alexandr.porunov@gmail.com>
porunov added a commit to porunov/janusgraph that referenced this pull request Oct 5, 2026
toLowerCase() and toUpperCase() without a locale follow the JVM's
default one, in which Turkish and Azerbaijani lower-case I to a dotless
ı and upper-case i to a dotted İ. On such a JVM JanusGraph didn't know
a mapping given in lower case, the type of a JTS geoshape, a time unit
or a backend shorthand in upper case; allowed a key named ID; sent
Elasticsearch geo queries with relations it rejects; named the
Elasticsearch index of a store otherwise than other JVMs; and missed
recovering Solr replicas. Names now convert in the root locale.

Text which JanusGraph lower-cases for text predicates, in memory, for
Elasticsearch queries and for Lucene, is lower-cased one code point at
a time, as Lucene's LowerCaseFilter does in the analyzers of all three
backends: String.toLowerCase differs from it for I in a Turkish locale
and, in any locale, for the dotted capital I and a word-final sigma.

Solr dates are formatted in the root locale. The one conversion left,
the action of an Elasticsearch bulk request, is JanusGraph#5025's.

Fixes JanusGraph#5026

Co-Authored-By: Oleksandr Porunov <alexandr.porunov@gmail.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Signed-off-by: Oleksandr Porunov <alexandr.porunov@gmail.com>
porunov added a commit to porunov/janusgraph that referenced this pull request Oct 5, 2026
toLowerCase() and toUpperCase() without a locale follow the JVM's
default one, in which Turkish and Azerbaijani lower-case I to a dotless
ı and upper-case i to a dotted İ. On such a JVM JanusGraph didn't know
a mapping given in lower case, the type of a JTS geoshape, a time unit
or a backend shorthand in upper case; allowed a key named ID; sent
Elasticsearch geo queries with relations it rejects; named the
Elasticsearch index of a store otherwise than other JVMs; and missed
recovering Solr replicas. Names now convert in the root locale.

Text which JanusGraph lower-cases for text predicates, in memory, for
Elasticsearch queries and for Lucene, is lower-cased one code point at
a time, as Lucene's LowerCaseFilter does in the analyzers of all three
backends: String.toLowerCase differs from it for I in a Turkish locale
and, in any locale, for the dotted capital I and a word-final sigma.

Solr dates are formatted in the root locale. The one conversion left,
the action of an Elasticsearch bulk request, is JanusGraph#5025's.

Fixes JanusGraph#5026

Co-Authored-By: Oleksandr Porunov <alexandr.porunov@gmail.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Signed-off-by: Oleksandr Porunov <alexandr.porunov@gmail.com>
porunov added a commit to porunov/janusgraph that referenced this pull request Oct 5, 2026
toLowerCase() and toUpperCase() without a locale follow the JVM's
default one, in which Turkish and Azerbaijani lower-case I to a dotless
ı and upper-case i to a dotted İ. On such a JVM JanusGraph didn't know
a mapping given in lower case, the type of a JTS geoshape, a time unit
or a backend shorthand in upper case; allowed a key named ID; sent
Elasticsearch geo queries with relations it rejects; named the
Elasticsearch index of a store otherwise than other JVMs; and missed
recovering Solr replicas. Names now convert in the root locale.

Text which JanusGraph lower-cases for text predicates, in memory, for
Elasticsearch queries and for Lucene, is lower-cased one code point at
a time, as Lucene's LowerCaseFilter does in the analyzers of all three
backends: String.toLowerCase differs from it for I in a Turkish locale
and, in any locale, for the dotted capital I and a word-final sigma.

Solr dates are formatted in the root locale. The one conversion left,
the action of an Elasticsearch bulk request, is JanusGraph#5025's.

Fixes JanusGraph#5026

Co-Authored-By: Oleksandr Porunov <alexandr.porunov@gmail.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Signed-off-by: Oleksandr Porunov <alexandr.porunov@gmail.com>

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Elasticsearch bulk requests hold their body several times over and read whole responses for a few fields

2 participants