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
133 changes: 133 additions & 0 deletions be/src/core/data_type_serde/data_type_variant_v2_serde.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
#include "core/data_type_serde/data_type_variant_v2_serde.h"

#include <arrow/array/builder_binary.h>
#include <arrow/array/builder_nested.h>

#include <algorithm>
#include <cstring>
Expand Down Expand Up @@ -175,6 +176,133 @@ void preflight_json(const IColumn& column, size_t start, size_t end,
});
}

void validate_paimon_variant_primitive(VariantPrimitiveId primitive_id) {
switch (primitive_id) {
case VariantPrimitiveId::NULL_VALUE:
case VariantPrimitiveId::TRUE_VALUE:
case VariantPrimitiveId::FALSE_VALUE:
case VariantPrimitiveId::INT8:
case VariantPrimitiveId::INT16:
case VariantPrimitiveId::INT32:
case VariantPrimitiveId::INT64:
case VariantPrimitiveId::DOUBLE:
case VariantPrimitiveId::DECIMAL4:
case VariantPrimitiveId::DECIMAL8:
case VariantPrimitiveId::DECIMAL16:
case VariantPrimitiveId::DATE:
case VariantPrimitiveId::TIMESTAMP_MICROS:
case VariantPrimitiveId::TIMESTAMP_NTZ_MICROS:
case VariantPrimitiveId::FLOAT:
case VariantPrimitiveId::BINARY:
case VariantPrimitiveId::STRING:
case VariantPrimitiveId::UUID:
return;
case VariantPrimitiveId::TIME_NTZ_MICROS:
case VariantPrimitiveId::TIMESTAMP_NANOS:
case VariantPrimitiveId::TIMESTAMP_NTZ_NANOS:
throw Exception(ErrorCode::NOT_IMPLEMENTED_ERROR,
"Paimon does not support Variant primitive id {}",
static_cast<uint8_t>(primitive_id));
}
throw Exception(ErrorCode::NOT_IMPLEMENTED_ERROR,
"Paimon does not support unknown Variant primitive id {}",
static_cast<uint8_t>(primitive_id));
}

void validate_paimon_variant_value(VariantRef value, uint32_t depth = 0) {
if (depth > VARIANT_MAX_NESTING_DEPTH) {
throw Exception(ErrorCode::CORRUPTION, "Variant value exceeds maximum nesting depth {}",
VARIANT_MAX_NESTING_DEPTH);
}
const size_t encoded_size = value.value_size();
if (encoded_size != value.value.size) {
throw Exception(ErrorCode::CORRUPTION,
"Variant value has {} trailing bytes after the encoded value",
value.value.size - encoded_size);
}

switch (value.basic_type()) {
case VariantBasicType::PRIMITIVE:
validate_paimon_variant_primitive(value.primitive_id());
return;
case VariantBasicType::SHORT_STRING:
return;
case VariantBasicType::OBJECT:
for (uint32_t i = 0; i < value.num_elements(); ++i) {
uint32_t field_id = 0;
VariantRef child = value.object_value_at(i, &field_id);
value.metadata.key_at(field_id);
validate_paimon_variant_value(child, depth + 1);
}
return;
case VariantBasicType::ARRAY:
for (uint32_t i = 0; i < value.num_elements(); ++i) {
validate_paimon_variant_value(value.array_at(i), depth + 1);
}
return;
}
}

void require_variant_arrow_status(const arrow::Status& status) {
if (!status.ok()) {
throw Exception(ErrorCode::INTERNAL_ERROR, "Variant V2 Arrow append failed: {}",
status.ToString());
}
}

