From 8bcb93cb243809589d2d304f3f8ee9428461eff3 Mon Sep 17 00:00:00 2001 From: Alexander Taepper Date: Wed, 23 Sep 2026 16:50:22 +0200 Subject: [PATCH 1/2] [C++] Delimit JSON documents without parsing them --- cpp/src/arrow/json/chunker.cc | 38 +++++++------------------------ cpp/src/arrow/json/reader_test.cc | 9 +------- 2 files changed, 9 insertions(+), 38 deletions(-) diff --git a/cpp/src/arrow/json/chunker.cc b/cpp/src/arrow/json/chunker.cc index f69c45b93867..5536a8670ae7 100644 --- a/cpp/src/arrow/json/chunker.cc +++ b/cpp/src/arrow/json/chunker.cc @@ -42,17 +42,8 @@ int64_t ConsumeWhitespace(std::string_view view) { return static_cast(ws_count); } -Status ConsumeDocument(simdjson::ondemand::document_stream::iterator& it) { - ARROW_ASSIGN_OR_RAISE( - auto document, internal::ResolveSimdjsonResult(*it, "Failed to get JSON document")); - ARROW_ASSIGN_OR_RAISE( - auto value, - internal::ResolveSimdjsonResult(document.get_value(), "Failed to get JSON value")); - return internal::ConsumeJsonValue(value); -} - // A BoundaryFinder implementation that assumes JSON objects can contain raw newlines, -// and uses actual JSON parsing to delimit them. +// and uses the structural indexes computed by simdjson to delimit them. class ParsingBoundaryFinder : public BoundaryFinder { public: explicit ParsingBoundaryFinder(MemoryPool* pool) : pool_(pool) {} @@ -122,39 +113,26 @@ class ParsingBoundaryFinder : public BoundaryFinder { } // Find the first or last JSON object (depending on `find_last`) - // and return the consumed JSON byte length, or 0 if no valid document - // can be parsed. + // and return the consumed JSON byte length, or 0 if no complete document + // can be found. Result FindDocument(simdjson::padded_string_view input, bool find_last) { simdjson::ondemand::document_stream stream; // XXX Should be pass a specific batch_size? // The default value used by simdjson is 1MB, probably enough for most purposes. RETURN_NOT_OK(ToStatus(parser_.iterate_many(input).get(stream))); - auto it = stream.begin(); - if (it == stream.end()) { - // Empty input (only whitespace?) - return 0; - } int64_t consumed_length = 0; - if (!find_last) { - // Parsing the first document only. - if (!ConsumeDocument(it).ok()) { + for (auto it = stream.begin(); it != stream.end(); ++it) { + if (it.error()) { // Could be either a partial document or invalid JSON, we'll let // followup chunker or parser calls decide. - return 0; + break; } // current_index() is the start of the current document; // source() is the complete source span of the current document. consumed_length = it.current_index() + it.source().size(); - } else { - while (it != stream.end()) { - if (!ConsumeDocument(it).ok()) { - // Could be either a partial document or invalid JSON, we'll let - // followup chunker or parser calls decide. - break; - } - consumed_length = it.current_index() + it.source().size(); - ++it; + if (!find_last) { + break; } } if (consumed_length > 0) { diff --git a/cpp/src/arrow/json/reader_test.cc b/cpp/src/arrow/json/reader_test.cc index 7dcd30a0eb0b..549a49480579 100644 --- a/cpp/src/arrow/json/reader_test.cc +++ b/cpp/src/arrow/json/reader_test.cc @@ -743,16 +743,9 @@ TEST_P(StreamingReaderTest, PropagateErrorsNonLinewiseChunker) { read_options_.block_size = 10; parse_options_.newlines_in_values = true; - ASSERT_OK_AND_ASSIGN(reader, MakeReader(bad_first_block)); - AssertReadNext(reader, &batch); - EXPECT_EQ(reader->bytes_processed(), 7); - ASSERT_BATCHES_EQUAL(*RecordBatchFromJSON(test_schema, "[{\"i\":0}]"), *batch); - EXPECT_RAISES_WITH_MESSAGE_THAT(Invalid, ::testing::StartsWith("Invalid: JSON parse error"), - reader->ReadNext(&batch)); - EXPECT_EQ(reader->bytes_processed(), 7); - AssertReadEnd(reader); + MakeReader(bad_first_block)); ASSERT_OK_AND_ASSIGN(reader, MakeReader(bad_middle_blocks)); AssertReadNext(reader, &batch); From 333f4fe75848e07442d792b8136b2dcb1f2003d7 Mon Sep 17 00:00:00 2001 From: Alexander Taepper Date: Wed, 23 Sep 2026 22:41:07 +0200 Subject: [PATCH 2/2] keep the structure more similar to the previous version --- cpp/src/arrow/json/chunker.cc | 21 +++++++++++++++++---- 1 file changed, 17 insertions(+), 4 deletions(-) diff --git a/cpp/src/arrow/json/chunker.cc b/cpp/src/arrow/json/chunker.cc index 5536a8670ae7..4bff06e2ca5c 100644 --- a/cpp/src/arrow/json/chunker.cc +++ b/cpp/src/arrow/json/chunker.cc @@ -120,19 +120,32 @@ class ParsingBoundaryFinder : public BoundaryFinder { // XXX Should be pass a specific batch_size? // The default value used by simdjson is 1MB, probably enough for most purposes. RETURN_NOT_OK(ToStatus(parser_.iterate_many(input).get(stream))); + auto it = stream.begin(); + if (it == stream.end()) { + // Empty input (only whitespace?) + return 0; + } int64_t consumed_length = 0; - for (auto it = stream.begin(); it != stream.end(); ++it) { + if (!find_last) { + // Delimiting the first document only. if (it.error()) { // Could be either a partial document or invalid JSON, we'll let // followup chunker or parser calls decide. - break; + return 0; } // current_index() is the start of the current document; // source() is the complete source span of the current document. consumed_length = it.current_index() + it.source().size(); - if (!find_last) { - break; + } else { + while (it != stream.end()) { + if (it.error()) { + // Could be either a partial document or invalid JSON, we'll let + // followup chunker or parser calls decide. + break; + } + consumed_length = it.current_index() + it.source().size(); + ++it; } } if (consumed_length > 0) {