Skip to content
Draft
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
1 change: 1 addition & 0 deletions ci/docker/ubuntu-24.04-cpp.dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -221,4 +221,5 @@ ENV absl_SOURCE=BUNDLED \
PARQUET_BUILD_EXECUTABLES=ON \
PATH=/usr/lib/ccache/:$PATH \
PYTHON=python3 \
simdjson_SOURCE=BUNDLED \
xsimd_SOURCE=BUNDLED
4 changes: 3 additions & 1 deletion cpp/src/arrow/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -733,7 +733,8 @@ if(ARROW_BUILD_INTEGRATION OR ARROW_BUILD_TESTS)
arrow_add_object_library(ARROW_INTEGRATION integration/json_integration.cc
integration/json_internal.cc)
foreach(ARROW_INTEGRATION_TARGET ${ARROW_INTEGRATION_TARGETS})
target_link_libraries(${ARROW_INTEGRATION_TARGET} PRIVATE RapidJSON)
target_link_libraries(${ARROW_INTEGRATION_TARGET} PRIVATE RapidJSON
simdjson::simdjson)
endforeach()
else()
set(ARROW_INTEGRATION_TARGET_SHARED)
Expand Down Expand Up @@ -1041,6 +1042,7 @@ if(ARROW_JSON)
json/from_string.cc
json/object_parser.cc
json/object_writer.cc
json/json_writer_internal.cc
json/parser.cc
json/reader.cc)
foreach(ARROW_JSON_TARGET ${ARROW_JSON_TARGETS})
Expand Down
7 changes: 6 additions & 1 deletion cpp/src/arrow/integration/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -21,12 +21,17 @@ arrow_install_all_headers("arrow/integration")
# - an executable that can be called to answer integration test requests
# - a self-(unit)test for the C++ side of integration testing
if(ARROW_BUILD_TESTS)
add_arrow_test(json_integration_test EXTRA_LINK_LIBS RapidJSON ${GFLAGS_LIBRARIES})
add_arrow_test(json_integration_test
EXTRA_LINK_LIBS
RapidJSON
simdjson::simdjson
${GFLAGS_LIBRARIES})
add_dependencies(arrow-integration arrow-json-integration-test)
elseif(ARROW_BUILD_INTEGRATION)
add_executable(arrow-json-integration-test json_integration_test.cc)
target_link_libraries(arrow-json-integration-test
RapidJSON
simdjson::simdjson
${ARROW_TEST_LINK_LIBS}
${GFLAGS_LIBRARIES}
${ARROW_GTEST_GTEST})
Expand Down
35 changes: 17 additions & 18 deletions cpp/src/arrow/integration/json_integration.cc
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@
#include "arrow/integration/json_internal.h"
#include "arrow/io/file.h"
#include "arrow/ipc/dictionary.h"
#include "arrow/json/json_writer_internal.h"
#include "arrow/record_batch.h"
#include "arrow/result.h"
#include "arrow/status.h"
Expand All @@ -36,6 +37,8 @@
using arrow::ipc::DictionaryFieldMapper;
using arrow::ipc::DictionaryMemo;

using JsonWriter = arrow::json::JsonWriter;