Status write_binary_variant_arrow(const IColumn& column, const NullMap* null_map,
arrow::StructBuilder& builder, size_t start, size_t end) {
// StructBuilder::type() returns a shared_ptr by value. Keep that owner alive while using the
// cast reference; otherwise the reference would dangle as soon as the temporary is destroyed.
const auto builder_type = builder.type();
const auto& struct_type = assert_cast<const arrow::StructType&>(*builder_type);
if (struct_type.num_fields() != 2 || struct_type.field(0)->name() != "value" ||
struct_type.field(1)->name() != "metadata" ||
struct_type.field(0)->type()->id() != arrow::Type::BINARY ||
struct_type.field(1)->type()->id() != arrow::Type::BINARY) {
return Status::InvalidArgument(
"Binary Variant V2 Arrow type must be "
"struct<value: binary, metadata: binary>, got {}",
struct_type.ToString());
}
auto* value_builder = dynamic_cast<arrow::BinaryBuilder*>(builder.field_builder(0));
auto* metadata_builder = dynamic_cast<arrow::BinaryBuilder*>(builder.field_builder(1));
if (value_builder == nullptr || metadata_builder == nullptr) {
return Status::InvalidArgument("Binary Variant V2 Arrow child builders must be binary");
}

// GenericVariant assumes its input is valid, and Paimon's unshredded writer copies these two
// buffers without inspecting them. Validate once at the Doris-to-Paimon boundary so a write
// cannot commit bytes which Paimon is unable to read later.
const auto outer_nulls = forced_nulls(null_map);
visit_variant_v2_values(
column, start, end, outer_nulls,
[&](size_t) { require_variant_arrow_status(builder.AppendNull()); },
[&](size_t row, VariantRef value) {
try {
constexpr size_t PAIMON_VARIANT_SIZE_LIMIT = 128 * 1024 * 1024;
Comment thread
suxiaogang223 marked this conversation as resolved.
if (value.value.size > PAIMON_VARIANT_SIZE_LIMIT ||
value.metadata.size > PAIMON_VARIANT_SIZE_LIMIT) {
throw Exception(ErrorCode::INVALID_ARGUMENT,
"exceeds the 128 MiB value/metadata limit");
}
value.metadata.validate();
Comment thread
suxiaogang223 marked this conversation as resolved.
validate_paimon_variant_value(value);
} catch (const Exception& e) {
throw Exception(e.code(), "Paimon Variant V2 row {} is incompatible: {}", row,
e.what());
}
require_variant_arrow_status(builder.Append());
Comment thread
suxiaogang223 marked this conversation as resolved.
require_variant_arrow_status(
value_builder->Append(reinterpret_cast<const uint8_t*>(value.value.data),
cast_set<int32_t, size_t, false>(value.value.size)));
require_variant_arrow_status(metadata_builder->Append(
reinterpret_cast<const uint8_t*>(value.metadata.data),
cast_set<int32_t, size_t, false>(value.metadata.size)));
});
return Status::OK();
}

} // namespace

DataTypeVariantV2SerDe::DataTypeVariantV2SerDe(int nesting_level) : DataTypeSerDe(nesting_level) {}
Expand Down Expand Up @@ -537,6 +665,11 @@ Status DataTypeVariantV2SerDe::write_column_to_arrow(const IColumn& column, cons
assert_cast<arrow::LargeStringBuilder&>(*array_builder), first, last,
options);
}
if (array_builder->type()->id() == arrow::Type::STRUCT) {
return write_binary_variant_arrow(column, null_map,
assert_cast<arrow::StructBuilder&>(*array_builder),
first, last);
}
return Status::InvalidArgument("Unsupported arrow type for variant column: {}",
array_builder->type()->name());
});
Expand Down
70 changes: 69 additions & 1 deletion be/src/exec/sink/writer/paimon/jni_paimon_write_backend.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,10 @@

#include "common/check.h"
#include "common/logging.h"
#include "core/data_type/data_type_agg_state.h"
#include "core/data_type/data_type_array.h"
#include "core/data_type/data_type_map.h"
#include "core/data_type/data_type_struct.h"
#include "exec/sink/writer/paimon/paimon_jni_memory_manager.h"
#include "format/arrow/arrow_block_convertor.h"
#include "format/arrow/arrow_row_batch.h"
Expand Down Expand Up @@ -69,6 +73,69 @@ void retain_memory_after_failed_close(std::unique_ptr<PaimonJniMemoryManager> ma
std::lock_guard<std::mutex> lock(retained_memory_managers_mutex());
retained_memory_managers().emplace_back(std::move(manager));
}

