pitrou commented on code in PR #51038:
URL: https://github.com/apache/arrow/pull/51038#discussion_r4080487231
##########
cpp/src/arrow/json/parser.cc:
##########
@@ -762,102 +744,204 @@ class HandlerBase : public BlockParser,
}
protected:
- template <typename Handler, typename Stream>
- Status DoParse(Handler& handler, Stream&& json, size_t json_size) {
- constexpr auto parse_flags = rj::kParseIterativeFlag |
rj::kParseNanAndInfFlag |
- rj::kParseStopWhenDoneFlag |
- rj::kParseNumbersAsStringsFlag;
-
- rj::Reader reader;
- // ensure that the loop can exit when the block too large.
- for (; num_rows_ < std::numeric_limits<int32_t>::max(); ++num_rows_) {
- auto ok = reader.Parse<parse_flags>(json, handler);
- switch (ok.Code()) {
- case rj::kParseErrorNone:
- // parse the next object
- continue;
- case rj::kParseErrorDocumentEmpty:
- if (json.Tell() < json_size) {
- return ParseError(rj::GetParseError_En(ok.Code()));
- }
- // parsed all objects, finish
- return Status::OK();
- case rj::kParseErrorTermination:
- // handler emitted an error
- return handler.Error();
- default:
- // rj emitted an error
- return ParseError(rj::GetParseError_En(ok.Code()), " in row ",
num_rows_);
+ Status Parse(const std::shared_ptr<Buffer>& json) override {
+ RETURN_NOT_OK(ReserveScalarStorage(json->size()));
+
+ const std::string_view input(reinterpret_cast<const char*>(json->data()),
+ json->size());
+
+ const int64_t input_size = input.size();
+ if (internal::ConsumeJsonWhitespace(input, /*trailing=*/false) ==
input_size) {
+ return Status::OK();
+ }
Review Comment:
As replied in another comment, let's find out if this is really necessary.
##########
cpp/src/arrow/json/parser.cc:
##########
@@ -876,25 +960,60 @@ class HandlerBase : public BlockParser,
return Status::OK();
}
- Status StartArrayImpl() {
- constexpr auto kind = Kind::kArray;
- if (ARROW_PREDICT_FALSE(builder_.kind != kind)) {
+ Status ParseObjectField(std::string_view key, sj::value value) {
+ bool duplicate_keys = false;
+
+ if (SetFieldBuilder(key, &duplicate_keys)) {
+ return ParseValue(value);
+ }
+
+ if (duplicate_keys) {
+ return status_;
+ }
+
+ return HandleUnexpectedField(key, value);
+ }
+
+ template <Kind::type kind>
+ Status AppendScalar(BuilderPtr builder, std::string_view scalar) {
+ if (ARROW_PREDICT_FALSE(builder.kind != kind)) {
return IllegallyChangedTo(kind);
}
- StartNested();
- // append to the list builder in EndArrayImpl
- builder_ = Cast<kind>(builder_)->value_builder();
+ auto index = static_cast<int32_t>(scalar_values_builder_.length());
+ auto value_length = static_cast<int32_t>(scalar.size());
+ RETURN_NOT_OK(Cast<kind>(builder)->Append(index, value_length));
+ RETURN_NOT_OK(scalar_values_builder_.Reserve(1));
+ scalar_values_builder_.UnsafeAppend(scalar);
return Status::OK();
}
- Status EndArrayImpl(rj::SizeType size) {
- EndNested();
- // append to list_builder here
- auto list_builder = Cast<Kind::kArray>(builder_);
- return list_builder->Append(size);
+ /// \brief helper for parsing object fields.
+ ///
+ /// Sets the field builder with the given name, or returns false if
+ /// there is no such field or the field was already specified.
+ bool SetFieldBuilder(std::string_view key, bool* duplicate_keys) {
+ auto parent = Cast<Kind::kObject>(builder_stack_.back());
+ field_index_ = parent->GetFieldIndex(key);
+ if (ARROW_PREDICT_FALSE(field_index_ == -1)) {
+ return false;
+ }
+ if (field_index_ < absent_fields_stack_.TopSize()) {
+ *duplicate_keys = !absent_fields_stack_[field_index_];
+ } else {
+ // When field_index is beyond the range of absent_fields_stack_ we have
a duplicated
+ // field that wasn't declared in schema or previous records.
+ *duplicate_keys = true;
+ }
+ if (*duplicate_keys) {
+ status_ = ParseError("Column(", Path(), ") was specified twice in row ",
num_rows_);
Review Comment:
Why keep this `status_` member? Let's just return a `Result<bool>` from here?
##########
cpp/src/arrow/json/parser.cc:
##########
@@ -762,150 +743,276 @@ class HandlerBase : public BlockParser,
}
protected:
- template <typename Handler, typename Stream>
- Status DoParse(Handler& handler, Stream&& json, size_t json_size) {
- constexpr auto parse_flags = rj::kParseIterativeFlag |
rj::kParseNanAndInfFlag |
- rj::kParseStopWhenDoneFlag |
- rj::kParseNumbersAsStringsFlag;
-
- rj::Reader reader;
- // ensure that the loop can exit when the block too large.
- for (; num_rows_ < std::numeric_limits<int32_t>::max(); ++num_rows_) {
- auto ok = reader.Parse<parse_flags>(json, handler);
- switch (ok.Code()) {
- case rj::kParseErrorNone:
- // parse the next object
- continue;
- case rj::kParseErrorDocumentEmpty:
- if (json.Tell() < json_size) {
- return ParseError(rj::GetParseError_En(ok.Code()));
- }
- // parsed all objects, finish
- return Status::OK();
- case rj::kParseErrorTermination:
- // handler emitted an error
- return handler.Error();
- default:
- // rj emitted an error
- return ParseError(rj::GetParseError_En(ok.Code()), " in row ",
num_rows_);
+ Status DoParse(const std::shared_ptr<Buffer>& json) {
+ RETURN_NOT_OK(ReserveScalarStorage(json->size()));
+
+ const std::string_view input(reinterpret_cast<const char*>(json->data()),
+ json->size());
+
+ const int64_t input_size = input.size();
+ if (internal::ConsumeJsonWhitespace(input, /*trailing=*/false) ==
input_size) {
+ return Status::OK();
+ }
+
+ auto parse = [&](const auto& input) -> Status {
+ ARROW_ASSIGN_OR_RAISE(auto stream,
arrow::internal::ResolveSimdjsonResult(
+ parser_.iterate_many(input),
+ "Failed to create JSON document
stream"));
+
+ for (auto document_result : stream) {
+ ARROW_ASSIGN_OR_RAISE(
+ auto document,
+ arrow::internal::ResolveSimdjsonResult(
+ document_result, "Failed to iterate JSON document stream"));
+
+ if (num_rows_ == std::numeric_limits<int32_t>::max()) {
+ return Status::Invalid("Row count overflowed int32_t");
+ }
+
+ ARROW_ASSIGN_OR_RAISE(
+ auto value,
+ arrow::internal::ResolveSimdjsonResult(
+ document.get_value(), "JSON parse error: Failed to get JSON
value"));
+
+ RETURN_NOT_OK(ParseValue(value));
+
+ ++num_rows_;
+ }
+
+ if (stream.truncated_bytes() != 0) {
+ return ParseError("The document is empty");
}
+
+ return Status::OK();
+ };
+
+ if (json->capacity() - json->size() >=
+ static_cast<int64_t>(simdjson::SIMDJSON_PADDING)) {
+ const auto padded_json = simdjson::padded_string_view(
+ reinterpret_cast<const char*>(json->data()), json->size(),
json->capacity());
+ return parse(padded_json);
}
- return Status::Invalid("Row count overflowed int32_t");
- }
- template <typename Handler>
- Status DoParse(Handler& handler, const std::shared_ptr<Buffer>& json) {
- RETURN_NOT_OK(ReserveScalarStorage(json->size()));
- rj::MemoryStream ms(reinterpret_cast<const char*>(json->data()),
json->size());
- using InputStream = rj::EncodedInputStream<rj::UTF8<>, rj::MemoryStream>;
- return DoParse(handler, InputStream(ms),
static_cast<size_t>(json->size()));
+ simdjson::padded_string padded_json(reinterpret_cast<const
char*>(json->data()),
+ json->size());
+ return parse(padded_json);
}
- /// \defgroup handlerbase-append-methods append non-nested values
- ///
- /// @{
-
template <Kind::type kind>
- Status AppendScalar(BuilderPtr builder, std::string_view scalar) {
- if (ARROW_PREDICT_FALSE(builder.kind != kind)) {
- return IllegallyChangedTo(kind);
+ Status MaybePromoteFromNull() {
+ if (builder_.kind != Kind::kNull) {
+ return Status::OK();
}
- auto index = static_cast<int32_t>(scalar_values_builder_.length());
- auto value_length = static_cast<int32_t>(scalar.size());
- RETURN_NOT_OK(Cast<kind>(builder)->Append(index, value_length));
- RETURN_NOT_OK(scalar_values_builder_.Reserve(1));
- scalar_values_builder_.UnsafeAppend(scalar);
+
+ auto parent = builder_stack_.back();
+
+ if (parent.kind == Kind::kArray) {
+ auto list_builder = Cast<Kind::kArray>(parent);
+ DCHECK_EQ(list_builder->value_builder(), builder_);
+
+ RETURN_NOT_OK(builder_set_.MakeBuilder<kind>(builder_.index, &builder_));
+
+ list_builder = Cast<Kind::kArray>(parent);
+ list_builder->value_builder(builder_);
+ } else {
+ auto struct_builder = Cast<Kind::kObject>(parent);
+ DCHECK_EQ(struct_builder->field_builder(field_index_), builder_);
+
+ RETURN_NOT_OK(builder_set_.MakeBuilder<kind>(builder_.index, &builder_));
+
+ struct_builder = Cast<Kind::kObject>(parent);
+ struct_builder->field_builder(field_index_, builder_);
+ }
+
return Status::OK();
}
- /// @}
+ Status ParseValue(sj::value value) {
+ ARROW_ASSIGN_OR_RAISE(auto type, arrow::internal::ResolveSimdjsonResult(
+ value.type(), "Failed to determine
JSON type"));
- Status StartObjectImpl() {
+ switch (type) {
+ case sj::json_type::null: {
+ ARROW_ASSIGN_OR_RAISE([[maybe_unused]] auto is_null,
+ arrow::internal::ResolveSimdjsonResult(
+ value.is_null(), "Failed to validate JSON
null"));
+ return Null();
+ }
+
+ case sj::json_type::boolean: {
+ RETURN_NOT_OK(MaybePromoteFromNull<Kind::kBoolean>());
+
+ ARROW_ASSIGN_OR_RAISE(auto boolean,
+ arrow::internal::ResolveSimdjsonResult(
+ value.get_bool(), "Failed to get JSON
boolean"));
+ return Bool(boolean);
+ }
+
+ case sj::json_type::string: {
+ RETURN_NOT_OK(MaybePromoteFromNull<Kind::kString>());
+
+ ARROW_ASSIGN_OR_RAISE(auto string,
+ arrow::internal::ResolveSimdjsonResult(
+ value.get_string(), "Failed to get JSON
string"));
+ return String(string);
+ }
+
+ case sj::json_type::number: {
+ RETURN_NOT_OK(MaybePromoteFromNull<Kind::kNumber>());
+ auto raw_number = value.raw_json_token();
+ RETURN_NOT_OK(arrow::internal::ResolveSimdjsonResult(
+ value.get_number(), "Failed to parse JSON number"));
+ raw_number.remove_suffix(
+ internal::ConsumeJsonWhitespace(raw_number, /*trailing=*/true));
+ return RawNumber(raw_number);
+ }
+
+ case sj::json_type::array:
+ RETURN_NOT_OK(MaybePromoteFromNull<Kind::kArray>());
+ return ParseArray(value);
+
+ case sj::json_type::object:
+ RETURN_NOT_OK(MaybePromoteFromNull<Kind::kObject>());
+ return ParseObject(value);
+
+ default:
+ return ParseError("Invalid value");
Review Comment:
Nit: make the message "Invalid JSON value"?
##########
cpp/src/arrow/util/simdjson_internal.cc:
##########
@@ -581,4 +581,18 @@ Status ValidateJsonDocument(simdjson::ondemand::parser&
parser,
return ConsumeJsonValue(value);
}
+/// Returns the number of leading whitespace characters when trailing is false,
+/// or the number of trailing whitespace characters when trailing is true.
+// XXX We could try to SIMD-accelerate this routine.
+int64_t ConsumeJsonWhitespace(std::string_view view, bool trailing) {
+ if (!trailing) {
+ const auto pos = view.find_first_not_of(" \t\r\n");
+ return static_cast<int64_t>(pos == std::string_view::npos ? view.size() :
pos);
Review Comment:
The docstring says "Returns the number of leading whitespace characters",
but this actually returns the view size (not 0) when there is no leading
whitespace. Why?
##########
cpp/src/arrow/json/parser.cc:
##########
@@ -876,25 +960,60 @@ class HandlerBase : public BlockParser,
return Status::OK();
}
- Status StartArrayImpl() {
- constexpr auto kind = Kind::kArray;
- if (ARROW_PREDICT_FALSE(builder_.kind != kind)) {
+ Status ParseObjectField(std::string_view key, sj::value value) {
+ bool duplicate_keys = false;
+
+ if (SetFieldBuilder(key, &duplicate_keys)) {
+ return ParseValue(value);
+ }
+
+ if (duplicate_keys) {
+ return status_;
+ }
+
+ return HandleUnexpectedField(key, value);
+ }
+
+ template <Kind::type kind>
+ Status AppendScalar(BuilderPtr builder, std::string_view scalar) {
+ if (ARROW_PREDICT_FALSE(builder.kind != kind)) {
return IllegallyChangedTo(kind);
}
- StartNested();
- // append to the list builder in EndArrayImpl
- builder_ = Cast<kind>(builder_)->value_builder();
+ auto index = static_cast<int32_t>(scalar_values_builder_.length());
+ auto value_length = static_cast<int32_t>(scalar.size());
+ RETURN_NOT_OK(Cast<kind>(builder)->Append(index, value_length));
+ RETURN_NOT_OK(scalar_values_builder_.Reserve(1));
+ scalar_values_builder_.UnsafeAppend(scalar);
return Status::OK();
}
- Status EndArrayImpl(rj::SizeType size) {
- EndNested();
- // append to list_builder here
- auto list_builder = Cast<Kind::kArray>(builder_);
- return list_builder->Append(size);
+ /// \brief helper for parsing object fields.
+ ///
+ /// Sets the field builder with the given name, or returns false if
+ /// there is no such field or the field was already specified.
+ bool SetFieldBuilder(std::string_view key, bool* duplicate_keys) {
Review Comment:
`bool* duplicate_keys` is only used for returning an error status, let's
remove it?
##########
cpp/src/arrow/json/parser.cc:
##########
@@ -762,102 +744,204 @@ class HandlerBase : public BlockParser,
}
protected:
- template <typename Handler, typename Stream>
- Status DoParse(Handler& handler, Stream&& json, size_t json_size) {
- constexpr auto parse_flags = rj::kParseIterativeFlag |
rj::kParseNanAndInfFlag |
- rj::kParseStopWhenDoneFlag |
- rj::kParseNumbersAsStringsFlag;
-
- rj::Reader reader;
- // ensure that the loop can exit when the block too large.
- for (; num_rows_ < std::numeric_limits<int32_t>::max(); ++num_rows_) {
- auto ok = reader.Parse<parse_flags>(json, handler);
- switch (ok.Code()) {
- case rj::kParseErrorNone:
- // parse the next object
- continue;
- case rj::kParseErrorDocumentEmpty:
- if (json.Tell() < json_size) {
- return ParseError(rj::GetParseError_En(ok.Code()));
- }
- // parsed all objects, finish
- return Status::OK();
- case rj::kParseErrorTermination:
- // handler emitted an error
- return handler.Error();
- default:
- // rj emitted an error
- return ParseError(rj::GetParseError_En(ok.Code()), " in row ",
num_rows_);
+ Status Parse(const std::shared_ptr<Buffer>& json) override {
+ RETURN_NOT_OK(ReserveScalarStorage(json->size()));
+
+ const std::string_view input(reinterpret_cast<const char*>(json->data()),
+ json->size());
+
+ const int64_t input_size = input.size();
+ if (internal::ConsumeJsonWhitespace(input, /*trailing=*/false) ==
input_size) {
+ return Status::OK();
+ }
+
+ auto parse = [&](const auto& input) -> Status {
+ ARROW_ASSIGN_OR_RAISE(auto stream,
arrow::internal::ResolveSimdjsonResult(
+ parser_.iterate_many(input),
+ "Failed to create JSON document
stream"));
+
+ for (auto document_result : stream) {
+ ARROW_ASSIGN_OR_RAISE(
+ auto document,
+ arrow::internal::ResolveSimdjsonResult(
+ document_result, "Failed to iterate JSON document stream"));
+
+ if (num_rows_ == std::numeric_limits<int32_t>::max()) {
+ return Status::Invalid("Row count overflowed int32_t");
+ }
+
+ ARROW_ASSIGN_OR_RAISE(
+ auto value,
+ arrow::internal::ResolveSimdjsonResult(
+ document.get_value(), "JSON parse error: Failed to get JSON
value"));
+
+ RETURN_NOT_OK(ParseValue(value));
+
+ ++num_rows_;
+ }
+
+ if (stream.truncated_bytes() != 0) {
+ return ParseError("The document is empty");
}
+
+ return Status::OK();
+ };
+
+ if (json->capacity() - json->size() >=
+ static_cast<int64_t>(simdjson::SIMDJSON_PADDING)) {
+ const auto padded_json = simdjson::padded_string_view(
+ reinterpret_cast<const char*>(json->data()), json->size(),
json->capacity());
+ return parse(padded_json);
}
- return Status::Invalid("Row count overflowed int32_t");
+
+ // padded_string makes a copy of the input buffer.
Review Comment:
Can you add a TODO pointing to the GH issue you created?
##########
cpp/src/arrow/json/parser.cc:
##########
@@ -762,102 +744,204 @@ class HandlerBase : public BlockParser,
}
protected:
- template <typename Handler, typename Stream>
- Status DoParse(Handler& handler, Stream&& json, size_t json_size) {
- constexpr auto parse_flags = rj::kParseIterativeFlag |
rj::kParseNanAndInfFlag |
- rj::kParseStopWhenDoneFlag |
- rj::kParseNumbersAsStringsFlag;
-
- rj::Reader reader;
- // ensure that the loop can exit when the block too large.
- for (; num_rows_ < std::numeric_limits<int32_t>::max(); ++num_rows_) {
- auto ok = reader.Parse<parse_flags>(json, handler);
- switch (ok.Code()) {
- case rj::kParseErrorNone:
- // parse the next object
- continue;
- case rj::kParseErrorDocumentEmpty:
- if (json.Tell() < json_size) {
- return ParseError(rj::GetParseError_En(ok.Code()));
- }
- // parsed all objects, finish
- return Status::OK();
- case rj::kParseErrorTermination:
- // handler emitted an error
- return handler.Error();
- default:
- // rj emitted an error
- return ParseError(rj::GetParseError_En(ok.Code()), " in row ",
num_rows_);
+ Status Parse(const std::shared_ptr<Buffer>& json) override {
+ RETURN_NOT_OK(ReserveScalarStorage(json->size()));
+
+ const std::string_view input(reinterpret_cast<const char*>(json->data()),
+ json->size());
+
+ const int64_t input_size = input.size();
+ if (internal::ConsumeJsonWhitespace(input, /*trailing=*/false) ==
input_size) {
+ return Status::OK();
+ }
+
+ auto parse = [&](const auto& input) -> Status {
+ ARROW_ASSIGN_OR_RAISE(auto stream,
arrow::internal::ResolveSimdjsonResult(
+ parser_.iterate_many(input),
+ "Failed to create JSON document
stream"));
+
+ for (auto document_result : stream) {
+ ARROW_ASSIGN_OR_RAISE(
+ auto document,
+ arrow::internal::ResolveSimdjsonResult(
+ document_result, "Failed to iterate JSON document stream"));
+
+ if (num_rows_ == std::numeric_limits<int32_t>::max()) {
+ return Status::Invalid("Row count overflowed int32_t");
+ }
+
+ ARROW_ASSIGN_OR_RAISE(
+ auto value,
+ arrow::internal::ResolveSimdjsonResult(
+ document.get_value(), "JSON parse error: Failed to get JSON
value"));
+
+ RETURN_NOT_OK(ParseValue(value));
+
+ ++num_rows_;
+ }
+
+ if (stream.truncated_bytes() != 0) {
+ return ParseError("The document is empty");
}
+
+ return Status::OK();
+ };
+
+ if (json->capacity() - json->size() >=
+ static_cast<int64_t>(simdjson::SIMDJSON_PADDING)) {
+ const auto padded_json = simdjson::padded_string_view(
+ reinterpret_cast<const char*>(json->data()), json->size(),
json->capacity());
+ return parse(padded_json);
}
- return Status::Invalid("Row count overflowed int32_t");
+
+ // padded_string makes a copy of the input buffer.
+ simdjson::padded_string padded_json(reinterpret_cast<const
char*>(json->data()),
+ json->size());
+ return parse(padded_json);
}
- template <typename Handler>
- Status DoParse(Handler& handler, const std::shared_ptr<Buffer>& json) {
- RETURN_NOT_OK(ReserveScalarStorage(json->size()));
- rj::MemoryStream ms(reinterpret_cast<const char*>(json->data()),
json->size());
- using InputStream = rj::EncodedInputStream<rj::UTF8<>, rj::MemoryStream>;
- return DoParse(handler, InputStream(ms),
static_cast<size_t>(json->size()));
+ template <Kind::type kind>
+ Status MaybePromoteFromNull() {
+ if (builder_.kind != Kind::kNull) {
+ return Status::OK();
+ }
+
+ auto parent = builder_stack_.back();
+
+ if (parent.kind == Kind::kArray) {
+ auto list_builder = Cast<Kind::kArray>(parent);
+ DCHECK_EQ(list_builder->value_builder(), builder_);
+
+ RETURN_NOT_OK(builder_set_.MakeBuilder<kind>(builder_.index, &builder_));
+
+ list_builder = Cast<Kind::kArray>(parent);
+ list_builder->value_builder(builder_);
+ } else {
+ auto struct_builder = Cast<Kind::kObject>(parent);
+ DCHECK_EQ(struct_builder->field_builder(field_index_), builder_);
+
+ RETURN_NOT_OK(builder_set_.MakeBuilder<kind>(builder_.index, &builder_));
+
+ struct_builder = Cast<Kind::kObject>(parent);
+ struct_builder->field_builder(field_index_, builder_);
+ }
+
+ return Status::OK();
}
- /// \defgroup handlerbase-append-methods append non-nested values
- ///
- /// @{
+ Status ParseValue(sj::value value) {
+ ARROW_ASSIGN_OR_RAISE(auto type, arrow::internal::ResolveSimdjsonResult(
+ value.type(), "Failed to determine
JSON type"));
- template <Kind::type kind>
- Status AppendScalar(BuilderPtr builder, std::string_view scalar) {
- if (ARROW_PREDICT_FALSE(builder.kind != kind)) {
- return IllegallyChangedTo(kind);
+ switch (type) {
+ case sj::json_type::null: {
+ ARROW_ASSIGN_OR_RAISE([[maybe_unused]] auto is_null,
+ arrow::internal::ResolveSimdjsonResult(
+ value.is_null(), "Failed to validate JSON
null"));
+ return Null();
+ }
+
+ case sj::json_type::boolean: {
+ RETURN_NOT_OK(MaybePromoteFromNull<Kind::kBoolean>());
+
+ ARROW_ASSIGN_OR_RAISE(auto boolean,
+ arrow::internal::ResolveSimdjsonResult(
+ value.get_bool(), "Failed to get JSON
boolean"));
+ return Bool(boolean);
+ }
+
+ case sj::json_type::string: {
+ RETURN_NOT_OK(MaybePromoteFromNull<Kind::kString>());
+
+ ARROW_ASSIGN_OR_RAISE(auto string,
+ arrow::internal::ResolveSimdjsonResult(
+ value.get_string(), "Failed to get JSON
string"));
+ return String(string);
+ }
+
+ case sj::json_type::number: {
+ RETURN_NOT_OK(MaybePromoteFromNull<Kind::kNumber>());
+ auto raw_number = value.raw_json_token();
+ RETURN_NOT_OK(arrow::internal::ResolveSimdjsonResult(
+ value.get_number(), "Failed to parse JSON number"));
+ raw_number.remove_suffix(
+ internal::ConsumeJsonWhitespace(raw_number, /*trailing=*/true));
Review Comment:
Can we add a test exercising this condition, i.e. a test with numbers with
whitespace around them?
##########
cpp/src/arrow/json/parser.cc:
##########
@@ -928,6 +1047,7 @@ class HandlerBase : public BlockParser,
return scalar_values_builder_.ReserveData(size - available_storage);
}
+ UnexpectedFieldBehavior unexpected_field_behavior_;
Status status_;
Review Comment:
Let's remove the `status_` field since we are returning `Status` from all
those APIs?
##########
cpp/src/arrow/util/simdjson_internal.cc:
##########
@@ -581,4 +581,18 @@ Status ValidateJsonDocument(simdjson::ondemand::parser&
parser,
return ConsumeJsonValue(value);
}
+/// Returns the number of leading whitespace characters when trailing is false,
+/// or the number of trailing whitespace characters when trailing is true.
+// XXX We could try to SIMD-accelerate this routine.
Review Comment:
The original comment says this, why modify it?
```c++
// XXX We could try to SIMD-accelerate this routine but it's called only
// once per chunk and also will presumably examine a minimal amount of bytes.
```
##########
cpp/src/arrow/json/parser_test.cc:
##########
@@ -321,5 +321,15 @@ TEST(BlockParser, AdHoc) {
R"([{"c":true, "d": "1991-02-03"}, {"c":false, "d":"2019-04-01"}])"});
}
+TEST(BlockParserWithSchema, ValidateIgnoredFields) {
+ auto options = ParseOptions::Defaults();
+ options.explicit_schema = schema({field("known", int64())});
+ options.unexpected_field_behavior = UnexpectedFieldBehavior::Ignore;
+
+ std::shared_ptr<Array> parsed;
+ ASSERT_RAISES(Invalid,
+ ParseFromString(options, R"({"known": 1, "ignored": [1,]})",
&parsed));
Review Comment:
Thanks. Can you also test the error message?
##########
cpp/src/arrow/json/parser.cc:
##########
@@ -762,102 +744,204 @@ class HandlerBase : public BlockParser,
}
protected:
- template <typename Handler, typename Stream>
- Status DoParse(Handler& handler, Stream&& json, size_t json_size) {
- constexpr auto parse_flags = rj::kParseIterativeFlag |
rj::kParseNanAndInfFlag |
- rj::kParseStopWhenDoneFlag |
- rj::kParseNumbersAsStringsFlag;
-
- rj::Reader reader;
- // ensure that the loop can exit when the block too large.
- for (; num_rows_ < std::numeric_limits<int32_t>::max(); ++num_rows_) {
- auto ok = reader.Parse<parse_flags>(json, handler);
- switch (ok.Code()) {
- case rj::kParseErrorNone:
- // parse the next object
- continue;
- case rj::kParseErrorDocumentEmpty:
- if (json.Tell() < json_size) {
- return ParseError(rj::GetParseError_En(ok.Code()));
- }
- // parsed all objects, finish
- return Status::OK();
- case rj::kParseErrorTermination:
- // handler emitted an error
- return handler.Error();
- default:
- // rj emitted an error
- return ParseError(rj::GetParseError_En(ok.Code()), " in row ",
num_rows_);
+ Status Parse(const std::shared_ptr<Buffer>& json) override {
+ RETURN_NOT_OK(ReserveScalarStorage(json->size()));
+
+ const std::string_view input(reinterpret_cast<const char*>(json->data()),
+ json->size());
+
+ const int64_t input_size = input.size();
+ if (internal::ConsumeJsonWhitespace(input, /*trailing=*/false) ==
input_size) {
+ return Status::OK();
+ }
+
+ auto parse = [&](const auto& input) -> Status {
+ ARROW_ASSIGN_OR_RAISE(auto stream,
arrow::internal::ResolveSimdjsonResult(
+ parser_.iterate_many(input),
+ "Failed to create JSON document
stream"));
+
+ for (auto document_result : stream) {
+ ARROW_ASSIGN_OR_RAISE(
+ auto document,
+ arrow::internal::ResolveSimdjsonResult(
+ document_result, "Failed to iterate JSON document stream"));
+
+ if (num_rows_ == std::numeric_limits<int32_t>::max()) {
+ return Status::Invalid("Row count overflowed int32_t");
+ }
+
+ ARROW_ASSIGN_OR_RAISE(
+ auto value,
+ arrow::internal::ResolveSimdjsonResult(
+ document.get_value(), "JSON parse error: Failed to get JSON
value"));
+
+ RETURN_NOT_OK(ParseValue(value));
+
+ ++num_rows_;
+ }
+
+ if (stream.truncated_bytes() != 0) {
+ return ParseError("The document is empty");
Review Comment:
I don't understand this error message. The doc says:
> If
[truncated_bytes()](https://simdjson.org/api/1.0.0/classsimdjson_1_1_s_i_m_d_j_s_o_n___i_m_p_l_e_m_e_n_t_a_t_i_o_n_1_1ondemand_1_1document__stream.html#acd3fd372dfd30fcce1d6a103c6682d0e)
differs from zero, then the input was truncated maybe because incomplete JSON
documents were found at the end of the stream.
##########
cpp/src/arrow/json/reader_test.cc:
##########
@@ -1032,5 +1030,26 @@ TEST_F(AsyncStreamingReaderTest,
StressSharedIoAndCpuExecutor) {
AssertBatchSequenceEquals(expected.batches, batches);
}
+TEST(ReaderTest, FailOnMalformedNumbers) {
+ auto read_options = ReadOptions::Defaults();
+ auto parse_options = ParseOptions::Defaults();
+
+ const std::vector<std::string> malformed = {
+ R"({"a": 01})",
+ R"({"a": 1.})",
+ };
+
+ // Malformed numbers should be rejected regardless of whether parsing is
threaded.
Review Comment:
```suggestion
// Malformed numbers should be rejected
```
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]