namespace arrow::internal::integration {

// ----------------------------------------------------------------------
Expand All @@ -44,13 +47,10 @@ namespace arrow::internal::integration {
class IntegrationJsonWriter::Impl {
public:
explicit Impl(const std::shared_ptr<Schema>& schema)
: schema_(schema), mapper_(*schema), first_batch_written_(false) {
writer_.reset(new RjWriter(string_buffer_));
}

: schema_(schema), mapper_(*schema), first_batch_written_(false) {}
Status Start() {
writer_->StartObject();
RETURN_NOT_OK(json::WriteSchema(*schema_, mapper_, writer_.get()));
writer_.StartObject();
RETURN_NOT_OK(json::WriteSchema(*schema_, mapper_, &writer_));
return Status::OK();
}

Expand All @@ -59,26 +59,26 @@ class IntegrationJsonWriter::Impl {

// Write dictionaries, if any
if (!dictionaries.empty()) {
writer_->Key("dictionaries");
writer_->StartArray();
writer_.Key("dictionaries");
writer_.StartArray();
for (const auto& entry : dictionaries) {
RETURN_NOT_OK(json::WriteDictionary(entry.first, entry.second, writer_.get()));
RETURN_NOT_OK(json::WriteDictionary(entry.first, entry.second, &writer_));
}
writer_->EndArray();
writer_.EndArray();
}

// Record batches
writer_->Key("batches");
writer_->StartArray();
writer_.Key("batches");
writer_.StartArray();
first_batch_written_ = true;
return Status::OK();
}

Result<std::string> Finish() {
writer_->EndArray(); // Record batches
writer_->EndObject();
writer_.EndArray(); // Record batches
writer_.EndObject();

return string_buffer_.GetString();
return std::string(writer_.GetString());
}

Status WriteRecordBatch(const RecordBatch& batch) {
Expand All @@ -87,7 +87,7 @@ class IntegrationJsonWriter::Impl {
if (!first_batch_written_) {
RETURN_NOT_OK(FirstRecordBatch(batch));
}
return json::WriteRecordBatch(batch, writer_.get());
return json::WriteRecordBatch(batch, &writer_);
}

private:
Expand All @@ -96,8 +96,7 @@ class IntegrationJsonWriter::Impl {

bool first_batch_written_;

rj::StringBuffer string_buffer_;
std::unique_ptr<RjWriter> writer_;
JsonWriter writer_;
};

IntegrationJsonWriter::IntegrationJsonWriter(const std::shared_ptr<Schema>& schema) {
Expand Down
12 changes: 5 additions & 7 deletions cpp/src/arrow/integration/json_integration_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@
#include "arrow/ipc/reader.h"
#include "arrow/ipc/test_common.h"
#include "arrow/ipc/writer.h"
#include "arrow/json/json_writer_internal.h"
#include "arrow/pretty_print.h"
#include "arrow/status.h"
#include "arrow/testing/builder.h"
Expand Down Expand Up @@ -723,16 +724,15 @@ static const char* json_example6 = R"example(
)example";

void TestSchemaRoundTrip(const std::shared_ptr<Schema>& schema) {
rj::StringBuffer sb;
rj::Writer<rj::StringBuffer> writer(sb);
arrow::json::JsonWriter writer;

DictionaryFieldMapper mapper(*schema);

writer.StartObject();
ASSERT_OK(json::WriteSchema(*schema, mapper, &writer));
writer.EndObject();

std::string json_schema = sb.GetString();
std::string json_schema(writer.GetString());

rj::Document d;
// Pass explicit size to avoid ASAN issues with
Expand All @@ -748,12 +748,11 @@ void TestSchemaRoundTrip(const std::shared_ptr<Schema>& schema) {
void TestArrayRoundTrip(const Array& array) {
static std::string name = "dummy";

rj::StringBuffer sb;
rj::Writer<rj::StringBuffer> writer(sb);
arrow::json::JsonWriter writer;

ASSERT_OK(json::WriteArray(name, array, &writer));

std::string array_as_json = sb.GetString();
std::string array_as_json(writer.GetString());

rj::Document d;
// Pass explicit size to avoid ASAN issues with
Expand All @@ -768,7 +767,6 @@ void TestArrayRoundTrip(const Array& array) {
json::ReadArray(default_memory_pool(), d, ::arrow::field(name, array.type())));
ASSERT_OK(result_array->ValidateFull());

// std::cout << array_as_json << std::endl;
CompareArraysDetailed(0, *result_array, array);
}

Expand Down
45 changes: 22 additions & 23 deletions cpp/src/arrow/integration/json_internal.cc
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@
#include "arrow/array/builder_time.h"
#include "arrow/extension_type.h"
#include "arrow/ipc/dictionary.h"
#include "arrow/json/json_writer_internal.h"
#include "arrow/record_batch.h"
#include "arrow/result.h"
#include "arrow/scalar.h"
Expand Down Expand Up @@ -64,6 +65,8 @@ using arrow::ipc::DictionaryFieldMapper;
using arrow::ipc::DictionaryMemo;
using arrow::ipc::internal::FieldPosition;

using JsonWriter = arrow::json::JsonWriter;

namespace arrow::internal::integration::json {

namespace {
Expand Down Expand Up @@ -118,7 +121,7 @@ Result<std::string_view> GetStringView(const rj::Value& str) {
class SchemaWriter {
public:
explicit SchemaWriter(const Schema& schema, const DictionaryFieldMapper& mapper,
RjWriter* writer)
JsonWriter* writer)
: schema_(schema), mapper_(mapper), writer_(writer) {}

Status Write() {
Expand Down Expand Up @@ -460,7 +463,7 @@ class SchemaWriter {
private:
const Schema& schema_;
const DictionaryFieldMapper& mapper_;
RjWriter* writer_;
JsonWriter* writer_;
};

Status SchemaWriter::VisitType(const DataType& type) {
Expand All @@ -469,7 +472,7 @@ Status SchemaWriter::VisitType(const DataType& type) {

class ArrayWriter {
public:
ArrayWriter(const std::string& name, const Array& array, RjWriter* writer)
ArrayWriter(const std::string& name, const Array& array, JsonWriter* writer)
: name_(name), array_(array), writer_(writer) {}
Comment on lines +475 to 476

Status Write() { return VisitArray(name_, array_); }
Expand All @@ -490,11 +493,7 @@ class ArrayWriter {
return Status::OK();
}

void WriteRawNumber(std::string_view v) {
// Avoid RawNumber() as it misleadingly adds quotes
// (see https://github.com/Tencent/rapidjson/pull/1155)
writer_->RawValue(v.data(), v.size(), rj::kNumberType);
}
void WriteRawNumber(std::string_view v) { writer_->RawValue(v); }

template <typename ArrayType, typename TypeClass = typename ArrayType::TypeClass,
typename CType = typename TypeClass::c_type>
Expand Down Expand Up @@ -522,11 +521,10 @@ class ArrayWriter {
for (int64_t i = 0; i < arr.length(); ++i) {
if (arr.IsValid(i)) {
fmt(arr.Value(i), [&](std::string_view repr) {
writer_->String(repr.data(), static_cast<rj::SizeType>(repr.size()));
writer_->String(std::string_view(repr.data(), repr.size()));
});
} else {
writer_->String(null_string.data(),
static_cast<rj::SizeType>(null_string.size()));
writer_->String(std::string_view(null_string.data(), null_string.size()));
}
}
}
Expand All @@ -553,7 +551,7 @@ class ArrayWriter {
if constexpr (Type::is_utf8) {
// UTF8 string, write as is
auto view = arr.GetView(i);
writer_->String(view.data(), static_cast<rj::SizeType>(view.size()));
writer_->String(std::string_view(view.data(), view.size()));
} else {
// Binary, encode to hexadecimal.
writer_->String(HexEncode(arr.GetView(i)));
Expand Down Expand Up @@ -598,7 +596,7 @@ class ArrayWriter {
const Decimal32 value(arr.GetValue(i));
writer_->String(value.ToIntegerString());
} else {
writer_->String(null_string, sizeof(null_string));
writer_->String(std::string_view(null_string));
}
}
}
Expand All @@ -610,7 +608,7 @@ class ArrayWriter {
const Decimal64 value(arr.GetValue(i));
writer_->String(value.ToIntegerString());
} else {
writer_->String(null_string, sizeof(null_string));
writer_->String(std::string_view(null_string));
}
}
}
Expand All @@ -622,7 +620,7 @@ class ArrayWriter {
const Decimal128 value(arr.GetValue(i));
writer_->String(value.ToIntegerString());
} else {
writer_->String(null_string, sizeof(null_string));
writer_->String(std::string_view(null_string));
}
}
}
Expand All @@ -634,7 +632,7 @@ class ArrayWriter {
const Decimal256 value(arr.GetValue(i));
writer_->String(value.ToIntegerString());
} else {
writer_->String(null_string, sizeof(null_string));
writer_->String(std::string_view(null_string));
}
}
}
Expand Down Expand Up @@ -670,7 +668,7 @@ class ArrayWriter {
// them exactly.
::arrow::internal::StringFormatter<typename CTypeTraits<T>::ArrowType> formatter;
auto append = [this](std::string_view v) {
writer_->String(v.data(), static_cast<rj::SizeType>(v.size()));
writer_->String(std::string_view(v.data(), v.size()));
return Status::OK();
};
for (int i = 0; i < length; ++i) {
Expand All @@ -692,7 +690,8 @@ class ArrayWriter {
if (s.is_inline()) {
writer_->Key("INLINED");
if constexpr (ArrayType::TypeClass::is_utf8) {
writer_->String(reinterpret_cast<const char*>(s.inline_data()), s.size());
writer_->String(
std::string_view(reinterpret_cast<const char*>(s.inline_data()), s.size()));
} else {
writer_->String(HexEncode(s.inline_data(), s.size()));
}
Expand Down Expand Up @@ -863,7 +862,7 @@ class ArrayWriter {
private:
const std::string& name_;
const Array& array_;
RjWriter* writer_;
JsonWriter* writer_;
};

Result<TimeUnit::type> GetUnitFromString(const std::string& unit_str) {
Expand Down Expand Up @@ -2035,13 +2034,13 @@ Result<std::shared_ptr<RecordBatch>> ReadRecordBatch(
}

Status WriteSchema(const Schema& schema, const DictionaryFieldMapper& mapper,
RjWriter* json_writer) {
JsonWriter* json_writer) {
SchemaWriter converter(schema, mapper, json_writer);
return converter.Write();
}

Status WriteDictionary(int64_t id, const std::shared_ptr<Array>& dictionary,
RjWriter* writer) {
JsonWriter* writer) {
writer->StartObject();
writer->Key("id");
writer->Int(static_cast<int32_t>(id));
Expand All @@ -2055,7 +2054,7 @@ Status WriteDictionary(int64_t id, const std::shared_ptr<Array>& dictionary,
return Status::OK();
}

Status WriteRecordBatch(const RecordBatch& batch, RjWriter* writer) {
Status WriteRecordBatch(const RecordBatch& batch, JsonWriter* writer) {
writer->StartObject();
writer->Key("count");
writer->Int(static_cast<int32_t>(batch.num_rows()));
Expand All @@ -2076,7 +2075,7 @@ Status WriteRecordBatch(const RecordBatch& batch, RjWriter* writer) {
return Status::OK();
}

Status WriteArray(const std::string& name, const Array& array, RjWriter* json_writer) {
Status WriteArray(const std::string& name, const Array& array, JsonWriter* json_writer) {
ArrayWriter converter(name, array, json_writer);
return converter.Write();
}
Expand Down
15 changes: 9 additions & 6 deletions cpp/src/arrow/integration/json_internal.h
Original file line number Diff line number Diff line change
Expand Up @@ -36,9 +36,12 @@
#include "arrow/util/visibility.h"

namespace rj = arrow::rapidjson;
using RjWriter = rj::Writer<rj::StringBuffer>;
using RjArray = rj::Value::ConstArray;
using RjObject = rj::Value::ConstObject;
using RjArray = rj::Value::ConstArray;

namespace arrow::json {
class JsonWriter;
} // namespace arrow::json

#define RETURN_NOT_FOUND(TOK, NAME, PARENT) \
if (NAME == (PARENT).MemberEnd()) { \
Expand Down Expand Up @@ -80,17 +83,17 @@ namespace arrow::internal::integration::json {
/// \brief Append integration test Schema format to rapidjson writer
ARROW_EXPORT
Status WriteSchema(const Schema& schema, const ipc::DictionaryFieldMapper& mapper,
RjWriter* writer);
arrow::json::JsonWriter*);

ARROW_EXPORT
Status WriteDictionary(int64_t id, const std::shared_ptr<Array>& dictionary,
RjWriter* writer);
arrow::json::JsonWriter*);

ARROW_EXPORT
Status WriteRecordBatch(const RecordBatch& batch, RjWriter* writer);
Status WriteRecordBatch(const RecordBatch& batch, arrow::json::JsonWriter*);

ARROW_EXPORT
Status WriteArray(const std::string& name, const Array& array, RjWriter* writer);
Status WriteArray(const std::string& name, const Array& array, arrow::json::JsonWriter*);

ARROW_EXPORT
Result<std::shared_ptr<Schema>> ReadSchema(const rj::Value& json_obj, MemoryPool* pool,
Expand Down
Loading
Loading