Status convert_to_paimon_arrow_type(const DataTypePtr& origin_type,
std::shared_ptr<arrow::DataType>* result,
const std::string& timezone) {
const DataTypePtr type = get_serialized_type(origin_type);
switch (type->get_primitive_type()) {
case TYPE_VARIANT:
// Paimon consumes the lossless Variant V2 representation. Keeping both children non-null
// distinguishes a SQL NULL struct from a non-null Variant value.
*result = arrow::struct_({arrow::field("value", arrow::binary(), false),
arrow::field("metadata", arrow::binary(), false)});
return Status::OK();
case TYPE_ARRAY: {
const auto& array_type = assert_cast<const DataTypeArray&>(*remove_nullable(type));
std::shared_ptr<arrow::DataType> element_type;
RETURN_IF_ERROR(convert_to_paimon_arrow_type(array_type.get_nested_type(), &element_type,
timezone));
*result = std::make_shared<arrow::ListType>(element_type);
return Status::OK();
}
case TYPE_MAP: {
const auto& map_type = assert_cast<const DataTypeMap&>(*remove_nullable(type));
std::shared_ptr<arrow::DataType> key_type;
std::shared_ptr<arrow::DataType> value_type;
RETURN_IF_ERROR(convert_to_paimon_arrow_type(map_type.get_key_type(), &key_type, timezone));
RETURN_IF_ERROR(
convert_to_paimon_arrow_type(map_type.get_value_type(), &value_type, timezone));
*result = std::make_shared<arrow::MapType>(key_type, value_type);
return Status::OK();
}
case TYPE_STRUCT: {
const auto& struct_type = assert_cast<const DataTypeStruct&>(*remove_nullable(type));
std::vector<std::shared_ptr<arrow::Field>> fields;
fields.reserve(struct_type.get_elements().size());
for (size_t i = 0; i < struct_type.get_elements().size(); ++i) {
const DataTypePtr& element = struct_type.get_element(i);
std::shared_ptr<arrow::DataType> field_type;
RETURN_IF_ERROR(convert_to_paimon_arrow_type(element, &field_type, timezone));
fields.push_back(arrow::field(struct_type.get_element_name(i), field_type,
element->is_nullable()));
}
*result = arrow::struct_(std::move(fields));
return Status::OK();
}
default:
return convert_to_arrow_type(origin_type, result, timezone);
}
}

Status get_paimon_arrow_schema_from_block(const Block& block,
std::shared_ptr<arrow::Schema>* result) {
std::vector<std::shared_ptr<arrow::Field>> fields;
fields.reserve(block.columns());
for (const auto& type_and_name : block) {
std::shared_ptr<arrow::DataType> arrow_type;
RETURN_IF_ERROR(convert_to_paimon_arrow_type(type_and_name.type, &arrow_type, ""));
fields.push_back(create_arrow_field_with_metadata(
type_and_name.name, arrow_type, type_and_name.type->is_nullable(),
type_and_name.type->get_primitive_type()));
}
*result = arrow::schema(std::move(fields));
return Status::OK();
}
} // namespace

// ────────────────────────────────────────────────────────────
Expand Down Expand Up @@ -359,8 +426,9 @@ Status JniPaimonWriter::_write_projected_block(RuntimeState* state, Block& block
// Step 1: Build Arrow schema from the projected Block.
// Paimon write timestamps are transported as civil-time fields. The Java writer uses the
// pinned Paimon target type to preserve NTZ values or convert LTZ values with the session zone.
// Variant V2 is transported losslessly as its value/metadata pair, including nested Variant.
std::shared_ptr<arrow::Schema> arrow_schema;
RETURN_IF_ERROR(get_arrow_schema_from_block(block, &arrow_schema, ""));
RETURN_IF_ERROR(get_paimon_arrow_schema_from_block(block, &arrow_schema));

// Step 2: Convert Doris Block columns to an Arrow RecordBatch.
std::shared_ptr<arrow::RecordBatch> record_batch;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -137,7 +137,9 @@ Status build_array_node_plan(const ColumnPtr& source, const DataTypePtr& source_

Status build_array_leaf_plan(const ColumnPtr& source, PrimitiveType primitive,
ArrayEncodePlan* plan) {
if (primitive == INVALID_TYPE && source->empty()) {
if (primitive == INVALID_TYPE) {
// DataTypeNothing is represented by the element null map, including non-empty
// expressions such as array(NULL).
return Status::OK();
} else if (primitive == TYPE_VARIANT) {
const auto* variant = check_and_get_column<ColumnVariantV2>(source.get());
Expand Down Expand Up @@ -196,7 +198,7 @@ void append_array_value(const ArrayEncodePlan& plan, size_t index, VariantBatchB
} else if (plan.jsonb_leaf != nullptr) {
jsonb_to_variant(plan.jsonb_leaf->get_data_at(index), *row);
} else {
DORIS_CHECK(false) << "empty Array leaf unexpectedly contains a value";
DORIS_CHECK(false) << "Array Variant V2 leaf has no encoder";
}
return;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -70,6 +70,7 @@
#include "core/data_type/data_type_quantilestate.h"
#include "core/data_type/data_type_string.h"
#include "core/data_type/data_type_struct.h"
#include "core/data_type/data_type_variant_v2.h"
#include "core/data_type/define_primitive_type.h"
#include "core/field.h"
#include "core/types.h"
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,9 @@
// specific language governing permissions and limitations
// under the License.

#include <arrow/array/array_nested.h>
#include <arrow/array/builder_binary.h>
#include <arrow/array/builder_nested.h>
#include <gtest/gtest.h>

#include <array>
Expand Down Expand Up @@ -198,6 +200,50 @@ std::vector<std::optional<std::string>> orc_values(const DataTypeVariantV2SerDe&
return result;
}

std::shared_ptr<arrow::DataType> binary_variant_arrow_type() {
return arrow::struct_({arrow::field("value", arrow::binary(), false),
arrow::field("metadata", arrow::binary(), false)});
}

std::unique_ptr<arrow::StructBuilder> binary_variant_arrow_builder() {
return std::make_unique<arrow::StructBuilder>(
binary_variant_arrow_type(), arrow::default_memory_pool(),
std::vector<std::shared_ptr<arrow::ArrayBuilder>> {
std::make_shared<arrow::BinaryBuilder>(arrow::default_memory_pool()),
std::make_shared<arrow::BinaryBuilder>(arrow::default_memory_pool())});
}

void expect_binary_variant_bytes(const DataTypeVariantV2SerDe& serde, const IColumn& column,
const ColumnVariantV2& encoded,
const NullMap* null_map = nullptr) {
auto builder = binary_variant_arrow_builder();
const Status status = serde.write_column_to_arrow(column, null_map, builder.get(), 0,
column.size(), cctz::utc_time_zone());
ASSERT_TRUE(status.ok()) << status;

std::shared_ptr<arrow::Array> output;
ASSERT_TRUE(builder->Finish(&output).ok());
const auto& array = assert_cast<const arrow::StructArray&>(*output);
const auto& values = assert_cast<const arrow::BinaryArray&>(*array.field(0));
const auto& metadata = assert_cast<const arrow::BinaryArray&>(*array.field(1));
ASSERT_EQ(array.length(), static_cast<int64_t>(column.size()));
const auto view = encoded.read_view();
for (size_t row = 0; row < column.size(); ++row) {
const bool expected_null = null_map != nullptr && (*null_map)[row] != 0;
EXPECT_EQ(array.IsNull(row), expected_null);
if (expected_null) {
continue;
}
const VariantRef expected = view.value_at(row);
const auto actual_value = values.GetView(row);
const auto actual_metadata = metadata.GetView(row);
EXPECT_EQ(std::string_view(actual_value.data(), actual_value.size()),
std::string_view(expected.value.data, expected.value.size));
EXPECT_EQ(std::string_view(actual_metadata.data(), actual_metadata.size()),
std::string_view(expected.metadata.data, expected.metadata.size));
}
}

// NOLINTNEXTLINE(readability-function-cognitive-complexity) -- GTest macros inflate the matrix.
void expect_text_surfaces(const DataTypeVariantV2SerDe& serde, const IColumn& encoded,
const ColumnVariantV2& typed,
Expand Down Expand Up @@ -379,4 +425,34 @@ TEST(DataTypeVariantV2SerdeOutputTest, ConstNullableAndOuterMasksPreserveBoundar
EXPECT_TRUE(invalid_dates->is_typed());
}

TEST(DataTypeVariantV2SerdeOutputTest, BinaryStructPreservesEncodedAndTypedBytesAndOuterNulls) {
DataTypeVariantV2SerDe serde;
auto documents = encoded_json({R"({"a":[1,null,"x"]})", R"({"hidden":true})", "null"});
NullMap mask {0, 1, 0};
expect_binary_variant_bytes(serde, *documents, *documents, &mask);

auto typed = typed_strings(
{std::string_view("plain"), std::nullopt, std::string_view(R"({"text":"value"})")});
ColumnPtr encoded = encoded_copy(*typed);
expect_binary_variant_bytes(serde, *typed, assert_cast<const ColumnVariantV2&>(*encoded));
EXPECT_TRUE(typed->is_typed());
}

TEST(DataTypeVariantV2SerdeOutputTest, BinaryStructRejectsUnsupportedPaimonPrimitive) {
DataTypeVariantV2SerDe serde;
VariantBatchBuilder builder(VariantBatchBuilder::ReserveHint {.rows = 1});
auto row = builder.begin_row();
row.add_time_ntz_micros(1'500'000);
row.finish();
auto encoded = ColumnVariantV2::create();
encoded->insert_encoded_batch(builder.finish_batch());
auto arrow_builder = binary_variant_arrow_builder();
const Status status = serde.write_column_to_arrow(*encoded, nullptr, arrow_builder.get(), 0,
encoded->size(), cctz::utc_time_zone());
EXPECT_EQ(status.code(), ErrorCode::NOT_IMPLEMENTED_ERROR);
EXPECT_NE(status.to_string().find("Paimon does not support Variant primitive id 17"),
std::string::npos);
EXPECT_EQ(arrow_builder->length(), 0);
}

} // namespace doris
Loading
Loading