Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
21 changes: 6 additions & 15 deletions cpp/src/arrow/json/chunker.cc
Original file line number Diff line number Diff line change
Expand Up @@ -42,17 +42,8 @@ int64_t ConsumeWhitespace(std::string_view view) {
return static_cast<int64_t>(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) {}
Expand Down Expand Up @@ -122,8 +113,8 @@ 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<size_t> FindDocument(simdjson::padded_string_view input, bool find_last) {
simdjson::ondemand::document_stream stream;
// XXX Should be pass a specific batch_size?
Expand All @@ -137,8 +128,8 @@ class ParsingBoundaryFinder : public BoundaryFinder {

int64_t consumed_length = 0;
if (!find_last) {
// Parsing the first document only.
if (!ConsumeDocument(it).ok()) {
// 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.
return 0;
Expand All @@ -148,7 +139,7 @@ class ParsingBoundaryFinder : public BoundaryFinder {
consumed_length = it.current_index() + it.source().size();
} else {
while (it != stream.end()) {
if (!ConsumeDocument(it).ok()) {
if (it.error()) {
// Could be either a partial document or invalid JSON, we'll let
// followup chunker or parser calls decide.
break;
Expand Down
9 changes: 1 addition & 8 deletions cpp/src/arrow/json/reader_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
Loading