From 9813fc281851f5b3bdb81671f2c6a75223aecb6a Mon Sep 17 00:00:00 2001 From: Aaditya Srinivasan Date: Sun, 20 Sep 2026 01:02:45 +0530 Subject: [PATCH 1/3] Replace example --- cpp/examples/arrow/CMakeLists.txt | 5 +- cpp/examples/arrow/rapidjson_row_converter.cc | 618 ------------------ cpp/examples/arrow/simdjson_row_converter.cc | 407 ++++++++++++ 3 files changed, 409 insertions(+), 621 deletions(-) delete mode 100644 cpp/examples/arrow/rapidjson_row_converter.cc create mode 100644 cpp/examples/arrow/simdjson_row_converter.cc diff --git a/cpp/examples/arrow/CMakeLists.txt b/cpp/examples/arrow/CMakeLists.txt index 82c075c51df1..a51ee9268148 100644 --- a/cpp/examples/arrow/CMakeLists.txt +++ b/cpp/examples/arrow/CMakeLists.txt @@ -17,9 +17,8 @@ add_arrow_example(row_wise_conversion_example) -if(ARROW_WITH_RAPIDJSON) - add_arrow_example(rapidjson_row_converter EXTRA_LINK_LIBS RapidJSON) - add_arrow_example(from_json_string_example EXTRA_LINK_LIBS RapidJSON) +if(ARROW_JSON) + add_arrow_example(simdjson_row_converter EXTRA_LINK_LIBS arrow::simdjson) endif() if(ARROW_ACERO) diff --git a/cpp/examples/arrow/rapidjson_row_converter.cc b/cpp/examples/arrow/rapidjson_row_converter.cc deleted file mode 100644 index a5fd1536af45..000000000000 --- a/cpp/examples/arrow/rapidjson_row_converter.cc +++ /dev/null @@ -1,618 +0,0 @@ -// Licensed to the Apache Software Foundation (ASF) under one -// or more contributor license agreements. See the NOTICE file -// distributed with this work for additional information -// regarding copyright ownership. The ASF licenses this file -// to you under the Apache License, Version 2.0 (the -// "License"); you may not use this file except in compliance -// with the License. You may obtain a copy of the License at -// -// http://www.apache.org/licenses/LICENSE-2.0 -// -// Unless required by applicable law or agreed to in writing, -// software distributed under the License is distributed on an -// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY -// KIND, either express or implied. See the License for the -// specific language governing permissions and limitations -// under the License. - -#define RAPIDJSON_HAS_STDSTRING 1 - -#include -#include -#include -#include -#include -#include -#include - -#include -#include -#include - -#include -#include -#include - -// Transforming dynamic row data into Arrow data -// When building connectors to other data systems, it's common to receive data in -// row-based structures. While the row_wise_conversion_example.cc shows how to -// handle this conversion for fixed schemas, this example demonstrates how to -// writer converters for arbitrary schemas. -// -// As an example, this conversion is between Arrow and rapidjson::Documents. -// -// We use the following helpers and patterns here: -// * arrow::VisitArrayInline and arrow::VisitTypeInline for implementing a visitor -// pattern with Arrow to handle different array types -// * arrow::enable_if_primitive_ctype to create a template method that handles -// conversion for Arrow types that have corresponding C types (bool, integer, -// float). - -const rapidjson::Value kNullJsonSingleton = rapidjson::Value(); - -/// \brief Builder that holds state for a single conversion. -/// -/// Implements Visit() methods for each type of Arrow Array that set the values -/// of the corresponding fields in each row. -class RowBatchBuilder { - public: - explicit RowBatchBuilder(int64_t num_rows) : field_(nullptr) { - // Reserve all of the space required up-front to avoid unnecessary resizing - rows_.reserve(num_rows); - - for (int64_t i = 0; i < num_rows; ++i) { - rows_.push_back(rapidjson::Document()); - rows_[i].SetObject(); - } - } - - /// \brief Set which field to convert. - void SetField(const arrow::Field* field) { field_ = field; } - - /// \brief Retrieve converted rows from builder. - std::vector Rows() && { return std::move(rows_); } - - // Default implementation - arrow::Status Visit(const arrow::Array& array) { - return arrow::Status::NotImplemented( - "Cannot convert to json document for array of type ", array.type()->ToString()); - } - - // Handles booleans, integers, floats - template - arrow::enable_if_primitive_ctype Visit( - const ArrayType& array) { - assert(static_cast(rows_.size()) == array.length()); - for (int64_t i = 0; i < array.length(); ++i) { - if (!array.IsNull(i)) { - rapidjson::Value str_key(field_->name(), rows_[i].GetAllocator()); - rows_[i].AddMember(str_key, array.Value(i), rows_[i].GetAllocator()); - } - } - return arrow::Status::OK(); - } - - arrow::Status Visit(const arrow::StringArray& array) { - assert(static_cast(rows_.size()) == array.length()); - for (int64_t i = 0; i < array.length(); ++i) { - if (!array.IsNull(i)) { - rapidjson::Value str_key(field_->name(), rows_[i].GetAllocator()); - std::string_view value_view = array.Value(i); - rapidjson::Value value; - value.SetString(value_view.data(), - static_cast(value_view.size()), - rows_[i].GetAllocator()); - rows_[i].AddMember(str_key, value, rows_[i].GetAllocator()); - } - } - return arrow::Status::OK(); - } - - arrow::Status Visit(const arrow::StructArray& array) { - const arrow::StructType* type = array.struct_type(); - - assert(static_cast(rows_.size()) == array.length()); - - RowBatchBuilder child_builder(rows_.size()); - for (int i = 0; i < type->num_fields(); ++i) { - const arrow::Field* child_field = type->field(i).get(); - child_builder.SetField(child_field); - ARROW_RETURN_NOT_OK(arrow::VisitArrayInline(*array.field(i).get(), &child_builder)); - } - std::vector rows = std::move(child_builder).Rows(); - - for (int64_t i = 0; i < array.length(); ++i) { - if (!array.IsNull(i)) { - rapidjson::Value str_key(field_->name(), rows_[i].GetAllocator()); - // Must copy value to new allocator - rapidjson::Value row_val; - row_val.CopyFrom(rows[i], rows_[i].GetAllocator()); - rows_[i].AddMember(str_key, row_val, rows_[i].GetAllocator()); - } - } - return arrow::Status::OK(); - } - - arrow::Status Visit(const arrow::ListArray& array) { - assert(static_cast(rows_.size()) == array.length()); - // First create rows from values - std::shared_ptr values = array.values(); - RowBatchBuilder child_builder(values->length()); - const arrow::Field* value_field = array.list_type()->value_field().get(); - std::string value_field_name = value_field->name(); - child_builder.SetField(value_field); - ARROW_RETURN_NOT_OK(arrow::VisitArrayInline(*values.get(), &child_builder)); - - std::vector rows = std::move(child_builder).Rows(); - - int64_t values_i = 0; - for (int64_t i = 0; i < array.length(); ++i) { - if (array.IsNull(i)) continue; - - rapidjson::Document::AllocatorType& allocator = rows_[i].GetAllocator(); - auto array_len = array.value_length(i); - - rapidjson::Value value; - value.SetArray(); - value.Reserve(array_len, allocator); - - for (int64_t j = 0; j < array_len; ++j) { - rapidjson::Value row_val; - // Must copy value to new allocator - row_val.CopyFrom(rows[values_i][value_field_name], allocator); - value.PushBack(row_val, allocator); - ++values_i; - } - - rapidjson::Value str_key(field_->name(), allocator); - rows_[i].AddMember(str_key, value, allocator); - } - - return arrow::Status::OK(); - } - - private: - const arrow::Field* field_; - std::vector rows_; -}; // RowBatchBuilder - -class ArrowToDocumentConverter { - public: - /// Convert a single batch of Arrow data into Documents - arrow::Result> ConvertToVector( - std::shared_ptr batch) { - RowBatchBuilder builder{batch->num_rows()}; - - for (int i = 0; i < batch->num_columns(); ++i) { - builder.SetField(batch->schema()->field(i).get()); - ARROW_RETURN_NOT_OK(arrow::VisitArrayInline(*batch->column(i).get(), &builder)); - } - - return std::move(builder).Rows(); - } - - /// Convert an Arrow table into an iterator of Documents - arrow::Iterator ConvertToIterator( - std::shared_ptr table, size_t batch_size) { - // Use TableBatchReader to divide table into smaller batches. The batches - // created are zero-copy slices with *at most* `batch_size` rows. - auto batch_reader = std::make_shared(*table); - batch_reader->set_chunksize(batch_size); - - auto read_batch = [this](const std::shared_ptr& batch) - -> arrow::Result> { - ARROW_ASSIGN_OR_RAISE(auto rows, ConvertToVector(batch)); - return arrow::MakeVectorIterator(std::move(rows)); - }; - - auto nested_iter = arrow::MakeMaybeMapIterator( - read_batch, arrow::MakeIteratorFromReader(std::move(batch_reader))); - - return arrow::MakeFlattenIterator(std::move(nested_iter)); - } -}; // ArrowToDocumentConverter - -/// \brief Iterator over rows values of a document for a given field -/// -/// path and array_levels are used to address each field in a JSON document. As -/// an example, consider this JSON document: -/// { -/// "x": 3, // path: ["x"], array_levels: 0 -/// "files": [ // path: ["files"], array_levels: 0 -/// { // path: ["files"], array_levels: 1 -/// "path": "my_str", // path: ["files", "path"], array_levels: 1 -/// "sizes": [ // path: ["files", "size"], array_levels: 1 -/// 20, // path: ["files", "size"], array_levels: 2 -/// 22 -/// ] -/// } -/// ] -/// }, -class DocValuesIterator { - public: - /// \param rows vector of rows - /// \param path field names to enter - /// \param array_levels number of arrays to enter - DocValuesIterator(const std::vector& rows, - std::vector path, int64_t array_levels) - : rows(rows), path(std::move(path)), array_levels(array_levels) {} - - const rapidjson::Value* NextArrayOrRow(const rapidjson::Value* value, size_t* path_i, - int64_t* arr_i) { - while (array_stack.size() > 0) { - ArrayPosition& pos = array_stack.back(); - // Try to get next position in Array - if (pos.index + 1 < pos.array_node->Size()) { - ++pos.index; - value = &(*pos.array_node)[pos.index]; - *path_i = pos.path_index; - *arr_i = array_stack.size(); - return value; - } else { - array_stack.pop_back(); - } - } - ++row_i; - if (row_i < rows.size()) { - value = static_cast(&rows[row_i]); - } else { - value = nullptr; - } - *path_i = 0; - *arr_i = 0; - return value; - } - - arrow::Result Next() { - const rapidjson::Value* value = nullptr; - size_t path_i; - int64_t arr_i; - // Can either start at document or at last array level - if (array_stack.size() > 0) { - auto pos = array_stack.back(); - value = pos.array_node; - path_i = pos.path_index; - arr_i = array_stack.size() - 1; - } - - value = NextArrayOrRow(value, &path_i, &arr_i); - - // Traverse to desired level (with possible backtracking as needed) - while (path_i < path.size() || arr_i < array_levels) { - if (value == nullptr) { - return value; - } else if (value->IsArray() && value->Size() > 0) { - ArrayPosition pos; - pos.array_node = value; - pos.path_index = path_i; - pos.index = 0; - array_stack.push_back(pos); - - value = &(*value)[0]; - ++arr_i; - } else if (value->IsArray()) { - // Empty array means we need to backtrack and go to next array or row - value = NextArrayOrRow(value, &path_i, &arr_i); - } else if (value->HasMember(path[path_i])) { - value = &(*value)[path[path_i]]; - ++path_i; - } else { - return &kNullJsonSingleton; - } - } - - // Return value - return value; - } - - private: - const std::vector& rows; - std::vector path; - int64_t array_levels; - size_t row_i = -1; // index of current row - - // Info about array position for one array level in array stack - struct ArrayPosition { - const rapidjson::Value* array_node; - int64_t path_index; - rapidjson::SizeType index; - }; - std::vector array_stack; -}; - -class JsonValueConverter { - public: - explicit JsonValueConverter(const std::vector& rows) - : rows_(rows), array_levels_(0) {} - - JsonValueConverter(const std::vector& rows, - const std::vector& root_path, int64_t array_levels) - : rows_(rows), root_path_(root_path), array_levels_(array_levels) {} - - /// \brief For field passed in, append corresponding values to builder - arrow::Status Convert(const arrow::Field& field, arrow::ArrayBuilder* builder) { - return Convert(field, field.name(), builder); - } - - /// \brief For field passed in, append corresponding values to builder - arrow::Status Convert(const arrow::Field& field, const std::string& field_name, - arrow::ArrayBuilder* builder) { - field_name_ = field_name; - builder_ = builder; - ARROW_RETURN_NOT_OK(arrow::VisitTypeInline(*field.type().get(), this)); - return arrow::Status::OK(); - } - - // Default implementation - arrow::Status Visit(const arrow::DataType& type) { - return arrow::Status::NotImplemented( - "Cannot convert json value to Arrow array of type ", type.ToString()); - } - - arrow::Status Visit(const arrow::Int64Type& type) { - arrow::Int64Builder* builder = static_cast(builder_); - for (const auto& maybe_value : FieldValues()) { - ARROW_ASSIGN_OR_RAISE(auto value, maybe_value); - if (value->IsNull()) { - ARROW_RETURN_NOT_OK(builder->AppendNull()); - } else { - if (value->IsUint()) { - ARROW_RETURN_NOT_OK(builder->Append(value->GetUint())); - } else if (value->IsInt()) { - ARROW_RETURN_NOT_OK(builder->Append(value->GetInt())); - } else if (value->IsUint64()) { - ARROW_RETURN_NOT_OK(builder->Append(value->GetUint64())); - } else if (value->IsInt64()) { - ARROW_RETURN_NOT_OK(builder->Append(value->GetInt64())); - } else { - return arrow::Status::Invalid("Value is not an integer"); - } - } - } - return arrow::Status::OK(); - } - - arrow::Status Visit(const arrow::DoubleType& type) { - arrow::DoubleBuilder* builder = static_cast(builder_); - for (const auto& maybe_value : FieldValues()) { - ARROW_ASSIGN_OR_RAISE(auto value, maybe_value); - if (value->IsNull()) { - ARROW_RETURN_NOT_OK(builder->AppendNull()); - } else { - ARROW_RETURN_NOT_OK(builder->Append(value->GetDouble())); - } - } - return arrow::Status::OK(); - } - - arrow::Status Visit(const arrow::StringType& type) { - arrow::StringBuilder* builder = static_cast(builder_); - for (const auto& maybe_value : FieldValues()) { - ARROW_ASSIGN_OR_RAISE(auto value, maybe_value); - if (value->IsNull()) { - ARROW_RETURN_NOT_OK(builder->AppendNull()); - } else { - ARROW_RETURN_NOT_OK(builder->Append(value->GetString())); - } - } - return arrow::Status::OK(); - } - - arrow::Status Visit(const arrow::BooleanType& type) { - arrow::BooleanBuilder* builder = static_cast(builder_); - for (const auto& maybe_value : FieldValues()) { - ARROW_ASSIGN_OR_RAISE(auto value, maybe_value); - if (value->IsNull()) { - ARROW_RETURN_NOT_OK(builder->AppendNull()); - } else { - ARROW_RETURN_NOT_OK(builder->Append(value->GetBool())); - } - } - return arrow::Status::OK(); - } - - arrow::Status Visit(const arrow::StructType& type) { - arrow::StructBuilder* builder = static_cast(builder_); - - std::vector child_path(root_path_); - if (field_name_.size() > 0) { - child_path.push_back(field_name_); - } - auto child_converter = JsonValueConverter(rows_, child_path, array_levels_); - - for (int i = 0; i < type.num_fields(); ++i) { - std::shared_ptr child_field = type.field(i); - std::shared_ptr child_builder = builder->child_builder(i); - - ARROW_RETURN_NOT_OK( - child_converter.Convert(*child_field.get(), child_builder.get())); - } - - // Make null bitmap - for (const auto& maybe_value : FieldValues()) { - ARROW_ASSIGN_OR_RAISE(auto value, maybe_value); - ARROW_RETURN_NOT_OK(builder->Append(!value->IsNull())); - } - - return arrow::Status::OK(); - } - - arrow::Status Visit(const arrow::ListType& type) { - arrow::ListBuilder* builder = static_cast(builder_); - - // Values and offsets needs to be interleaved in ListBuilder, so first collect the - // values - std::unique_ptr tmp_value_builder; - ARROW_ASSIGN_OR_RAISE(tmp_value_builder, - arrow::MakeBuilder(builder->value_builder()->type())); - std::vector child_path(root_path_); - child_path.push_back(field_name_); - auto child_converter = JsonValueConverter(rows_, child_path, array_levels_ + 1); - ARROW_RETURN_NOT_OK( - child_converter.Convert(*type.value_field().get(), "", tmp_value_builder.get())); - - std::shared_ptr values_array; - ARROW_RETURN_NOT_OK(tmp_value_builder->Finish(&values_array)); - std::shared_ptr values_data = values_array->data(); - - arrow::ArrayBuilder* value_builder = builder->value_builder(); - int64_t offset = 0; - for (const auto& maybe_value : FieldValues()) { - ARROW_ASSIGN_OR_RAISE(auto value, maybe_value); - ARROW_RETURN_NOT_OK(builder->Append(!value->IsNull())); - if (!value->IsNull() && value->Size() > 0) { - ARROW_RETURN_NOT_OK( - value_builder->AppendArraySlice(*values_data.get(), offset, value->Size())); - offset += value->Size(); - } - } - - return arrow::Status::OK(); - } - - private: - std::string field_name_; - arrow::ArrayBuilder* builder_; - const std::vector& rows_; - std::vector root_path_; - int64_t array_levels_; - - /// Return a flattened iterator over values at nested location - arrow::Iterator FieldValues() { - std::vector path(root_path_); - if (field_name_.size() > 0) { - path.push_back(field_name_); - } - auto iter = DocValuesIterator(rows_, std::move(path), array_levels_); - auto fn = [iter]() mutable -> arrow::Result { - return iter.Next(); - }; - - return arrow::MakeFunctionIterator(fn); - } -}; // JsonValueConverter - -arrow::Result> ConvertToRecordBatch( - const std::vector& rows, std::shared_ptr schema) { - // RecordBatchBuilder will create array builders for us for each field in our - // schema. By passing the number of output rows (`rows.size()`) we can - // pre-allocate the correct size of arrays, except of course in the case of - // string, byte, and list arrays, which have dynamic lengths. - std::unique_ptr batch_builder; - ARROW_ASSIGN_OR_RAISE( - batch_builder, - arrow::RecordBatchBuilder::Make(schema, arrow::default_memory_pool(), rows.size())); - - // Inner converter will take rows and be responsible for appending values - // to provided array builders. - JsonValueConverter converter(rows); - for (int i = 0; i < batch_builder->num_fields(); ++i) { - std::shared_ptr field = schema->field(i); - arrow::ArrayBuilder* builder = batch_builder->GetField(i); - ARROW_RETURN_NOT_OK(converter.Convert(*field.get(), builder)); - } - - std::shared_ptr batch; - ARROW_ASSIGN_OR_RAISE(batch, batch_builder->Flush()); - - // Use RecordBatch::ValidateFull() to make sure arrays were correctly constructed. - ARROW_RETURN_NOT_OK(batch->ValidateFull()); - return batch; -} // ConvertToRecordBatch - -arrow::Status DoRowConversion(int32_t num_rows, int32_t batch_size) { - //(Doc section: Convert to Arrow) - // Write JSON records - std::vector json_records = { - R"({"pk": 1, "date_created": "2020-10-01", "data": {"deleted": true, "metrics": [{"key": "x", "value": 1}]}})", - R"({"pk": 2, "date_created": "2020-10-03", "data": {"deleted": false, "metrics": []}})", - R"({"pk": 3, "date_created": "2020-10-05", "data": {"deleted": false, "metrics": [{"key": "x", "value": 33}, {"key": "x", "value": 42}]}})"}; - - std::vector records; - records.reserve(num_rows); - for (int32_t i = 0; i < num_rows; ++i) { - rapidjson::Document document; - document.Parse(json_records[i % json_records.size()]); - records.push_back(std::move(document)); - } - - for (const rapidjson::Document& doc : records) { - rapidjson::StringBuffer sb; - rapidjson::Writer writer(sb); - doc.Accept(writer); - std::cout << sb.GetString() << std::endl; - } - auto tags_schema = arrow::list(arrow::struct_({ - arrow::field("key", arrow::utf8()), - arrow::field("value", arrow::int64()), - })); - auto schema = arrow::schema( - {arrow::field("pk", arrow::int64()), arrow::field("date_created", arrow::utf8()), - arrow::field("data", arrow::struct_({arrow::field("deleted", arrow::boolean()), - arrow::field("metrics", tags_schema)}))}); - - // Convert records into a table - ARROW_ASSIGN_OR_RAISE(std::shared_ptr batch, - ConvertToRecordBatch(records, schema)); - - ARROW_ASSIGN_OR_RAISE(std::shared_ptr table, - arrow::Table::FromRecordBatches({batch})); - - // Print table - std::cout << table->ToString() << std::endl; - ARROW_RETURN_NOT_OK(table->ValidateFull()); - //(Doc section: Convert to Arrow) - - //(Doc section: Convert to Rows) - // Create converter - ArrowToDocumentConverter to_doc_converter; - - // Convert table into document (row) iterator - arrow::Iterator document_iter = - to_doc_converter.ConvertToIterator(table, batch_size); - - // Print each row - for (arrow::Result doc_result : document_iter) { - ARROW_ASSIGN_OR_RAISE(rapidjson::Document doc, std::move(doc_result)); - - assert(doc.HasMember("pk")); - assert(doc["pk"].IsInt64()); - assert(doc.HasMember("date_created")); - assert(doc["date_created"].IsString()); - assert(doc.HasMember("data")); - assert(doc["data"].IsObject()); - assert(doc["data"].HasMember("deleted")); - assert(doc["data"]["deleted"].IsBool()); - assert(doc["data"].HasMember("metrics")); - assert(doc["data"]["metrics"].IsArray()); - if (doc["data"]["metrics"].Size() > 0) { - auto metric = &doc["data"]["metrics"][0]; - assert(metric->IsObject()); - assert(metric->HasMember("key")); - assert((*metric)["key"].IsString()); - assert(metric->HasMember("value")); - assert((*metric)["value"].IsInt64()); - } - - rapidjson::StringBuffer sb; - rapidjson::Writer writer(sb); - doc.Accept(writer); - std::cout << sb.GetString() << std::endl; - } - //(Doc section: Convert to Rows) - - return arrow::Status::OK(); -} - -int main(int argc, char** argv) { - int32_t num_rows = argc > 1 ? std::atoi(argv[1]) : 100; - int32_t batch_size = argc > 2 ? std::atoi(argv[2]) : 100; - - arrow::Status status = DoRowConversion(num_rows, batch_size); - - if (!status.ok()) { - std::cerr << "Error occurred: " << status.message() << std::endl; - return EXIT_FAILURE; - } - return EXIT_SUCCESS; -} diff --git a/cpp/examples/arrow/simdjson_row_converter.cc b/cpp/examples/arrow/simdjson_row_converter.cc new file mode 100644 index 000000000000..26cc6727f98f --- /dev/null +++ b/cpp/examples/arrow/simdjson_row_converter.cc @@ -0,0 +1,407 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +#include "arrow/api.h" +#include "arrow/record_batch.h" +#include "arrow/result.h" +#include "arrow/table_builder.h" +#include "arrow/util/iterator.h" +#include "arrow/util/logging.h" +#include "arrow/util/simdjson_internal.h" + +#include + +#include +#include +#include +#include + +// Transforming dynamic row data into Arrow data +// When building connectors to other data systems, it's common to receive data in +// row-based structures. While the row_wise_conversion_example.cc shows how to +// handle this conversion for fixed schemas, this example demonstrates how to +// convert between row-based JSON data and Arrow data for arbitrary schemas. +// +// As an example, this conversion is between JSON strings and Arrow tables. +// +// We use the following helpers and patterns here: +// * arrow::internal::JsonWriter for writing JSON values +// * arrow::internal::ParseJsonObject and related helpers for parsing JSON +// * arrow::RecordBatchBuilder for constructing Arrow arrays from row data +// * arrow::TableBatchReader and Arrow iterators for converting Arrow tables back +// into row-based JSON data + +namespace arrow { + +namespace { + +namespace sj = simdjson; + +// Append a JSON value to an Arrow builder according to the expected Arrow type. +// This example only handles the types used by the example schema; extend this +// switch when adapting it to schemas with additional Arrow types. +Status AppendJsonValue(const sj::dom::element& value, + const std::shared_ptr& type, ArrayBuilder* builder); + +Status AppendJsonStruct(const sj::dom::element& value, const StructType& type, + StructBuilder* builder) { + if (value.is_null()) { + for (int i = 0; i < type.num_fields(); ++i) { + ARROW_RETURN_NOT_OK(builder->child_builder(i)->AppendNull()); + } + return builder->AppendNull(); + } + + if (!value.is_object()) { + return Status::TypeError("Expected JSON object for struct"); + } + + ARROW_ASSIGN_OR_RAISE( + auto object, + internal::ResolveSimdjsonResult(value.get_object(), "Failed to get JSON object")); + + for (int i = 0; i < type.num_fields(); ++i) { + const auto& field = type.field(i); + + ARROW_ASSIGN_OR_RAISE(auto child, + internal::GetOptionalJsonField(object, field->name())); + + if (!child.has_value()) { + ARROW_RETURN_NOT_OK(builder->child_builder(i)->AppendNull()); + } else { + ARROW_RETURN_NOT_OK( + AppendJsonValue(*child, field->type(), builder->child_builder(i).get())); + } + } + + return builder->Append(); +} + +Status AppendJsonList(const sj::dom::element& value, const ListType& type, + ListBuilder* builder) { + if (value.is_null()) { + return builder->AppendNull(); + } + + ARROW_ASSIGN_OR_RAISE(auto array, internal::GetJsonArray(value, "JSON value")); + + ARROW_RETURN_NOT_OK(builder->Append()); + + for (auto element : array) { + ARROW_RETURN_NOT_OK( + AppendJsonValue(element, type.value_field()->type(), builder->value_builder())); + } + + return Status::OK(); +} + +Status AppendJsonValue(const sj::dom::element& value, + const std::shared_ptr& type, ArrayBuilder* builder) { + if (value.is_null()) { + return builder->AppendNull(); + } + + switch (type->id()) { + case Type::INT64: { + ARROW_ASSIGN_OR_RAISE(auto number, + internal::GetJsonInt(value, "JSON value", "integers")); + return static_cast(builder)->Append(number); + } + + case Type::DOUBLE: { + ARROW_ASSIGN_OR_RAISE(auto number, + internal::ResolveSimdjsonResult(value.get_double(), + "Failed to get JSON double")); + return static_cast(builder)->Append(number); + } + + case Type::STRING: { + ARROW_ASSIGN_OR_RAISE(auto string, + internal::ResolveSimdjsonResult(value.get_string(), + "Failed to get JSON string")); + return static_cast(builder)->Append(string); + } + + case Type::BOOL: { + ARROW_ASSIGN_OR_RAISE( + auto boolean, internal::ResolveSimdjsonResult(value.get_bool(), + "Failed to get JSON boolean")); + return static_cast(builder)->Append(boolean); + } + + case Type::STRUCT: + return AppendJsonStruct(value, *static_cast(type.get()), + static_cast(builder)); + + case Type::LIST: + return AppendJsonList(value, *static_cast(type.get()), + static_cast(builder)); + + default: + return Status::NotImplemented("Cannot convert JSON value to Arrow array of type ", + type->ToString()); + } +} // AppendJsonValue + +// RecordBatchBuilder will create array builders for us for each field in our +// schema. By passing the number of output rows (`rows.size()`), we pre-allocate +// the correct size of arrays, except of course in the case of string and list +// arrays, which have dynamic lengths. +Result> ConvertToRecordBatch( + const std::vector& rows, const std::shared_ptr& schema) { + std::unique_ptr batch_builder; + + ARROW_ASSIGN_OR_RAISE(batch_builder, RecordBatchBuilder::Make( + schema, default_memory_pool(), rows.size())); + + sj::dom::parser parser; + + // Parse each row and append its values to the corresponding Arrow builders. + for (const auto& json : rows) { + ARROW_ASSIGN_OR_RAISE(auto object, internal::ParseJsonObject(parser, json)); + + for (int i = 0; i < schema->num_fields(); ++i) { + const auto& field = schema->field(i); + auto builder = batch_builder->GetField(i); + + ARROW_ASSIGN_OR_RAISE(auto value, + internal::GetOptionalJsonField(object, field->name())); + + if (!value.has_value()) { + ARROW_RETURN_NOT_OK(builder->AppendNull()); + } else { + ARROW_RETURN_NOT_OK(AppendJsonValue(*value, field->type(), builder)); + } + } + } + + ARROW_ASSIGN_OR_RAISE(std::shared_ptr batch, batch_builder->Flush()); + + ARROW_RETURN_NOT_OK(batch->ValidateFull()); + return batch; +} + +// Write an Arrow value as JSON according to its Arrow type. +// This example only handles the types used by the example schema; extend this +// switch when adapting it to schemas with additional Arrow types. +Status WriteJsonValue(const Array& array, int64_t index, + const std::shared_ptr& type, + internal::JsonWriter* writer); + +Status WriteJsonStruct(const StructArray& array, int64_t index, const StructType& type, + internal::JsonWriter* writer) { + writer->StartObject(); + + for (int i = 0; i < type.num_fields(); ++i) { + const auto& field = type.field(i); + const auto& child = array.field(i); + + writer->Key(field->name()); + + if (child->IsNull(index)) { + writer->Null(); + } else { + ARROW_RETURN_NOT_OK(WriteJsonValue(*child, index, field->type(), writer)); + } + } + + writer->EndObject(); + return Status::OK(); +} + +Status WriteJsonList(const ListArray& array, int64_t index, const ListType& type, + internal::JsonWriter* writer) { + writer->StartArray(); + + const int64_t offset = array.value_offset(index); + const int64_t length = array.value_length(index); + const auto& values = *array.values(); + + for (int64_t i = 0; i < length; ++i) { + ARROW_RETURN_NOT_OK( + WriteJsonValue(values, offset + i, type.value_field()->type(), writer)); + } + + writer->EndArray(); + return Status::OK(); +} + +Status WriteJsonValue(const Array& array, int64_t index, + const std::shared_ptr& type, + internal::JsonWriter* writer) { + if (array.IsNull(index)) { + writer->Null(); + return Status::OK(); + } + + switch (type->id()) { + case Type::INT64: + writer->Int64(static_cast(array).Value(index)); + return Status::OK(); + + case Type::DOUBLE: + writer->Double(static_cast(array).Value(index)); + return Status::OK(); + + case Type::STRING: + writer->String(static_cast(array).GetView(index)); + return Status::OK(); + + case Type::BOOL: + writer->Bool(static_cast(array).Value(index)); + return Status::OK(); + + case Type::STRUCT: + return WriteJsonStruct(static_cast(array), index, + *static_cast(type.get()), writer); + + case Type::LIST: + return WriteJsonList(static_cast(array), index, + *static_cast(type.get()), writer); + + default: + return Status::NotImplemented("Cannot convert Arrow array of type ", + type->ToString(), " to JSON"); + } +} + +// Convert a single row of an Arrow record batch into a JSON object. +Result ConvertRowToJson(const RecordBatch& batch, int64_t row) { + internal::JsonWriter writer; + + writer.StartObject(); + + for (int i = 0; i < batch.num_columns(); ++i) { + const auto& field = batch.schema()->field(i); + const auto& column = batch.column(i); + + writer.Key(field->name()); + ARROW_RETURN_NOT_OK(WriteJsonValue(*column, row, field->type(), &writer)); + } + + writer.EndObject(); + + ARROW_ASSIGN_OR_RAISE(auto json, writer.GetString()); + + return std::string(json); +} + +// Convert a single batch of Arrow data into JSON rows. +Result>> ConvertToVector( + const std::shared_ptr& batch) { + std::vector> rows; + rows.reserve(batch->num_rows()); + + for (int64_t i = 0; i < batch->num_rows(); ++i) { + ARROW_ASSIGN_OR_RAISE(auto row, ConvertRowToJson(*batch, i)); + rows.push_back(std::make_shared(std::move(row))); + } + + return rows; +} + +// Convert an Arrow table into an iterator of JSON rows. +class ArrowToJsonConverter { + public: + Iterator> ConvertToIterator(std::shared_ptr table, + size_t batch_size) { + // Use TableBatchReader to divide the table into smaller batches. The batches + // created are zero-copy slices with *at most* `batch_size` rows. + auto batch_reader = std::make_shared(*table); + batch_reader->set_chunksize(batch_size); + + auto read_batch = [](const std::shared_ptr& batch) + -> Result>> { + ARROW_ASSIGN_OR_RAISE(auto rows, ConvertToVector(batch)); + return MakeVectorIterator(std::move(rows)); + }; + + auto nested_iter = + MakeMaybeMapIterator(read_batch, MakeIteratorFromReader(std::move(batch_reader))); + + return MakeFlattenIterator(std::move(nested_iter)); + } +}; + +Status DoRowConversion(int32_t num_rows, int32_t batch_size) { + //(Doc section: Convert to Arrow) + // Write JSON records + std::vector json_records = { + R"({"pk": 1, "date_created": "2020-10-01", "data": {"deleted": true, "metrics": [{"key": "x", "value": 1}]}})", + R"({"pk": 2, "date_created": "2020-10-03", "data": {"deleted": false, "metrics": []}})", + R"({"pk": 3, "date_created": "2020-10-05", "data": {"deleted": false, "metrics": [{"key": "x", "value": 33}, {"key": "x", "value": 42}]}})"}; + + std::vector records; + records.reserve(num_rows); + + for (int32_t i = 0; i < num_rows; ++i) { + records.push_back(json_records[i % json_records.size()]); + } + + for (const auto& json : records) { + std::cout << json << std::endl; + } + + auto tags_schema = list(struct_({ + field("key", utf8()), + field("value", int64()), + })); + + auto table_schema = schema({field("pk", int64()), field("date_created", utf8()), + field("data", struct_({field("deleted", boolean()), + field("metrics", tags_schema)}))}); + + // Convert records into a table + ARROW_ASSIGN_OR_RAISE(std::shared_ptr batch, + ConvertToRecordBatch(records, table_schema)); + + ARROW_ASSIGN_OR_RAISE(std::shared_ptr
table, Table::FromRecordBatches({batch})); + + std::cout << table->ToString() << std::endl; + ARROW_RETURN_NOT_OK(table->ValidateFull()); + + //(Doc section: Convert to Rows) + ArrowToJsonConverter to_json_converter; + + auto json_iter = to_json_converter.ConvertToIterator(table, batch_size); + + for (Result> json_result : json_iter) { + ARROW_ASSIGN_OR_RAISE(auto json, std::move(json_result)); + std::cout << *json << std::endl; + } + //(Doc section: Convert to Rows) + + return Status::OK(); +} + +} // namespace + +} // namespace arrow + +int main(int argc, char** argv) { + int32_t num_rows = argc > 1 ? std::atoi(argv[1]) : 100; + int32_t batch_size = argc > 2 ? std::atoi(argv[2]) : 100; + + arrow::Status status = arrow::DoRowConversion(num_rows, batch_size); + + if (!status.ok()) { + std::cerr << "Error occurred: " << status.message() << std::endl; + return EXIT_FAILURE; + } + + return EXIT_SUCCESS; +} From d0a9b25a830e20ecff972899d9b57ee63f6c6597 Mon Sep 17 00:00:00 2001 From: Aaditya Srinivasan Date: Sun, 20 Sep 2026 01:53:30 +0530 Subject: [PATCH 2/3] Update rst doc --- .../cpp/examples/row_columnar_conversion.rst | 42 ++++++++----------- 1 file changed, 17 insertions(+), 25 deletions(-) diff --git a/docs/source/cpp/examples/row_columnar_conversion.rst b/docs/source/cpp/examples/row_columnar_conversion.rst index 8c95e3bedd94..69b3368a6ede 100644 --- a/docs/source/cpp/examples/row_columnar_conversion.rst +++ b/docs/source/cpp/examples/row_columnar_conversion.rst @@ -47,9 +47,9 @@ provides several utilities: * :class:`arrow::TableBatchReader`: read a table in a batch at a time, with each batch being a zero-copy slice. -The following example shows how to implement conversion between ``rapidjson::Document`` +The following example shows how to implement conversion between JSON rows and Arrow objects. You can read the full code example at -https://github.com/apache/arrow/blob/main/cpp/examples/arrow/rapidjson_row_converter.cc +https://github.com/apache/arrow/blob/main/cpp/examples/arrow/simdjson_row_converter.cc Writing conversions to Arrow ~~~~~~~~~~~~~~~~~~~~~~~~~~~~ @@ -63,7 +63,7 @@ check the first N rows to infer a schema if there is none already available. At the top level, we define a function ``ConvertToRecordBatch``: -.. literalinclude:: ../../../../cpp/examples/arrow/rapidjson_row_converter.cc +.. literalinclude:: ../../../../cpp/examples/arrow/simdjson_row_converter.cc :language: cpp :start-at: arrow::Result> ConvertToRecordBatch( :end-at: } // ConvertToRecordBatch @@ -72,28 +72,20 @@ At the top level, we define a function ``ConvertToRecordBatch``: First we use :class:`arrow::RecordBatchBuilder`, which conveniently creates builders for each field in the schema. Then we iterate over the fields of the schema, get -the builder, and call ``Convert()`` on our ``JsonValueConverter`` (to be discussed -next). At the end, we call ``batch->ValidateFull()``, which checks the integrity +the builder, and append the corresponding JSON value using ``AppendJsonValue``. +At the end, we call ``batch->ValidateFull()``, which checks the integrity of our arrays to make sure the conversion was performed correctly, which is useful for debugging new conversion implementations. -One level down, the ``JsonValueConverter`` is responsible for appending row values -for the provided field to a provided array builder. In order to specialize logic -for each data type, it implements ``Visit`` methods and calls :func:`arrow::VisitTypeInline`. -(See more about type visitors in :ref:`cpp-visitor-pattern`.) +The ``AppendJsonValue`` function is responsible for appending a JSON value +to an Arrow array builder according to its Arrow data type. It handles the +types used by the example schema and recursively processes nested structs and +lists. -At the end of that class is the private method ``FieldValues()``, which returns -an iterator of the column values for the current field across the rows. In -row-based structures that are flat (such as a vector of values) this may be -trivial to implement. But if the schema is nested, as in the case of JSON documents, -a special iterator is needed to navigate the levels of nesting. See the -`full example `_ -for the implementation details of ``DocValuesIterator``. - -.. literalinclude:: ../../../../cpp/examples/arrow/rapidjson_row_converter.cc +.. literalinclude:: ../../../../cpp/examples/arrow/simdjson_row_converter.cc :language: cpp - :start-at: class JsonValueConverter - :end-at: }; // JsonValueConverter + :start-at: arrow::Status AppendJsonValue + :end-at: } // AppendJsonValue :linenos: :lineno-match: @@ -104,7 +96,7 @@ To convert into rows *from* Arrow record batches, we'll process the table in smaller batches, visiting each field of the batch and filling the output rows column-by-column. -At the top-level, we define ``ArrowToDocumentConverter`` that provides the API +At the top-level, we define ``ArrowToJsonConverter`` that provides the API for converting Arrow batches and tables to rows. In many cases, it's more optimal to perform conversions to rows in smaller batches, rather than doing the entire table at once. So we define one ``ConvertToVector`` method to convert a single @@ -113,10 +105,10 @@ to iterate over slices of a table. This returns Arrow's iterator type (:class:`arrow::Iterator`) so rows could then be processed either one-at-a-time or be collected into a container. -.. literalinclude:: ../../../../cpp/examples/arrow/rapidjson_row_converter.cc +.. literalinclude:: ../../../../cpp/examples/arrow/simdjson_row_converter.cc :language: cpp - :start-at: class ArrowToDocumentConverter - :end-at: }; // ArrowToDocumentConverter + :start-at: class ArrowToJsonConverter + :end-at: }; // ArrowToJsonConverter :linenos: :lineno-match: @@ -126,7 +118,7 @@ write a template method for array types that have primitive C equivalents (booleans, integers, and floats) using ``arrow::enable_if_primitive_ctype``. See :ref:`type-traits` for other type predicates. -.. literalinclude:: ../../../../cpp/examples/arrow/rapidjson_row_converter.cc +.. literalinclude:: ../../../../cpp/examples/arrow/simdjson_row_converter.cc :language: cpp :start-at: class RowBatchBuilder :end-at: }; // RowBatchBuilder From ebf072fae415f27ac09ba90b51dfc79e41281cba Mon Sep 17 00:00:00 2001 From: Aaditya Srinivasan Date: Mon, 21 Sep 2026 00:14:37 +0530 Subject: [PATCH 3/3] Minor fixes --- cpp/examples/arrow/CMakeLists.txt | 1 + cpp/examples/arrow/simdjson_row_converter.cc | 6 +++--- .../cpp/examples/row_columnar_conversion.rst | 17 ++++++++--------- 3 files changed, 12 insertions(+), 12 deletions(-) diff --git a/cpp/examples/arrow/CMakeLists.txt b/cpp/examples/arrow/CMakeLists.txt index a51ee9268148..85989011854b 100644 --- a/cpp/examples/arrow/CMakeLists.txt +++ b/cpp/examples/arrow/CMakeLists.txt @@ -19,6 +19,7 @@ add_arrow_example(row_wise_conversion_example) if(ARROW_JSON) add_arrow_example(simdjson_row_converter EXTRA_LINK_LIBS arrow::simdjson) + add_arrow_example(from_json_string_example EXTRA_LINK_LIBS arrow::simdjson) endif() if(ARROW_ACERO) diff --git a/cpp/examples/arrow/simdjson_row_converter.cc b/cpp/examples/arrow/simdjson_row_converter.cc index 26cc6727f98f..50114042e38f 100644 --- a/cpp/examples/arrow/simdjson_row_converter.cc +++ b/cpp/examples/arrow/simdjson_row_converter.cc @@ -193,7 +193,7 @@ Result> ConvertToRecordBatch( ARROW_RETURN_NOT_OK(batch->ValidateFull()); return batch; -} +} // ConvertToRecordBatch // Write an Arrow value as JSON according to its Arrow type. // This example only handles the types used by the example schema; extend this @@ -277,7 +277,7 @@ Status WriteJsonValue(const Array& array, int64_t index, return Status::NotImplemented("Cannot convert Arrow array of type ", type->ToString(), " to JSON"); } -} +} // WriteJsonValue // Convert a single row of an Arrow record batch into a JSON object. Result ConvertRowToJson(const RecordBatch& batch, int64_t row) { @@ -335,7 +335,7 @@ class ArrowToJsonConverter { return MakeFlattenIterator(std::move(nested_iter)); } -}; +}; // ArrowToJsonConverter Status DoRowConversion(int32_t num_rows, int32_t batch_size) { //(Doc section: Convert to Arrow) diff --git a/docs/source/cpp/examples/row_columnar_conversion.rst b/docs/source/cpp/examples/row_columnar_conversion.rst index 69b3368a6ede..dea1de3f0589 100644 --- a/docs/source/cpp/examples/row_columnar_conversion.rst +++ b/docs/source/cpp/examples/row_columnar_conversion.rst @@ -65,7 +65,7 @@ At the top level, we define a function ``ConvertToRecordBatch``: .. literalinclude:: ../../../../cpp/examples/arrow/simdjson_row_converter.cc :language: cpp - :start-at: arrow::Result> ConvertToRecordBatch( + :start-at: Result> ConvertToRecordBatch( :end-at: } // ConvertToRecordBatch :linenos: :lineno-match: @@ -84,7 +84,7 @@ lists. .. literalinclude:: ../../../../cpp/examples/arrow/simdjson_row_converter.cc :language: cpp - :start-at: arrow::Status AppendJsonValue + :start-at: Status AppendJsonValue :end-at: } // AppendJsonValue :linenos: :lineno-match: @@ -112,15 +112,14 @@ or be collected into a container. :linenos: :lineno-match: -One level down, the output rows are filled in by ``RowBatchBuilder``. -The ``RowBatchBuilder`` implements ``Visit()`` methods, but to save on code we -write a template method for array types that have primitive C equivalents -(booleans, integers, and floats) using ``arrow::enable_if_primitive_ctype``. -See :ref:`type-traits` for other type predicates. +One level down, the ``WriteJsonValue`` function is responsible for writing +an Arrow value as JSON according to its Arrow data type. It handles the +types used by the example schema and recursively processes nested structs +and lists. .. literalinclude:: ../../../../cpp/examples/arrow/simdjson_row_converter.cc :language: cpp - :start-at: class RowBatchBuilder - :end-at: }; // RowBatchBuilder + :start-at: Status WriteJsonValue + :end-at: } // WriteJsonValue :linenos: :lineno-match: