diff --git a/be/src/core/data_type_serde/data_type_variant_v2_serde.cpp b/be/src/core/data_type_serde/data_type_variant_v2_serde.cpp index 57521e7f34aa0d..b5da6bb25cd2f5 100644 --- a/be/src/core/data_type_serde/data_type_variant_v2_serde.cpp +++ b/be/src/core/data_type_serde/data_type_variant_v2_serde.cpp @@ -18,6 +18,7 @@ #include "core/data_type_serde/data_type_variant_v2_serde.h" #include +#include #include #include @@ -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(primitive_id)); + } + throw Exception(ErrorCode::NOT_IMPLEMENTED_ERROR, + "Paimon does not support unknown Variant primitive id {}", + static_cast(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(*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, got {}", + struct_type.ToString()); + } + auto* value_builder = dynamic_cast(builder.field_builder(0)); + auto* metadata_builder = dynamic_cast(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; + 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(); + 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()); + require_variant_arrow_status( + value_builder->Append(reinterpret_cast(value.value.data), + cast_set(value.value.size))); + require_variant_arrow_status(metadata_builder->Append( + reinterpret_cast(value.metadata.data), + cast_set(value.metadata.size))); + }); + return Status::OK(); +} + } // namespace DataTypeVariantV2SerDe::DataTypeVariantV2SerDe(int nesting_level) : DataTypeSerDe(nesting_level) {} @@ -537,6 +665,11 @@ Status DataTypeVariantV2SerDe::write_column_to_arrow(const IColumn& column, cons assert_cast(*array_builder), first, last, options); } + if (array_builder->type()->id() == arrow::Type::STRUCT) { + return write_binary_variant_arrow(column, null_map, + assert_cast(*array_builder), + first, last); + } return Status::InvalidArgument("Unsupported arrow type for variant column: {}", array_builder->type()->name()); }); diff --git a/be/src/exec/sink/writer/paimon/jni_paimon_write_backend.cpp b/be/src/exec/sink/writer/paimon/jni_paimon_write_backend.cpp index 5f6a7930da6d24..e59668759df65b 100644 --- a/be/src/exec/sink/writer/paimon/jni_paimon_write_backend.cpp +++ b/be/src/exec/sink/writer/paimon/jni_paimon_write_backend.cpp @@ -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" @@ -69,6 +73,69 @@ void retain_memory_after_failed_close(std::unique_ptr ma std::lock_guard 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* 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(*remove_nullable(type)); + std::shared_ptr element_type; + RETURN_IF_ERROR(convert_to_paimon_arrow_type(array_type.get_nested_type(), &element_type, + timezone)); + *result = std::make_shared(element_type); + return Status::OK(); + } + case TYPE_MAP: { + const auto& map_type = assert_cast(*remove_nullable(type)); + std::shared_ptr key_type; + std::shared_ptr 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(key_type, value_type); + return Status::OK(); + } + case TYPE_STRUCT: { + const auto& struct_type = assert_cast(*remove_nullable(type)); + std::vector> 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 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* result) { + std::vector> fields; + fields.reserve(block.columns()); + for (const auto& type_and_name : block) { + std::shared_ptr 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 // ──────────────────────────────────────────────────────────── @@ -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; - 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 record_batch; diff --git a/be/src/exprs/function/cast/variant_v2/cast_array_to_variant.cpp b/be/src/exprs/function/cast/variant_v2/cast_array_to_variant.cpp index e35c9b1983b896..66a61f2da9623d 100644 --- a/be/src/exprs/function/cast/variant_v2/cast_array_to_variant.cpp +++ b/be/src/exprs/function/cast/variant_v2/cast_array_to_variant.cpp @@ -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(source.get()); @@ -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; } diff --git a/be/test/core/data_type_serde/data_type_serde_arrow_test.cpp b/be/test/core/data_type_serde/data_type_serde_arrow_test.cpp index 7b20c2f82d0c65..27e128b1a6d7a1 100644 --- a/be/test/core/data_type_serde/data_type_serde_arrow_test.cpp +++ b/be/test/core/data_type_serde/data_type_serde_arrow_test.cpp @@ -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" diff --git a/be/test/core/data_type_serde/data_type_variant_v2_serde_output_test.cpp b/be/test/core/data_type_serde/data_type_variant_v2_serde_output_test.cpp index fa900a392c0d6d..660ba5409038c9 100644 --- a/be/test/core/data_type_serde/data_type_variant_v2_serde_output_test.cpp +++ b/be/test/core/data_type_serde/data_type_variant_v2_serde_output_test.cpp @@ -15,7 +15,9 @@ // specific language governing permissions and limitations // under the License. +#include #include +#include #include #include @@ -198,6 +200,50 @@ std::vector> orc_values(const DataTypeVariantV2SerDe& return result; } +std::shared_ptr binary_variant_arrow_type() { + return arrow::struct_({arrow::field("value", arrow::binary(), false), + arrow::field("metadata", arrow::binary(), false)}); +} + +std::unique_ptr binary_variant_arrow_builder() { + return std::make_unique( + binary_variant_arrow_type(), arrow::default_memory_pool(), + std::vector> { + std::make_shared(arrow::default_memory_pool()), + std::make_shared(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 output; + ASSERT_TRUE(builder->Finish(&output).ok()); + const auto& array = assert_cast(*output); + const auto& values = assert_cast(*array.field(0)); + const auto& metadata = assert_cast(*array.field(1)); + ASSERT_EQ(array.length(), static_cast(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, @@ -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(*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 diff --git a/be/test/exprs/function/cast/cast_variant_v2_from_test.cpp b/be/test/exprs/function/cast/cast_variant_v2_from_test.cpp index a7858522cf1410..c28a3adb1e1b91 100644 --- a/be/test/exprs/function/cast/cast_variant_v2_from_test.cpp +++ b/be/test/exprs/function/cast/cast_variant_v2_from_test.cpp @@ -29,6 +29,7 @@ #include "core/data_type/data_type_date_or_datetime_v2.h" #include "core/data_type/data_type_decimal.h" #include "core/data_type/data_type_jsonb.h" +#include "core/data_type/data_type_nothing.h" #include "core/data_type/data_type_nullable.h" #include "core/data_type/data_type_number.h" #include "core/data_type/data_type_string.h" @@ -484,6 +485,29 @@ TEST(CastVariantV2FromTest, NestedArrayRoundTripPreservesNullAndEmptyArray) { EXPECT_EQ(assert_cast(values.get_nested_column()).get_data()[0], 1); } +TEST(CastVariantV2FromTest, NullOnlyArrayEncodesNonEmptyElements) { + auto array_type = std::make_shared(std::make_shared()); + MutableColumnPtr source = array_type->create_column(); + Array values {Field::create_field(Null()), Field::create_field(Null())}; + source->insert(Field::create_field(std::move(values))); + + auto variant_type = std::make_shared(); + Block block {{source->get_ptr(), array_type, "source"}, + {variant_type->create_column(), variant_type, "result"}}; + RuntimeState state; + auto context = FunctionContext::create_context(&state, {}, {}); + Status status = + create_cast_to_variant_v2_wrapper(array_type)(context.get(), block, {0}, 1, 1, nullptr); + ASSERT_TRUE(status.ok()) << status; + + VariantRef encoded = + assert_cast(*block.get_by_position(1).column).get_value_ref(0); + ASSERT_EQ(encoded.basic_type(), VariantBasicType::ARRAY); + ASSERT_EQ(encoded.num_elements(), 2); + EXPECT_TRUE(encoded.array_at(0).is_null()); + EXPECT_TRUE(encoded.array_at(1).is_null()); +} + TEST(CastVariantV2FromTest, DecimalScale38CastsAndScale39IsRejectedAtEncodingBoundary) { VariantBatchBuilder builder(VariantBatchBuilder::ReserveHint {.rows = 1}); auto row = builder.begin_row(); diff --git a/fe/be-java-extensions/java-udf/src/main/resources/package.xml b/fe/be-java-extensions/java-udf/src/main/resources/package.xml index 9ef59a2ee2df98..06b30f57e5cbf6 100644 --- a/fe/be-java-extensions/java-udf/src/main/resources/package.xml +++ b/fe/be-java-extensions/java-udf/src/main/resources/package.xml @@ -42,6 +42,7 @@ under the License. META-INF/services/org.apache.paimon* + META-INF/services/java.time.chrono.Chronology diff --git a/fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonArrowConverter.java b/fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonArrowConverter.java index 9502f8bb0e2bd2..01ee9e2709ea72 100644 --- a/fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonArrowConverter.java +++ b/fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonArrowConverter.java @@ -68,6 +68,9 @@ /** Converts Arrow columns into Paimon internal values without owning writer state. */ final class PaimonArrowConverter { + static final String VARIANT_VALUE_FIELD = "value"; + static final String VARIANT_METADATA_FIELD = "metadata"; + private final ZoneId sessionTimeZone; PaimonArrowConverter(ZoneId sessionTimeZone) { @@ -99,8 +102,16 @@ private RowReader( Object[] values(int rowIndex) { Object[] values = new Object[vectors.size()]; for (int column = 0; column < vectors.size(); column++) { - values[column] = convertVectorValue( - vectors.get(column), rowIndex, fields.get(column), targetTypes[column]); + try { + values[column] = convertVectorValue( + vectors.get(column), rowIndex, fields.get(column), targetTypes[column]); + } catch (RuntimeException e) { + throw new IllegalArgumentException( + "Failed to convert Arrow column '" + fields.get(column).getName() + + "' at row " + rowIndex + " to Paimon " + + targetTypes[column], + e); + } } return values; } @@ -108,6 +119,17 @@ Object[] values(int rowIndex) { private Object convertVectorValue( FieldVector vector, int index, Field arrowField, DataType targetType) { + if (targetType instanceof VariantType) { + if (!(vector instanceof StructVector)) { + throw new IllegalArgumentException( + "Paimon VARIANT write only supports Variant V2 Arrow " + + "struct, but got " + + vector.getField()); + } + return vector.isNull(index) + ? null + : convertVariantVector((StructVector) vector, index); + } if (vector.isNull(index)) { return null; } @@ -172,27 +194,6 @@ private Object convertToPaimonType(Object value, Field arrowField, DataType targ if (value == null) { return null; } - if (targetType instanceof VariantType) { - if (value instanceof byte[]) { - return toVariant((byte[]) value); - } - if (value instanceof BinaryString) { - return toVariant(((BinaryString) value).toBytes()); - } - if (value instanceof org.apache.arrow.vector.util.Text) { - return toVariant(((org.apache.arrow.vector.util.Text) value).copyBytes()); - } - if (value instanceof org.apache.hadoop.io.Text) { - org.apache.hadoop.io.Text text = (org.apache.hadoop.io.Text) value; - return GenericVariant.fromJson(text.toString()); - } - if (value instanceof CharSequence) { - return GenericVariant.fromJson(value.toString()); - } - throw new IllegalArgumentException( - "Paimon VARIANT requires Arrow UTF-8 JSON, but got " - + value.getClass().getName()); - } if (targetType instanceof BinaryType || targetType instanceof VarBinaryType) { if (value instanceof byte[]) { return value; @@ -249,17 +250,33 @@ private Object convertToPaimonType(Object value, Field arrowField, DataType targ } static Object convertText(byte[] value, DataType targetType) { - if (targetType instanceof VariantType) { - return toVariant(value); - } if (targetType instanceof BinaryType || targetType instanceof VarBinaryType) { return value; } return BinaryString.fromBytes(value); } - private static GenericVariant toVariant(byte[] json) { - return GenericVariant.fromJson(new String(json, StandardCharsets.UTF_8)); + private GenericVariant convertVariantVector(StructVector vector, int index) { + List children = vector.getChildrenFromFields(); + if (children.size() != 2 + || !VARIANT_VALUE_FIELD.equals(children.get(0).getName()) + || !VARIANT_METADATA_FIELD.equals(children.get(1).getName()) + || !(children.get(0) instanceof VarBinaryVector) + || !(children.get(1) instanceof VarBinaryVector)) { + throw new IllegalArgumentException( + "Paimon VARIANT binary transport requires Arrow " + + "struct, but got " + + vector.getField()); + } + VarBinaryVector valueVector = (VarBinaryVector) children.get(0); + VarBinaryVector metadataVector = (VarBinaryVector) children.get(1); + if (valueVector.isNull(index) || metadataVector.isNull(index)) { + throw new IllegalArgumentException( + "A non-null Paimon VARIANT struct requires non-null value and metadata"); + } + // BE validates the complete Variant value for Paimon before Arrow transport. + // Keep this JNI boundary allocation-only instead of traversing the same tree again. + return new GenericVariant(valueVector.get(index), metadataVector.get(index)); } private GenericRow convertStructVector( diff --git a/fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonJniScanner.java b/fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonJniScanner.java index 03a12bc82456c2..3507920163049f 100644 --- a/fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonJniScanner.java +++ b/fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonJniScanner.java @@ -44,6 +44,8 @@ import org.apache.paimon.types.RowType; import org.apache.paimon.types.TimestampType; import org.apache.paimon.types.VariantType; +import org.apache.paimon.utils.ChainTableUtils; +import org.apache.paimon.utils.StringUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -775,7 +777,7 @@ private static Table applyManifestParallelismBound( // Each branch owns an independent planner setting; a smaller sibling is not an // execution ceiling and must never throttle the other branch. return new FallbackReadFileStoreTable( - main, other, wrappedBranchHasReadPriority(pair)); + main, other, isWrappedFirst(pair)); } if (table instanceof DelegatedFileStoreTable) { @@ -817,6 +819,22 @@ private static Table applyManifestParallelismBound( return ((FileStoreTable) table).copyWithoutTimeTravel(cap); } + static boolean isWrappedFirst(FallbackReadFileStoreTable table) { + Map options = table.options(); + // Mirror Paimon 1.4.2 FileStoreTableFactory. There is no public accessor for the wrapper's + // private ordering flag, so reconstruction must recover the factory decision from options. + // Remove this inference after upgrading to an SDK which exposes the original ordering. + if (ChainTableUtils.isChainTable(options)) { + return true; + } + if (!StringUtils.isNullOrWhitespaceOnly( + options.get(CoreOptions.SCAN_FALLBACK_BRANCH.key()))) { + return true; + } + return StringUtils.isNullOrWhitespaceOnly( + options.get(CoreOptions.SCAN_PRIMARY_BRANCH.key())); + } + private static Table rebuildHiddenSystemTable(Table wrapper, FileStoreTable normalizedSource) { try { // Old FE payloads have no type/source side channel. Every Paimon 1.3 system wrapper @@ -838,11 +856,6 @@ private static FileStoreTable applyManifestParallelismBound( (Table) table, safeBound, materializeAbsent); } - private static boolean wrappedBranchHasReadPriority(FallbackReadFileStoreTable table) { - // Paimon's factory sets wrappedFirst=false only for scan.primary-branch tables. - return !table.wrapped().options().containsKey(CoreOptions.SCAN_PRIMARY_BRANCH.key()); - } - private static FileStoreTable unwrapSystemPlanningSource(FileStoreTable table) { FileStoreTable current = table; // System wrappers dispatch fallback reads only when the fallback pair is their direct diff --git a/fe/be-java-extensions/paimon-connector/src/test/java/org/apache/doris/paimon/PaimonArrowConverterTest.java b/fe/be-java-extensions/paimon-connector/src/test/java/org/apache/doris/paimon/PaimonArrowConverterTest.java index 4c75e9b27a41c2..45b7e233ecb1f2 100644 --- a/fe/be-java-extensions/paimon-connector/src/test/java/org/apache/doris/paimon/PaimonArrowConverterTest.java +++ b/fe/be-java-extensions/paimon-connector/src/test/java/org/apache/doris/paimon/PaimonArrowConverterTest.java @@ -17,8 +17,16 @@ package org.apache.doris.paimon; +import org.apache.arrow.memory.RootAllocator; +import org.apache.arrow.vector.VarBinaryVector; +import org.apache.arrow.vector.VarCharVector; +import org.apache.arrow.vector.VectorSchemaRoot; +import org.apache.arrow.vector.complex.StructVector; import org.apache.arrow.vector.types.TimeUnit; import org.apache.arrow.vector.types.pojo.ArrowType; +import org.apache.arrow.vector.types.pojo.Field; +import org.apache.arrow.vector.types.pojo.FieldType; +import org.apache.arrow.vector.types.pojo.Schema; import org.apache.paimon.data.Timestamp; import org.apache.paimon.data.variant.GenericVariant; import org.apache.paimon.types.DataTypes; @@ -34,6 +42,7 @@ import java.time.ZoneId; import java.time.ZoneOffset; import java.util.Arrays; +import java.util.Collections; public class PaimonArrowConverterTest { @@ -87,17 +96,118 @@ public void testPaimonWriteRejectsTimezoneInArrowType() { } @Test - public void testVariantJsonKindsUsePaimonVariant() { - String[] jsonValues = { - "\"scalar\"", - "[1,true,null]", - "{\"id\":1,\"name\":\"doris\"}" - }; - for (String json : jsonValues) { - Object value = PaimonArrowConverter.convertText( - json.getBytes(StandardCharsets.UTF_8), new VariantType()); - Assertions.assertInstanceOf(GenericVariant.class, value); - Assertions.assertEquals(json, ((GenericVariant) value).toJson()); + public void testVariantJsonTransportIsRejected() { + Field variantField = new Field( + "payload", FieldType.nullable(new ArrowType.Utf8()), null); + try (RootAllocator allocator = new RootAllocator(); + VectorSchemaRoot root = VectorSchemaRoot.create( + new Schema(Collections.singletonList(variantField)), allocator)) { + VarCharVector vector = (VarCharVector) root.getVector("payload"); + root.allocateNew(); + vector.setSafe(0, "{\"legacy\":true}".getBytes(StandardCharsets.UTF_8)); + root.setRowCount(1); + + PaimonArrowConverter.RowReader rows = + new PaimonArrowConverter(ZoneId.of("UTC")).rows( + root, new org.apache.paimon.types.DataType[] {new VariantType()}); + IllegalArgumentException exception = Assertions.assertThrows( + IllegalArgumentException.class, () -> rows.values(0)); + Assertions.assertTrue(exception.getCause().getMessage().contains( + "only supports Variant V2")); + } + } + + @Test + public void testVariantBinaryTransportPreservesValueAndMetadata() { + GenericVariant expected = GenericVariant.fromJson( + "{\"id\":1,\"nested\":[true,null,\"doris\"]}"); + Field variantField = variantField(); + + try (RootAllocator allocator = new RootAllocator(); + VectorSchemaRoot root = VectorSchemaRoot.create( + new Schema(Collections.singletonList(variantField)), allocator)) { + StructVector vector = (StructVector) root.getVector("payload"); + VarBinaryVector values = (VarBinaryVector) vector.getChild( + PaimonArrowConverter.VARIANT_VALUE_FIELD); + VarBinaryVector metadata = (VarBinaryVector) vector.getChild( + PaimonArrowConverter.VARIANT_METADATA_FIELD); + root.allocateNew(); + values.setSafe(0, expected.value()); + metadata.setSafe(0, expected.metadata()); + vector.setIndexDefined(0); + vector.setNull(1); + root.setRowCount(2); + + PaimonArrowConverter converter = new PaimonArrowConverter(ZoneId.of("UTC")); + PaimonArrowConverter.RowReader rows = converter.rows( + root, new org.apache.paimon.types.DataType[] {new VariantType()}); + GenericVariant actual = (GenericVariant) rows.values(0)[0]; + + Assertions.assertArrayEquals(expected.value(), actual.value()); + Assertions.assertArrayEquals(expected.metadata(), actual.metadata()); + Assertions.assertNull(rows.values(1)[0]); + } + } + + @Test + public void testVariantBinaryTransportReliesOnBeValidation() { + GenericVariant array = GenericVariant.fromJson("[0]"); + byte[] unsupportedValue = array.value().clone(); + int primitiveHeaderOffset = unsupportedValue.length - 2; + // This unit test constructs Arrow directly and therefore bypasses the BE compatibility + // check. The Java boundary intentionally remains allocation-only instead of repeating it. + // Replace INT1 (id 3) with TIME_NTZ_MICROS (id 17) to prove no second traversal occurs. + Assertions.assertEquals(3 << 2, unsupportedValue[primitiveHeaderOffset] & 0xff); + unsupportedValue[primitiveHeaderOffset] = (byte) (17 << 2); + + try (RootAllocator allocator = new RootAllocator(); + VectorSchemaRoot root = VectorSchemaRoot.create( + new Schema(Collections.singletonList(variantField())), allocator)) { + StructVector vector = (StructVector) root.getVector("payload"); + VarBinaryVector values = (VarBinaryVector) vector.getChild( + PaimonArrowConverter.VARIANT_VALUE_FIELD); + VarBinaryVector metadata = (VarBinaryVector) vector.getChild( + PaimonArrowConverter.VARIANT_METADATA_FIELD); + root.allocateNew(); + values.setSafe(0, unsupportedValue); + metadata.setSafe(0, array.metadata()); + vector.setIndexDefined(0); + root.setRowCount(1); + + PaimonArrowConverter.RowReader rows = + new PaimonArrowConverter(ZoneId.of("UTC")).rows( + root, new org.apache.paimon.types.DataType[] {new VariantType()}); + GenericVariant actual = (GenericVariant) rows.values(0)[0]; + Assertions.assertArrayEquals(unsupportedValue, actual.value()); + Assertions.assertArrayEquals(array.metadata(), actual.metadata()); + } + } + + @Test + public void testVariantBinaryTransportRejectsMissingMetadata() { + Field variantField = new Field( + "payload", + FieldType.nullable(new ArrowType.Struct()), + Arrays.asList( + new Field("value", FieldType.nullable(new ArrowType.Binary()), null), + new Field("metadata", FieldType.nullable(new ArrowType.Binary()), null))); + try (RootAllocator allocator = new RootAllocator(); + VectorSchemaRoot root = VectorSchemaRoot.create( + new Schema(Collections.singletonList(variantField)), allocator)) { + StructVector vector = (StructVector) root.getVector("payload"); + VarBinaryVector values = (VarBinaryVector) vector.getChild("value"); + root.allocateNew(); + values.setSafe(0, GenericVariant.fromJson("1").value()); + vector.setIndexDefined(0); + root.setRowCount(1); + + PaimonArrowConverter.RowReader rows = + new PaimonArrowConverter(ZoneId.of("UTC")).rows( + root, new org.apache.paimon.types.DataType[] {new VariantType()}); + IllegalArgumentException exception = Assertions.assertThrows( + IllegalArgumentException.class, () -> rows.values(0)); + Assertions.assertTrue(exception.getMessage().contains("payload")); + Assertions.assertTrue(exception.getCause().getMessage().contains("metadata")); } } @@ -116,4 +226,19 @@ private static RowType mixedCaseRowType() { DataTypes.FIELD(0, "Foo", DataTypes.INT()), DataTypes.FIELD(1, "foo", DataTypes.INT())); } + + private static Field variantField() { + return new Field( + "payload", + FieldType.nullable(new ArrowType.Struct()), + Arrays.asList( + new Field( + PaimonArrowConverter.VARIANT_VALUE_FIELD, + FieldType.notNullable(new ArrowType.Binary()), + null), + new Field( + PaimonArrowConverter.VARIANT_METADATA_FIELD, + FieldType.notNullable(new ArrowType.Binary()), + null))); + } } diff --git a/fe/be-java-extensions/paimon-connector/src/test/java/org/apache/doris/paimon/PaimonJniScannerTest.java b/fe/be-java-extensions/paimon-connector/src/test/java/org/apache/doris/paimon/PaimonJniScannerTest.java index e9bd2e1a5a0931..bce0c5950e7863 100644 --- a/fe/be-java-extensions/paimon-connector/src/test/java/org/apache/doris/paimon/PaimonJniScannerTest.java +++ b/fe/be-java-extensions/paimon-connector/src/test/java/org/apache/doris/paimon/PaimonJniScannerTest.java @@ -19,6 +19,7 @@ import org.apache.doris.common.jni.vec.ColumnType; +import com.google.common.collect.ImmutableMap; import org.apache.logging.log4j.Level; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.core.LogEvent; @@ -324,6 +325,21 @@ public void testBackendCapPreservesPrimaryBranchPriority() throws Exception { Assert.assertFalse(wrappedFirst.getBoolean(safe)); } + @Test + public void testFallbackWrapperOrderMatchesPaimonFactoryPrecedence() { + Assert.assertTrue(PaimonJniScanner.isWrappedFirst(fallbackTable(ImmutableMap.of( + CoreOptions.SCAN_FALLBACK_BRANCH.key(), "fallback", + CoreOptions.SCAN_PRIMARY_BRANCH.key(), "primary")))); + Assert.assertFalse(PaimonJniScanner.isWrappedFirst(fallbackTable(ImmutableMap.of( + CoreOptions.SCAN_PRIMARY_BRANCH.key(), "primary")))); + Assert.assertTrue(PaimonJniScanner.isWrappedFirst(fallbackTable(ImmutableMap.of( + CoreOptions.SCAN_PRIMARY_BRANCH.key(), " ")))); + Assert.assertTrue(PaimonJniScanner.isWrappedFirst(fallbackTable(ImmutableMap.of( + CoreOptions.CHAIN_TABLE_ENABLED.key(), "true", + CoreOptions.SCAN_PRIMARY_BRANCH.key(), "primary")))); + Assert.assertTrue(PaimonJniScanner.isWrappedFirst(fallbackTable(Collections.emptyMap()))); + } + @Test public void testBackendCapTraversesPrivilegeDelegate() { FileStoreTable main = serializableFileStoreTable(Collections.singletonMap( @@ -442,6 +458,12 @@ private static FileStoreTable serializableFileStoreTable(Map opt new Class[] {FileStoreTable.class}, new SerializableTableHandler(options)); } + private static FallbackReadFileStoreTable fallbackTable(Map options) { + return new FallbackReadFileStoreTable( + serializableFileStoreTable(options), + serializableFileStoreTable(Collections.emptyMap()), true); + } + @Test public void testSerializedReadBatchSizeReachesReaderTableInitialization() throws Exception { Table configuredTable = (Table) Proxy.newProxyInstance(Table.class.getClassLoader(), diff --git a/fe/fe-core/pom.xml b/fe/fe-core/pom.xml index e129de778e7644..a880c5ffad4db4 100644 --- a/fe/fe-core/pom.xml +++ b/fe/fe-core/pom.xml @@ -474,12 +474,10 @@ under the License. org.apache.paimon paimon-core - org.apache.paimon paimon-common - org.apache.paimon paimon-hive-connector-3.1 @@ -488,7 +486,6 @@ under the License. org.apache.paimon paimon-format - org.apache.paimon paimon-s3 diff --git a/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/PaimonReaderOptions.java b/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/PaimonReaderOptions.java index af352f37439197..5a422f6af5a9f2 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/PaimonReaderOptions.java +++ b/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/PaimonReaderOptions.java @@ -27,6 +27,8 @@ import org.apache.paimon.table.FallbackReadFileStoreTable; import org.apache.paimon.table.FileStoreTable; import org.apache.paimon.table.Table; +import org.apache.paimon.utils.ChainTableUtils; +import org.apache.paimon.utils.StringUtils; import java.util.Collections; import java.util.LinkedHashMap; @@ -187,7 +189,7 @@ private static Table normalizeManifestParallelism( pair.other(), safeBound, materializeAbsent); return main == pair.wrapped() && other == pair.other() ? table : new FallbackReadFileStoreTable( - main, other, wrappedBranchHasReadPriority(pair)); + main, other, isWrappedFirst(pair)); } if (table instanceof DelegatedFileStoreTable) { @@ -230,17 +232,28 @@ private static Table normalizeManifestParallelism( return ((FileStoreTable) table).copyWithoutTimeTravel(cap); } + static boolean isWrappedFirst(FallbackReadFileStoreTable table) { + Map options = table.options(); + // Keep this in the exact order used by Paimon 1.4.2 FileStoreTableFactory. The same + // wrapper represents chain, fallback and primary reads, but the SDK exposes no accessor + // for its private wrappedFirst flag. Remove this inference when Paimon exposes one. + if (ChainTableUtils.isChainTable(options)) { + return true; + } + if (!StringUtils.isNullOrWhitespaceOnly( + options.get(CoreOptions.SCAN_FALLBACK_BRANCH.key()))) { + return true; + } + return StringUtils.isNullOrWhitespaceOnly( + options.get(CoreOptions.SCAN_PRIMARY_BRANCH.key())); + } + private static FileStoreTable normalizeManifestParallelism( FileStoreTable table, int safeBound, boolean materializeAbsent) { return (FileStoreTable) normalizeManifestParallelism( (Table) table, safeBound, materializeAbsent); } - private static boolean wrappedBranchHasReadPriority(FallbackReadFileStoreTable table) { - // Paimon's factory sets wrappedFirst=false only for scan.primary-branch tables. - return !table.wrapped().options().containsKey(CoreOptions.SCAN_PRIMARY_BRANCH.key()); - } - public static void validateEffectiveTable(Table table) { validateEffectiveTableOptions(table.options()); if (table instanceof FallbackReadFileStoreTable) { diff --git a/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/PaimonScanParams.java b/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/PaimonScanParams.java index 9dda81ee1c1b45..7a1b9718c00192 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/PaimonScanParams.java +++ b/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/PaimonScanParams.java @@ -110,7 +110,7 @@ public static void validateOptions(Map options) { String scanMode = options.get(CoreOptions.SCAN_MODE.key()); if ("from-creation-timestamp".equalsIgnoreCase(scanMode) && options.get(CoreOptions.SCAN_CREATION_TIME_MILLIS.key()) == null) { - // Paimon 1.3.1 does not validate this newer mode, but its starting scanner + // Paimon 1.4.2 does not validate this mode, but its starting scanner // requires the creation timestamp and otherwise fails after analysis. throw new IllegalArgumentException("Paimon scan mode 'from-creation-timestamp' requires query option '" + CoreOptions.SCAN_CREATION_TIME_MILLIS.key() + "'."); diff --git a/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/PaimonTransaction.java b/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/PaimonTransaction.java index b2416b853e3aac..ddf58b0cbdedd6 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/PaimonTransaction.java +++ b/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/PaimonTransaction.java @@ -26,9 +26,7 @@ import com.google.common.collect.Lists; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; -import org.apache.paimon.CoreOptions; import org.apache.paimon.io.DataInputDeserializer; -import org.apache.paimon.table.FileStoreTable; import org.apache.paimon.table.sink.CommitMessage; import org.apache.paimon.table.sink.CommitMessageSerializer; import org.apache.paimon.table.sink.InnerTableCommit; @@ -234,7 +232,7 @@ int getPayloadCount() { private void doCommit(PaimonWriteBinding writeBinding, List messages) throws Exception { ops.dorisCatalog.getExecutionAuthenticator().execute(() -> { - InnerTableCommit committer = getCommitTable(writeBinding).newCommit(commitUser); + InnerTableCommit committer = writeBinding.getTable().newCommit(commitUser); Exception commitFailure = null; try { if (writeBinding.isOverwrite()) { @@ -332,15 +330,6 @@ private void doAbort(PaimonWriteBinding writeBinding, List messag }); } - private FileStoreTable getCommitTable(PaimonWriteBinding writeBinding) { - FileStoreTable paimonTable = writeBinding.getTable(); - if (!writeBinding.isOverwrite() || writeBinding.getStaticPartition().isEmpty()) { - return paimonTable; - } - return paimonTable.copy(Collections.singletonMap( - CoreOptions.DYNAMIC_PARTITION_OVERWRITE.key(), Boolean.FALSE.toString())); - } - private synchronized void markPreparedTransactionCommitted() { Preconditions.checkState(state == CommitState.PREPARED, "Only a prepared Paimon transaction can complete without a commit"); diff --git a/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/PaimonVariantWriteAnalyzer.java b/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/PaimonVariantWriteAnalyzer.java new file mode 100644 index 00000000000000..ad2e0e12e759b4 --- /dev/null +++ b/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/PaimonVariantWriteAnalyzer.java @@ -0,0 +1,155 @@ +// 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. + +package org.apache.doris.datasource.paimon; + +import org.apache.doris.catalog.Column; +import org.apache.doris.catalog.Type; +import org.apache.doris.nereids.exceptions.AnalysisException; +import org.apache.doris.nereids.exceptions.UnboundException; +import org.apache.doris.nereids.trees.expressions.NamedExpression; +import org.apache.doris.nereids.types.ArrayType; +import org.apache.doris.nereids.types.DataType; +import org.apache.doris.nereids.types.MapType; +import org.apache.doris.nereids.types.StructField; +import org.apache.doris.nereids.types.StructType; +import org.apache.doris.nereids.types.VariantType; + +import java.util.List; +import java.util.Map; +import java.util.Optional; + +/** Analysis checks for the V2-only Paimon Variant write protocol. */ +public final class PaimonVariantWriteAnalyzer { + private PaimonVariantWriteAnalyzer() { + } + + /** + * Rejects a disabled V2 protocol and unsupported Variant V2 conversions before sink coercion. + */ + public static void validate( + PaimonWriteTarget writeTarget, + List writeColumns, + Map columnToOutput, + boolean enableVariantV2) throws AnalysisException { + for (Column column : writeColumns) { + Type targetCatalogType = writeTarget.getColumnTypes().get(column.getName()); + if (targetCatalogType == null) { + continue; + } + DataType targetType = DataType.fromCatalogType(targetCatalogType); + if (!VariantType.containsVariant(targetType)) { + continue; + } + if (!enableVariantV2) { + throw new AnalysisException( + "Paimon VARIANT write only supports Variant V2; " + + "set enable_variant_v2=true"); + } + NamedExpression output = columnToOutput.get(column.getName()); + if (output != null) { + validateVariantConversion(output.getDataType(), targetType, column.getName()); + } + } + } + + /** + * Selects the target used before an inline table computes a common type for each VALUES + * column. An empty result defers coercion until sink binding can validate the resolved source. + */ + public static Optional resolveInlineCoercionTarget( + DataType targetType, NamedExpression value, boolean enableVariantV2) { + if (!VariantType.containsVariant(targetType)) { + return Optional.of(targetType); + } + if (!enableVariantV2) { + return Optional.empty(); + } + try { + if (VariantType.containsVariant(value.getDataType())) { + // Preserve the source layout so final analysis can distinguish V1 from V2. + return Optional.empty(); + } + } catch (UnboundException ignored) { + // Expression analysis will resolve this source after the target cast is attached. + } + // Preserve each non-Variant VALUES row before common-type coercion. Otherwise (1), ('x') + // would first become STRING and the integer would be encoded as a Variant string. + return Optional.of(VariantType.toComputeV2(targetType)); + } + + private static void validateVariantConversion( + DataType sourceType, DataType targetType, String path) throws AnalysisException { + if (!VariantType.containsVariant(targetType)) { + return; + } + if (targetType instanceof VariantType) { + validateVariantSource(sourceType, path); + return; + } + if (sourceType instanceof ArrayType && targetType instanceof ArrayType) { + validateVariantConversion( + ((ArrayType) sourceType).getItemType(), + ((ArrayType) targetType).getItemType(), + path + "[]"); + return; + } + if (sourceType instanceof MapType && targetType instanceof MapType) { + MapType sourceMap = (MapType) sourceType; + MapType targetMap = (MapType) targetType; + validateVariantConversion( + sourceMap.getKeyType(), targetMap.getKeyType(), path + ".key"); + validateVariantConversion( + sourceMap.getValueType(), targetMap.getValueType(), path + ".value"); + return; + } + if (sourceType instanceof StructType && targetType instanceof StructType) { + List sourceFields = ((StructType) sourceType).getFields(); + List targetFields = ((StructType) targetType).getFields(); + int fieldCount = Math.min(sourceFields.size(), targetFields.size()); + for (int i = 0; i < fieldCount; i++) { + validateVariantConversion( + sourceFields.get(i).getDataType(), + targetFields.get(i).getDataType(), + path + "." + targetFields.get(i).getName()); + } + return; + } + + // A shape-changing cast can bypass the matching container branches above. Validate its + // complete source against the leaf conversion contract before sink coercion adds a cast. + validateVariantSource(sourceType, path); + } + + private static void validateVariantSource( + DataType sourceType, String path) throws AnalysisException { + if (VariantType.isLegacyVariant(sourceType)) { + throw new AnalysisException( + "Paimon VARIANT write only supports Variant V2, but input column '" + + path + "' is Variant V1"); + } + if (sourceType instanceof ArrayType) { + validateVariantSource(((ArrayType) sourceType).getItemType(), path + "[]"); + return; + } + if (!VariantType.isSupportedComputeV2CastSource(sourceType)) { + throw new AnalysisException( + "Paimon VARIANT write cannot convert input column '" + path + + "' from " + sourceType.toSql() + " to Variant V2"); + } + } +} diff --git a/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/PaimonWriteBinding.java b/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/PaimonWriteBinding.java index cddb2864a70443..bbde84f0de1b73 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/PaimonWriteBinding.java +++ b/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/PaimonWriteBinding.java @@ -79,6 +79,7 @@ public static PaimonWriteBinding create(PaimonWriteTarget writeTarget, writeTarget.getColumnTypes(), typedStaticPartition, context.isOverwrite()); + table = configureTableForWrite(table, context.isOverwrite(), staticPartition); return new PaimonWriteBinding( dorisTable, table, @@ -87,6 +88,27 @@ public static PaimonWriteBinding create(PaimonWriteTarget writeTarget, staticPartition); } + static FileStoreTable configureTableForWrite(FileStoreTable table, boolean overwrite, + Map staticPartition) { + if (!overwrite) { + return table; + } + + String dynamicOverwriteKey = CoreOptions.DYNAMIC_PARTITION_OVERWRITE.key(); + boolean explicitlyDynamic = staticPartition.isEmpty() + && Boolean.parseBoolean(table.options().get(dynamicOverwriteKey)); + if (explicitlyDynamic) { + return table; + } + + // Doris INSERT OVERWRITE without PARTITION replaces the whole table by default, + // whereas Paimon's default is dynamic partition overwrite. Read the raw table option + // above so only an explicit "true" opts into Paimon's dynamic behavior. A static + // PARTITION clause always requires static overwrite, regardless of the table option. + return table.copy(Collections.singletonMap( + dynamicOverwriteKey, Boolean.FALSE.toString())); + } + public FileStoreTable getTable() { return table; } diff --git a/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/PaimonWriteTarget.java b/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/PaimonWriteTarget.java index 86fe665a28a065..4ed18bd7c94cda 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/PaimonWriteTarget.java +++ b/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/PaimonWriteTarget.java @@ -20,6 +20,8 @@ import org.apache.doris.catalog.Column; import org.apache.doris.catalog.Type; import org.apache.doris.common.AnalysisException; +import org.apache.doris.nereids.types.DataType; +import org.apache.doris.nereids.types.VariantType; import com.google.common.collect.ImmutableList; import org.apache.paimon.table.FileStoreTable; @@ -66,6 +68,10 @@ private PaimonWriteTarget(PaimonExternalTable dorisTable, FileStoreTable table) } Type type = PaimonUtil.paimonTypeToDorisType( field.type(), catalog.getEnableMappingVarbinary(), false); + // A Paimon schema has one logical VARIANT type. Doris uses the compute-V2 + // representation for the Paimon write protocol, including nested VARIANT nodes. + DataType writeType = VariantType.toComputeV2(DataType.fromCatalogType(type)); + type = writeType.toCatalogDataType(); // Doris exposes external-table columns as nullable. The real Paimon nullability and // defaults remain in the pinned FileStoreTable and are enforced by the Paimon writer. Column column = new Column(field.name(), type, true, null, true, diff --git a/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/source/PaimonScanNode.java b/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/source/PaimonScanNode.java index bcf5be93649c68..927ddeeba2882f 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/source/PaimonScanNode.java +++ b/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/source/PaimonScanNode.java @@ -538,6 +538,8 @@ public List getSplits(int numBackends) throws UserException { } Optional> optRawFiles = dataSplit.convertToRawFiles(); Optional> optDeletionFiles = dataSplit.deletionFiles(); + // Only evaluate merged row counts when COUNT(*) pushdown is active; ordinary scans + // must not pay the planning cost of merging Paimon manifest statistics. OptionalLong mergedRowCount = applyCountPushdown ? dataSplit.mergedRowCount() : OptionalLong.empty(); if (applyCountPushdown && mergedRowCount.isPresent()) { diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/BindSink.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/BindSink.java index e0e92e9ff31bcf..ed55f8204280d7 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/BindSink.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/BindSink.java @@ -46,6 +46,7 @@ import org.apache.doris.datasource.mvcc.MvccSnapshot; import org.apache.doris.datasource.paimon.PaimonExternalDatabase; import org.apache.doris.datasource.paimon.PaimonExternalTable; +import org.apache.doris.datasource.paimon.PaimonVariantWriteAnalyzer; import org.apache.doris.datasource.paimon.PaimonWriteTarget; import org.apache.doris.dictionary.Dictionary; import org.apache.doris.nereids.CascadesContext; @@ -973,7 +974,7 @@ private Plan bindPaimonTableSink(MatchingContext> c if (bindColumns.size() != child.getOutput().size()) { throw new AnalysisException("insert into cols should be corresponding to the query output"); } - Map columnToOutput = getJdbcColumnToOutput(bindColumns, child); + Map columnToOutput = getPaimonColumnToOutput(bindColumns, child); List writeColumns = new ArrayList<>(bindColumns); if (!staticPartitionColNames.isEmpty()) { for (Column column : writeTarget.getSchema()) { @@ -993,6 +994,11 @@ private Plan bindPaimonTableSink(MatchingContext> c .map(NamedExpression.class::cast) .collect(ImmutableList.toImmutableList()), sink.getDMLCommandType(), Optional.empty(), Optional.empty(), child); + ConnectContext connectContext = ctx.cascadesContext.getConnectContext(); + boolean enableVariantV2 = connectContext != null + && connectContext.getSessionVariable().isEnableVariantV2(); + PaimonVariantWriteAnalyzer.validate( + writeTarget, writeColumns, columnToOutput, enableVariantV2); LogicalProject outputProject = getOutputProjectByCoercion( writeColumns, child, columnToOutput, writeTarget.getColumnTypes()); return boundSink.withChildAndUpdateOutput(outputProject); @@ -1103,6 +1109,18 @@ private static Map getJdbcColumnToOutput( return columnToOutput; } + private static Map getPaimonColumnToOutput( + List bindColumns, LogicalPlan child) { + Map columnToOutput = + Maps.newTreeMap(String.CASE_INSENSITIVE_ORDER); + for (int i = 0; i < bindColumns.size(); i++) { + Column column = bindColumns.get(i); + columnToOutput.put(column.getName(), + new Alias(child.getOutput().get(i), column.getName())); + } + return columnToOutput; + } + private Plan bindDictionarySink(MatchingContext> ctx) { UnboundDictionarySink sink = ctx.root; Pair pair = bind(ctx.cascadesContext, sink); diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/expression/check/CheckCast.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/expression/check/CheckCast.java index 3b975b60055369..749c8aa247ba15 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/expression/check/CheckCast.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/expression/check/CheckCast.java @@ -355,8 +355,11 @@ public static boolean check(DataType originalType, DataType targetType, boolean */ public static boolean check(DataType originalType, DataType targetType, boolean isStrictMode, boolean looseAggState) { + if (targetType instanceof VariantType && ((VariantType) targetType).isComputeV2()) { + return VariantType.isSupportedComputeV2CastSource(originalType); + } if (originalType.isVariantType() && targetType.isVariantType()) { - return originalType.equals(targetType); + return ((VariantType) originalType).isCastCompatibleWith((VariantType) targetType); } if (originalType.isVariantType() && (targetType instanceof PrimitiveType || targetType.isArrayType())) { // variant could cast to primitive types and array diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/expressions/functions/scalar/Array.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/expressions/functions/scalar/Array.java index 8da0b6f3b45778..8170a5cd56bbb1 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/expressions/functions/scalar/Array.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/expressions/functions/scalar/Array.java @@ -26,6 +26,7 @@ import org.apache.doris.nereids.trees.expressions.visitor.ExpressionVisitor; import org.apache.doris.nereids.types.ArrayType; import org.apache.doris.nereids.types.DataType; +import org.apache.doris.nereids.types.VariantType; import org.apache.doris.nereids.types.coercion.FollowToArgumentType; import org.apache.doris.nereids.util.TypeCoercionUtils; @@ -66,12 +67,14 @@ private Array(ScalarFunctionParams functionParams) { @Override public void checkLegalityBeforeTypeCoercion() { - if (children.isEmpty()) { + if (arity() == 0) { return; } - DataType firstChildType = getArgument(0).getDataType(); - if (firstChildType.isJsonType() || firstChildType.isVariantType()) { - throw new AnalysisException("array does not support jsonb/variant type"); + for (Expression argument : getArguments()) { + DataType childType = argument.getDataType(); + if (childType.isJsonType() || VariantType.isLegacyVariant(childType)) { + throw new AnalysisException("array does not support jsonb/variant type"); + } } } diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/expressions/functions/scalar/CreateMap.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/expressions/functions/scalar/CreateMap.java index 6bca4dbda42557..d00cfc8e27c57c 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/expressions/functions/scalar/CreateMap.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/expressions/functions/scalar/CreateMap.java @@ -26,6 +26,7 @@ import org.apache.doris.nereids.types.ArrayType; import org.apache.doris.nereids.types.DataType; import org.apache.doris.nereids.types.MapType; +import org.apache.doris.nereids.types.VariantType; import org.apache.doris.nereids.util.TypeCoercionUtils; import com.google.common.collect.ImmutableList; @@ -85,11 +86,15 @@ public void checkLegalityBeforeTypeCoercion() { if (arity() % 2 != 0) { throw new AnalysisException("map can't be odd parameters, need even parameters " + this.toSql()); } - children.forEach(child -> { - if (child.getDataType().isJsonType() || child.getDataType().isVariantType()) { + for (int i = 0; i < arity(); i++) { + DataType childType = getArgument(i).getDataType(); + boolean isKey = i % 2 == 0; + if (childType.isJsonType() + || (isKey && childType.isVariantType()) + || (!isKey && VariantType.isLegacyVariant(childType))) { throw new AnalysisException("map does not support jsonb/variant type"); } - }); + } } /** diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/expressions/functions/scalar/CreateNamedStruct.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/expressions/functions/scalar/CreateNamedStruct.java index b0ab4725908f5d..90399ac162817c 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/expressions/functions/scalar/CreateNamedStruct.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/expressions/functions/scalar/CreateNamedStruct.java @@ -29,6 +29,7 @@ import org.apache.doris.nereids.types.DataType; import org.apache.doris.nereids.types.StructField; import org.apache.doris.nereids.types.StructType; +import org.apache.doris.nereids.types.VariantType; import com.google.common.collect.ImmutableList; import com.google.common.collect.Sets; @@ -65,11 +66,12 @@ public void checkLegalityBeforeTypeCoercion() { } Set names = Sets.newHashSet(); for (int i = 0; i < arity(); i = i + 2) { - if (!(child(i) instanceof StringLikeLiteral)) { + Expression nameArgument = getArgument(i); + if (!(nameArgument instanceof StringLikeLiteral)) { throw new AnalysisException("named_struct only allows" + " constant string parameter in odd position: " + this); } else { - String name = ((StringLikeLiteral) child(i)).getStringValue().toLowerCase(); + String name = ((StringLikeLiteral) nameArgument).getStringValue().toLowerCase(); if (names.contains(name)) { throw new AnalysisException("The name of the struct field cannot be repeated." + " same name fields are " + name); @@ -77,8 +79,8 @@ public void checkLegalityBeforeTypeCoercion() { names.add(name); } } - // i+1 is value, check if it is not jsonb/variant type - if (child(i + 1).getDataType().isJsonType() || child(i + 1).getDataType().isVariantType()) { + DataType valueType = getArgument(i + 1).getDataType(); + if (valueType.isJsonType() || VariantType.isLegacyVariant(valueType)) { throw new AnalysisException("named_struct does not support jsonb/variant type"); } } diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/expressions/functions/scalar/CreateStruct.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/expressions/functions/scalar/CreateStruct.java index 9e89da0fd87319..ed983fac1675b8 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/expressions/functions/scalar/CreateStruct.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/expressions/functions/scalar/CreateStruct.java @@ -27,6 +27,7 @@ import org.apache.doris.nereids.trees.expressions.visitor.ExpressionVisitor; import org.apache.doris.nereids.types.DataType; import org.apache.doris.nereids.types.StructType; +import org.apache.doris.nereids.types.VariantType; import com.google.common.collect.ImmutableList; @@ -59,9 +60,9 @@ public void checkLegalityBeforeTypeCoercion() { if (arity() == 0) { throw new AnalysisException("struct requires at least one argument, like: struct(1)"); } - // for all field we do not support struct field with jsonb/variant type - children.forEach(child -> { - if (child.getDataType().isJsonType() || child.getDataType().isVariantType()) { + getArguments().forEach(argument -> { + if (argument.getDataType().isJsonType() + || VariantType.isLegacyVariant(argument.getDataType())) { throw new AnalysisException("struct does not support jsonb/variant type"); } }); diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/InsertUtils.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/InsertUtils.java index 6d4df3ca89f6a2..528198c2b70fee 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/InsertUtils.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/InsertUtils.java @@ -30,6 +30,7 @@ import org.apache.doris.datasource.hive.HMSExternalTable; import org.apache.doris.datasource.jdbc.JdbcExternalTable; import org.apache.doris.datasource.mvcc.MvccTable; +import org.apache.doris.datasource.paimon.PaimonVariantWriteAnalyzer; import org.apache.doris.foundation.format.FormatOptions; import org.apache.doris.nereids.CascadesContext; import org.apache.doris.nereids.StatementContext; @@ -407,6 +408,9 @@ private static Plan normalizePlanWithoutLock(LogicalPlan plan, TableIf table, } ConnectContext context = ConnectContext.get(); + boolean enableVariantV2 = context != null + && context.getSessionVariable().isEnableVariantV2(); + boolean isPaimonSink = unboundLogicalSink instanceof UnboundPaimonTableSink; ExpressionRewriteContext rewriteContext = null; if (context != null && context.getStatementContext() != null) { rewriteContext = new ExpressionRewriteContext( @@ -460,7 +464,8 @@ private static Plan normalizePlanWithoutLock(LogicalPlan plan, TableIf table, addColumnValue(analyzer, optimizedRowConstructor, defaultExpression, null, rewriteContext, strictCast); } else { - DataType targetType = DataType.fromCatalogType(sameNameColumn.getType()); + DataType targetType = targetTypeForInlineValue( + sameNameColumn, values.get(i), isPaimonSink, enableVariantV2); addColumnValue(analyzer, optimizedRowConstructor, values.get(i), targetType, rewriteContext, strictCast); } @@ -481,7 +486,8 @@ private static Plan normalizePlanWithoutLock(LogicalPlan plan, TableIf table, addColumnValue(analyzer, optimizedRowConstructor, defaultExpression, null, rewriteContext, strictCast); } else { - DataType targetType = DataType.fromCatalogType(columns.get(i).getType()); + DataType targetType = targetTypeForInlineValue( + columns.get(i), values.get(i), isPaimonSink, enableVariantV2); addColumnValue(analyzer, optimizedRowConstructor, values.get(i), targetType, rewriteContext, strictCast); } @@ -493,6 +499,15 @@ private static Plan normalizePlanWithoutLock(LogicalPlan plan, TableIf table, return plan.withChildren(new LogicalInlineTable(optimizedRowConstructors.build())); } + private static DataType targetTypeForInlineValue( + Column column, NamedExpression value, boolean isPaimonSink, boolean enableVariantV2) { + DataType targetType = DataType.fromCatalogType(column.getType()); + return isPaimonSink + ? PaimonVariantWriteAnalyzer.resolveInlineCoercionTarget( + targetType, value, enableVariantV2).orElse(null) + : targetType; + } + /** buildAnalyzer */ public static ExpressionAnalyzer buildExprAnalyzer(Plan plan, CascadesContext analyzeContext) { return new ExpressionAnalyzer(plan, new Scope(ImmutableList.of()), diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/types/VariantType.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/types/VariantType.java index ae33c47040f026..c5280ef24d7aaf 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/types/VariantType.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/types/VariantType.java @@ -336,6 +336,42 @@ public static boolean containsVariant(DataType dataType) { return false; } + /** Whether this is a legacy Variant leaf rather than the compute-only V2 representation. */ + public static boolean isLegacyVariant(DataType dataType) { + return dataType instanceof VariantType && !((VariantType) dataType).isComputeV2(); + } + + /** + * Whether the Variant V2 execution kernel can convert this source type. + * + *

This mirrors the BE {@code execute_to_variant} contract: encoded JSON, nested arrays, + * compute V2 values, and the scalar types supported by the typed Variant representation. + * MAP, STRUCT, TIMEV2 and DECIMAL256 are intentionally excluded until their BE conversions + * are implemented.

+ */ + public static boolean isSupportedComputeV2CastSource(DataType dataType) { + if (dataType.isNullType() || dataType.isJsonType()) { + return true; + } + if (dataType instanceof VariantType) { + return ((VariantType) dataType).isComputeV2(); + } + if (dataType instanceof ArrayType) { + return isSupportedComputeV2CastSource(((ArrayType) dataType).getItemType()); + } + if (dataType instanceof DecimalV3Type) { + return ((DecimalV3Type) dataType).getPrecision() + <= DecimalV3Type.MAX_DECIMAL128_PRECISION; + } + return dataType.isBooleanType() + || dataType.isIntegralType() + || dataType.isFloatLikeType() + || dataType.isDecimalV2Type() + || dataType.isDateLikeType() + || dataType.isStringLikeType() + || dataType.isIPType(); + } + /** Selects the compute-only Variant representation in a possibly nested type. */ public static DataType toComputeV2(DataType dataType) { if (dataType instanceof VariantType) { @@ -359,10 +395,21 @@ public static DataType toComputeV2(DataType dataType) { *

Legacy Variant values retain their existing common-type behavior. Compute-only Variant * V2 values share one physical representation, independent of source layout properties.

*/ - public boolean isExecutionCompatibleWith(VariantType other) { + public boolean hasCommonExecutionTypeWith(VariantType other) { return computeV2 == other.computeV2; } + /** + * Whether a cast between two Variant types is safe. + * + *

Variant V1 embeds layout properties in its execution type, so V1 casts still require + * exact type equality. All compute-only V2 types share the same physical value/metadata + * representation, therefore layout-property differences do not require conversion.

+ */ + public boolean isCastCompatibleWith(VariantType other) { + return (computeV2 && other.computeV2) || equals(other); + } + /** Returns this Variant type with the requested compute-only physical representation. */ public VariantType withComputeV2(boolean enabled) { if (computeV2 == enabled) { diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/util/TypeCoercionUtils.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/util/TypeCoercionUtils.java index af18d291cc6dc5..196762604b36cd 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/util/TypeCoercionUtils.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/util/TypeCoercionUtils.java @@ -479,7 +479,7 @@ private static boolean matchesType(DataType input, DataType target) { public static Expression castIfNotSameType(Expression input, DataType targetType) { if (input.isNullLiteral()) { return new NullLiteral(targetType); - } else if (input.getDataType().equals(targetType) + } else if (isNoOpCastCompatible(input.getDataType(), targetType) || (input.getDataType().isStringLikeType()) && targetType.isStringLikeType()) { return input; } else { @@ -488,6 +488,45 @@ public static Expression castIfNotSameType(Expression input, DataType targetType } } + /** Whether two types share an execution layout and therefore need no runtime cast. */ + public static boolean isNoOpCastCompatible(DataType left, DataType right) { + if (left.equals(right)) { + return true; + } + if (left instanceof VariantType && right instanceof VariantType) { + return ((VariantType) left).isCastCompatibleWith((VariantType) right); + } + if (left instanceof ArrayType && right instanceof ArrayType) { + return isNoOpCastCompatible( + ((ArrayType) left).getItemType(), ((ArrayType) right).getItemType()); + } + if (left instanceof MapType && right instanceof MapType) { + MapType leftMap = (MapType) left; + MapType rightMap = (MapType) right; + return isNoOpCastCompatible(leftMap.getKeyType(), rightMap.getKeyType()) + && isNoOpCastCompatible(leftMap.getValueType(), rightMap.getValueType()); + } + if (left instanceof StructType && right instanceof StructType) { + List leftFields = ((StructType) left).getFields(); + List rightFields = ((StructType) right).getFields(); + if (leftFields.size() != rightFields.size()) { + return false; + } + for (int i = 0; i < leftFields.size(); i++) { + StructField leftField = leftFields.get(i); + StructField rightField = rightFields.get(i); + if (leftField.isNullable() != rightField.isNullable() + || !leftField.getName().equals(rightField.getName()) + || !isNoOpCastCompatible( + leftField.getDataType(), rightField.getDataType())) { + return false; + } + } + return true; + } + return false; + } + /** * Wrap {@code expression} in a cast to {@code targetType} when the source type can * already be resolved (Literal, or any expression whose {@code getDataType()} does @@ -1219,7 +1258,7 @@ private static Optional findWiderComplexTypeForTwo( } private static Optional findCommonVariantType(VariantType left, VariantType right) { - if (!left.isExecutionCompatibleWith(right)) { + if (!left.hasCommonExecutionTypeWith(right)) { return Optional.empty(); } return Optional.of(left.isComputeV2() ? VariantType.COMPUTE_V2_INSTANCE : right); diff --git a/fe/fe-core/src/test/java/org/apache/doris/datasource/paimon/PaimonReaderOptionsTest.java b/fe/fe-core/src/test/java/org/apache/doris/datasource/paimon/PaimonReaderOptionsTest.java index f6856c78e2f7da..aa1d6415adedd8 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/datasource/paimon/PaimonReaderOptionsTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/datasource/paimon/PaimonReaderOptionsTest.java @@ -164,6 +164,23 @@ void testRuntimeCapPreservesPrimaryBranchPriority() throws Exception { Assertions.assertFalse(wrappedFirst.getBoolean(safe)); } + @Test + void testFallbackWrapperOrderMatchesPaimonFactoryPrecedence() { + Assertions.assertTrue(PaimonReaderOptions.isWrappedFirst(fallbackTable(ImmutableMap.of( + CoreOptions.SCAN_FALLBACK_BRANCH.key(), "fallback", + CoreOptions.SCAN_PRIMARY_BRANCH.key(), "primary")))); + Assertions.assertFalse(PaimonReaderOptions.isWrappedFirst(fallbackTable(ImmutableMap.of( + CoreOptions.SCAN_PRIMARY_BRANCH.key(), "primary")))); + Assertions.assertTrue(PaimonReaderOptions.isWrappedFirst(fallbackTable(ImmutableMap.of( + CoreOptions.SCAN_PRIMARY_BRANCH.key(), " ")))); + Assertions.assertTrue(PaimonReaderOptions.isWrappedFirst(fallbackTable(ImmutableMap.of( + CoreOptions.CHAIN_TABLE_ENABLED.key(), "true", + CoreOptions.SCAN_PRIMARY_BRANCH.key(), "primary")))); + // Old FE payloads only represented fallback reads and may not carry either selector. + Assertions.assertTrue(PaimonReaderOptions.isWrappedFirst( + fallbackTable(Collections.emptyMap()))); + } + @Test void testSystemSourceKeepsFallbackAsOutermostPlanningDecorator() { FileStoreTable main = newFileStoreTable("main", Collections.emptyMap()); @@ -291,4 +308,10 @@ private FileStoreTable newFileStoreTable(String name, Map option return new AppendOnlyFileStoreTable( Mockito.mock(FileIO.class), new Path("memory://" + name), schema, CatalogEnvironment.empty()); } + + private FallbackReadFileStoreTable fallbackTable(Map options) { + return new FallbackReadFileStoreTable( + newFileStoreTable("order_main", options), + newFileStoreTable("order_other", Collections.emptyMap()), true); + } } diff --git a/fe/fe-core/src/test/java/org/apache/doris/datasource/paimon/PaimonVariantWriteAnalyzerTest.java b/fe/fe-core/src/test/java/org/apache/doris/datasource/paimon/PaimonVariantWriteAnalyzerTest.java new file mode 100644 index 00000000000000..75437cfbb3b3ff --- /dev/null +++ b/fe/fe-core/src/test/java/org/apache/doris/datasource/paimon/PaimonVariantWriteAnalyzerTest.java @@ -0,0 +1,201 @@ +// 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. + +package org.apache.doris.datasource.paimon; + +import org.apache.doris.catalog.Column; +import org.apache.doris.nereids.exceptions.AnalysisException; +import org.apache.doris.nereids.trees.expressions.Alias; +import org.apache.doris.nereids.trees.expressions.Cast; +import org.apache.doris.nereids.trees.expressions.NamedExpression; +import org.apache.doris.nereids.trees.expressions.literal.IntegerLiteral; +import org.apache.doris.nereids.trees.expressions.literal.StringLiteral; +import org.apache.doris.nereids.types.ArrayType; +import org.apache.doris.nereids.types.DataType; +import org.apache.doris.nereids.types.IntegerType; +import org.apache.doris.nereids.types.MapType; +import org.apache.doris.nereids.types.StringType; +import org.apache.doris.nereids.types.TimeV2Type; +import org.apache.doris.nereids.types.VariantType; + +import org.apache.paimon.table.FileStoreTable; +import org.apache.paimon.types.DataField; +import org.apache.paimon.types.DataTypes; +import org.junit.Assert; +import org.junit.Test; +import org.mockito.Mockito; + +import java.util.Collections; +import java.util.Map; +import java.util.TreeMap; + +public class PaimonVariantWriteAnalyzerTest { + + @Test + public void testDisabledVariantV2IsRejectedDuringAnalysis() throws Exception { + PaimonWriteTarget target = createTarget( + DataTypes.FIELD(0, "payload", DataTypes.VARIANT())); + + AnalysisException exception = Assert.assertThrows( + AnalysisException.class, + () -> validate(target, VariantType.COMPUTE_V2_INSTANCE, false)); + Assert.assertTrue(exception.getMessage().contains("enable_variant_v2=true")); + } + + @Test + public void testLegacyVariantInputIsRejectedDuringAnalysis() throws Exception { + PaimonWriteTarget target = createTarget( + DataTypes.FIELD(0, "payload", DataTypes.VARIANT())); + + AnalysisException exception = Assert.assertThrows( + AnalysisException.class, + () -> validate(target, VariantType.INSTANCE, true)); + Assert.assertTrue(exception.getMessage().contains("Variant V1")); + Assert.assertTrue(exception.getMessage().contains("payload")); + } + + @Test + public void testComputeV2InputIsAccepted() throws Exception { + PaimonWriteTarget target = createTarget( + DataTypes.FIELD(0, "payload", DataTypes.VARIANT())); + + PaimonVariantWriteAnalyzer.validate( + target, + Collections.singletonList(target.getColumn("payload")), + outputs("payload", VariantType.COMPUTE_V2_INSTANCE), + true); + } + + @Test + public void testInlineCoercionPreservesValuesBeforeCommonTypeResolution() { + Alias integerValue = new Alias(new IntegerLiteral(1), "payload"); + Assert.assertEquals( + VariantType.COMPUTE_V2_INSTANCE, + PaimonVariantWriteAnalyzer.resolveInlineCoercionTarget( + VariantType.INSTANCE, integerValue, true).get()); + Assert.assertFalse(PaimonVariantWriteAnalyzer.resolveInlineCoercionTarget( + VariantType.INSTANCE, integerValue, false).isPresent()); + + Alias variantValue = new Alias( + new Cast(new StringLiteral("{}"), VariantType.INSTANCE), "payload"); + Assert.assertFalse(PaimonVariantWriteAnalyzer.resolveInlineCoercionTarget( + VariantType.INSTANCE, variantValue, true).isPresent()); + + Assert.assertEquals( + ArrayType.of(VariantType.COMPUTE_V2_INSTANCE), + PaimonVariantWriteAnalyzer.resolveInlineCoercionTarget( + ArrayType.of(VariantType.INSTANCE), integerValue, true).get()); + } + + @Test + public void testNestedLegacyVariantInputIsRejected() throws Exception { + PaimonWriteTarget target = createTarget( + DataTypes.FIELD(0, "payload", DataTypes.ARRAY(DataTypes.VARIANT()))); + + AnalysisException exception = Assert.assertThrows( + AnalysisException.class, + () -> validate(target, ArrayType.of(VariantType.INSTANCE), true)); + Assert.assertTrue(exception.getMessage().contains("payload[]")); + } + + @Test + public void testNonVariantTableDoesNotRequireVariantV2() throws Exception { + PaimonWriteTarget target = createTarget( + DataTypes.FIELD(0, "payload", DataTypes.STRING())); + + validate(target, org.apache.doris.nereids.types.StringType.INSTANCE, false); + } + + @Test + public void testOmittedVariantColumnDoesNotRequireVariantV2() throws Exception { + PaimonWriteTarget target = createTarget( + DataTypes.FIELD(0, "id", DataTypes.INT()), + DataTypes.FIELD(1, "payload", DataTypes.VARIANT())); + Column id = target.getColumn("id"); + + PaimonVariantWriteAnalyzer.validate( + target, + Collections.singletonList(id), + outputs("id", IntegerType.INSTANCE), + false); + } + + @Test + public void testLegacyVariantNestedInShapeChangingSourceIsRejected() throws Exception { + PaimonWriteTarget scalarTarget = createTarget( + DataTypes.FIELD(0, "payload", DataTypes.VARIANT())); + AnalysisException arraySourceException = Assert.assertThrows( + AnalysisException.class, + () -> validate( + scalarTarget, ArrayType.of(VariantType.INSTANCE), true)); + Assert.assertTrue(arraySourceException.getMessage().contains("payload[]")); + + PaimonWriteTarget arrayTarget = createTarget( + DataTypes.FIELD(0, "payload", DataTypes.ARRAY(DataTypes.VARIANT()))); + AnalysisException scalarSourceException = Assert.assertThrows( + AnalysisException.class, + () -> validate(arrayTarget, VariantType.INSTANCE, true)); + Assert.assertTrue(scalarSourceException.getMessage().contains("payload")); + } + + @Test + public void testUnsupportedComputeV2SourcesAreRejectedDuringAnalysis() throws Exception { + PaimonWriteTarget target = createTarget( + DataTypes.FIELD(0, "payload", DataTypes.VARIANT())); + + AnalysisException mapException = Assert.assertThrows( + AnalysisException.class, + () -> validate( + target, MapType.of(StringType.INSTANCE, IntegerType.INSTANCE), true)); + Assert.assertTrue(mapException.getMessage().contains("MAP")); + + AnalysisException timeException = Assert.assertThrows( + AnalysisException.class, + () -> validate(target, TimeV2Type.MAX, true)); + Assert.assertTrue(timeException.getMessage().contains("TIME")); + } + + private static void validate( + PaimonWriteTarget target, DataType sourceType, boolean enableVariantV2) + throws AnalysisException { + Column column = target.getSchema().get(0); + PaimonVariantWriteAnalyzer.validate( + target, + Collections.singletonList(column), + outputs(column.getName(), sourceType), + enableVariantV2); + } + + private static Map outputs(String name, DataType dataType) { + NamedExpression expression = Mockito.mock(NamedExpression.class); + Mockito.when(expression.getDataType()).thenReturn(dataType); + Map outputs = new TreeMap<>(String.CASE_INSENSITIVE_ORDER); + outputs.put(name, expression); + return outputs; + } + + private static PaimonWriteTarget createTarget(DataField... fields) throws Exception { + PaimonExternalTable dorisTable = Mockito.mock(PaimonExternalTable.class); + PaimonExternalCatalog catalog = Mockito.mock(PaimonExternalCatalog.class); + FileStoreTable table = Mockito.mock(FileStoreTable.class); + Mockito.when(dorisTable.getCatalog()).thenReturn(catalog); + Mockito.when(dorisTable.getPaimonTableForWrite()).thenReturn(table); + Mockito.when(table.rowType()).thenReturn(DataTypes.ROW(fields)); + Mockito.when(table.partitionKeys()).thenReturn(Collections.emptyList()); + return PaimonWriteTarget.create(dorisTable); + } +} diff --git a/fe/fe-core/src/test/java/org/apache/doris/datasource/paimon/PaimonWriteBindingTest.java b/fe/fe-core/src/test/java/org/apache/doris/datasource/paimon/PaimonWriteBindingTest.java index f81426fe0484bc..2437ace3a768ba 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/datasource/paimon/PaimonWriteBindingTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/datasource/paimon/PaimonWriteBindingTest.java @@ -26,6 +26,7 @@ import org.apache.doris.nereids.trees.expressions.literal.StringLiteral; import org.apache.doris.qe.ConnectContext; +import org.apache.paimon.CoreOptions; import org.apache.paimon.table.FileStoreTable; import org.apache.paimon.types.DataTypes; import org.apache.paimon.utils.TypeUtils; @@ -42,6 +43,44 @@ public class PaimonWriteBindingTest { + @Test + public void testFullOverwriteUsesStaticSemanticsByDefault() { + FileStoreTable table = Mockito.mock(FileStoreTable.class); + FileStoreTable configuredTable = Mockito.mock(FileStoreTable.class); + Mockito.when(table.options()).thenReturn(Collections.emptyMap()); + Mockito.when(table.copy(Collections.singletonMap( + CoreOptions.DYNAMIC_PARTITION_OVERWRITE.key(), Boolean.FALSE.toString()))) + .thenReturn(configuredTable); + + Assert.assertSame(configuredTable, PaimonWriteBinding.configureTableForWrite( + table, true, Collections.emptyMap())); + } + + @Test + public void testFullOverwriteHonorsExplicitDynamicSemantics() { + FileStoreTable table = Mockito.mock(FileStoreTable.class); + Mockito.when(table.options()).thenReturn(Collections.singletonMap( + CoreOptions.DYNAMIC_PARTITION_OVERWRITE.key(), Boolean.TRUE.toString())); + + Assert.assertSame(table, PaimonWriteBinding.configureTableForWrite( + table, true, Collections.emptyMap())); + Mockito.verify(table, Mockito.never()).copy(Mockito.anyMap()); + } + + @Test + public void testStaticPartitionAlwaysUsesStaticSemantics() { + FileStoreTable table = Mockito.mock(FileStoreTable.class); + FileStoreTable configuredTable = Mockito.mock(FileStoreTable.class); + Mockito.when(table.options()).thenReturn(Collections.singletonMap( + CoreOptions.DYNAMIC_PARTITION_OVERWRITE.key(), Boolean.TRUE.toString())); + Mockito.when(table.copy(Collections.singletonMap( + CoreOptions.DYNAMIC_PARTITION_OVERWRITE.key(), Boolean.FALSE.toString()))) + .thenReturn(configuredTable); + + Assert.assertSame(configuredTable, PaimonWriteBinding.configureTableForWrite( + table, true, Collections.singletonMap("pt", "1"))); + } + @Test public void testStaticPartitionUsesValueAfterTargetTypeCast() throws Exception { diff --git a/fe/fe-core/src/test/java/org/apache/doris/datasource/paimon/PaimonWriteTargetTest.java b/fe/fe-core/src/test/java/org/apache/doris/datasource/paimon/PaimonWriteTargetTest.java index d4ccecf1fb11e6..a8aea73e22fc41 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/datasource/paimon/PaimonWriteTargetTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/datasource/paimon/PaimonWriteTargetTest.java @@ -18,6 +18,10 @@ package org.apache.doris.datasource.paimon; import org.apache.doris.common.AnalysisException; +import org.apache.doris.nereids.types.ArrayType; +import org.apache.doris.nereids.types.DataType; +import org.apache.doris.nereids.types.StructType; +import org.apache.doris.nereids.types.VariantType; import org.apache.paimon.table.FileStoreTable; import org.apache.paimon.types.DataTypes; @@ -34,8 +38,24 @@ public void testVariantColumnIsAvailableToWriteBinding() throws Exception { PaimonWriteTarget target = createTarget( DataTypes.FIELD(0, "payload", DataTypes.VARIANT())); - Assert.assertTrue(target.getColumn("payload").getType().isVariantType()); - Assert.assertTrue(target.getColumnTypes().get("payload").isVariantType()); + DataType type = DataType.fromCatalogType(target.getColumnTypes().get("payload")); + Assert.assertTrue(type instanceof VariantType); + Assert.assertTrue(((VariantType) type).isComputeV2()); + } + + @Test + public void testNestedVariantUsesComputeV2() throws Exception { + PaimonWriteTarget target = createTarget( + DataTypes.FIELD(0, "items", DataTypes.ARRAY(DataTypes.VARIANT())), + DataTypes.FIELD(1, "record", DataTypes.ROW( + DataTypes.FIELD(2, "payload", DataTypes.VARIANT())))); + + ArrayType arrayType = (ArrayType) DataType.fromCatalogType( + target.getColumnTypes().get("items")); + StructType structType = (StructType) DataType.fromCatalogType( + target.getColumnTypes().get("record")); + Assert.assertTrue(((VariantType) arrayType.getItemType()).isComputeV2()); + Assert.assertTrue(((VariantType) structType.getFields().get(0).getDataType()).isComputeV2()); } @Test diff --git a/fe/fe-core/src/test/java/org/apache/doris/datasource/paimon/source/PaimonScanNodeTest.java b/fe/fe-core/src/test/java/org/apache/doris/datasource/paimon/source/PaimonScanNodeTest.java index c933c627de66c2..3407fd4fc03129 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/datasource/paimon/source/PaimonScanNodeTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/datasource/paimon/source/PaimonScanNodeTest.java @@ -148,6 +148,21 @@ public void testSerializedTableCacheKeyIsStablePerScanNode() { Assert.assertNotEquals(firstKey, second.getSerializedTableCacheKey().orElse("")); } + @Test + public void testRegularScanDoesNotComputeMergedRowCount() throws UserException { + PaimonScanNode node = Mockito.spy(newTestNode(new PlanNodeId(1), new TupleId(3), sv)); + node.setSource(mockPaimonSourceWithPartitionKeys(Collections.emptyList())); + DataSplit dataSplit = Mockito.spy(createDataSplit("regular.parquet")); + Mockito.doReturn(Collections.singletonList(dataSplit)).when(node).getPaimonSplitFromAPI(); + Mockito.when(sv.isForceJniScanner()).thenReturn(true); + Mockito.when(sv.getIgnoreSplitType()).thenReturn("NONE"); + + List splits = node.getSplits(1); + + Assert.assertEquals(1, splits.size()); + Mockito.verify(dataSplit, Mockito.never()).mergedRowCount(); + } + @Test public void testCountColumnKeepsAllSplitsWhileCountStarUsesMergedRowCount() throws UserException { PaimonScanNode node = Mockito.spy(newTestNode(new PlanNodeId(1), new TupleId(3), sv)); diff --git a/fe/fe-core/src/test/java/org/apache/doris/nereids/rules/analysis/FunctionRegistryTest.java b/fe/fe-core/src/test/java/org/apache/doris/nereids/rules/analysis/FunctionRegistryTest.java index 5494b003ebecbf..6f8f3a76061a7c 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/nereids/rules/analysis/FunctionRegistryTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/nereids/rules/analysis/FunctionRegistryTest.java @@ -156,11 +156,58 @@ public void testVariantV2SessionSelectsComputeResultType() { return true; }) ); + + PlanChecker.from(connectContext) + .analyze("select array(parse_to_variant('1')), " + + "map('k', parse_to_variant('2')), " + + "struct(parse_to_variant('3')), " + + "named_struct('v', parse_to_variant('4'))") + .matches( + logicalOneRowRelation().when(oneRowRelation -> { + ArrayType array = (ArrayType) oneRowRelation.getProjects().get(0) + .child(0).getDataType(); + MapType map = (MapType) oneRowRelation.getProjects().get(1) + .child(0).getDataType(); + StructType struct = (StructType) oneRowRelation.getProjects().get(2) + .child(0).getDataType(); + StructType namedStruct = (StructType) oneRowRelation.getProjects().get(3) + .child(0).getDataType(); + Assertions.assertTrue(((VariantType) array.getItemType()).isComputeV2()); + Assertions.assertTrue(((VariantType) map.getValueType()).isComputeV2()); + Assertions.assertTrue(((VariantType) struct.getFields().get(0) + .getDataType()).isComputeV2()); + Assertions.assertTrue(((VariantType) namedStruct.getFields().get(0) + .getDataType()).isComputeV2()); + return true; + }) + ); + + AnalysisException mapKeyException = Assertions.assertThrowsExactly( + AnalysisException.class, + () -> PlanChecker.from(connectContext) + .analyze("select map(parse_to_variant('1'), 1)")); + Assertions.assertTrue(mapKeyException.getMessage() + .contains("map does not support jsonb/variant type")); } finally { connectContext.getSessionVariable().enableVariantV2 = false; } } + @Test + public void testLegacyVariantContainerConstructorsRemainRejected() { + List sqls = ImmutableList.of( + "select array(1, parse_to_variant('2'))", + "select map('k', parse_to_variant('2'))", + "select struct(parse_to_variant('3'))", + "select named_struct('v', parse_to_variant('4'))"); + for (String sql : sqls) { + AnalysisException exception = Assertions.assertThrowsExactly( + AnalysisException.class, + () -> PlanChecker.from(connectContext).analyze(sql)); + Assertions.assertTrue(exception.getMessage().contains("jsonb/variant type")); + } + } + @Test public void testOverrideArity() { // the substring function has 2 override functions: diff --git a/fe/fe-core/src/test/java/org/apache/doris/nereids/rules/expression/check/CheckCastTest.java b/fe/fe-core/src/test/java/org/apache/doris/nereids/rules/expression/check/CheckCastTest.java index a7c4a4d0184752..f5a64c1776a77a 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/nereids/rules/expression/check/CheckCastTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/nereids/rules/expression/check/CheckCastTest.java @@ -60,9 +60,39 @@ public void testCastBetweenVariantTypes() { VariantType v1Source = new VariantType(100); VariantType v1SameProperties = new VariantType(100); VariantType v1DifferentProperties = new VariantType(200); + VariantType v2Source = v1Source.withComputeV2(true); + VariantType v2DifferentProperties = v1DifferentProperties.withComputeV2(true); Assertions.assertTrue(CheckCast.check(v1Source, v1SameProperties, true)); Assertions.assertFalse(CheckCast.check(v1Source, v1DifferentProperties, true)); + Assertions.assertTrue(CheckCast.check(v2Source, v2DifferentProperties, true)); + Assertions.assertTrue(CheckCast.check( + ArrayType.of(v2Source), ArrayType.of(v2DifferentProperties), true)); + Assertions.assertFalse(CheckCast.check(v1Source, v2Source, true)); + Assertions.assertFalse(CheckCast.check(v2Source, v1Source, true)); + } + + @Test + public void testComputeVariantV2CastSourcesMatchExecutionKernel() { + VariantType target = VariantType.COMPUTE_V2_INSTANCE; + Assertions.assertTrue(CheckCast.check(IntegerType.INSTANCE, target, true)); + Assertions.assertTrue(CheckCast.check(JsonType.INSTANCE, target, true)); + Assertions.assertTrue(CheckCast.check(ArrayType.of(IntegerType.INSTANCE), target, true)); + Assertions.assertTrue(CheckCast.check(DecimalV3Type.SYSTEM_DEFAULT, target, true)); + + Assertions.assertFalse(CheckCast.check(VariantType.INSTANCE, target, true)); + Assertions.assertFalse(CheckCast.check(TimeV2Type.MAX, target, true)); + Assertions.assertFalse(CheckCast.check( + DecimalV3Type.createDecimalV3TypeNoCheck(39, 0), target, true)); + Assertions.assertFalse(CheckCast.check( + MapType.of(StringType.INSTANCE, IntegerType.INSTANCE), target, true)); + Assertions.assertFalse(CheckCast.check( + new StructType(Lists.newArrayList( + new StructField("field", IntegerType.INSTANCE, true, ""))), + target, true)); + Assertions.assertFalse(CheckCast.check( + ArrayType.of(MapType.of(StringType.INSTANCE, IntegerType.INSTANCE)), + target, true)); } @Test diff --git a/fe/fe-core/src/test/java/org/apache/doris/nereids/util/TypeCoercionUtilsTest.java b/fe/fe-core/src/test/java/org/apache/doris/nereids/util/TypeCoercionUtilsTest.java index b440eaa24210d5..bc1da8711559dc 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/nereids/util/TypeCoercionUtilsTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/nereids/util/TypeCoercionUtilsTest.java @@ -174,6 +174,11 @@ public void testVariantCommonTypePreservesLegacyBehavior() { VariantType v2 = v1.withComputeV2(true); VariantType anotherV2 = anotherV1.withComputeV2(true); + Assertions.assertTrue(v1.hasCommonExecutionTypeWith(anotherV1)); + Assertions.assertFalse(v1.isCastCompatibleWith(anotherV1)); + Assertions.assertFalse(TypeCoercionUtils.isNoOpCastCompatible(v1, anotherV1)); + Assertions.assertTrue(TypeCoercionUtils.isNoOpCastCompatible(v2, anotherV2)); + Assertions.assertEquals(anotherV1, TypeCoercionUtils.findWiderTypeForTwo(v1, anotherV1, false, true).get()); Assertions.assertEquals(VariantType.COMPUTE_V2_INSTANCE, diff --git a/regression-test/data/paimon_write/test_paimon_write_variant.out b/regression-test/data/paimon_write/test_paimon_write_variant.out new file mode 100644 index 00000000000000..cdf87fb1b01626 --- /dev/null +++ b/regression-test/data/paimon_write/test_paimon_write_variant.out @@ -0,0 +1,26 @@ +-- This file is automatically generated. You should know what you did if you want to edit this +-- !variant_heterogeneous -- +20 1 "row-one" +21 "row-two" 2 + +-- !variant_object -- +doris Hangzhou 1 true {} [] 中文😀 2 + +-- !variant_nulls -- +4 null false true +5 \N true true + +-- !variant_scalars -- +10 true false +11 -128 32767 +12 -2147483648 9223372036854775807 +13 1.25 -2.5 +14 123456.789 -0.000001 +15 "plain-string" "中文😀" +16 "2024-02-29" "2024-02-29 12:34:56.123456" + +-- !variant_long_string -- +45056 large-string + +-- !variant_row_count -- +15 diff --git a/regression-test/data/paimon_write/test_paimon_write_variant_dml.out b/regression-test/data/paimon_write/test_paimon_write_variant_dml.out new file mode 100644 index 00000000000000..c5b69cc6d55f39 --- /dev/null +++ b/regression-test/data/paimon_write/test_paimon_write_variant_dml.out @@ -0,0 +1,28 @@ +-- This file is automatically generated. You should know what you did if you want to edit this +-- !variant_dml_rows -- +1 {"n":1,"source":"direct"} default-note p1 false +10 {"mode":"reordered"} reordered p3 false +11 {"mode":"default"} default-note p3 false +12 \N default-note p3 true +2 ["direct",2] default-note p1 false +20 {"mode":"cte"} cte-note p4 false +21 {"mode":"union-a"} union p4 false +22 22 union p4 false +3 null default-note p2 false +30 {"generated":0} generated p5 false +37 {"generated":7} generated p5 false +4 \N default-note p2 true +50 {"partition":"static"} static-note static false +51 {"partition":"dynamic-a"} dynamic dynamic-a false +52 {"partition":"dynamic-b"} dynamic dynamic-b false + +-- !variant_partition_overwrite -- +10 {"state":"new-east"} east +2 {"state":"old-west"} west + +-- !variant_full_overwrite -- +20 {"state":"full-a"} all +21 {"state":"full-b"} all + +-- !variant_empty_overwrite -- +0 diff --git a/regression-test/data/paimon_write/test_paimon_write_variant_errors.out b/regression-test/data/paimon_write/test_paimon_write_variant_errors.out new file mode 100644 index 00000000000000..3a9b9e241d513a --- /dev/null +++ b/regression-test/data/paimon_write/test_paimon_write_variant_errors.out @@ -0,0 +1,5 @@ +-- This file is automatically generated. You should know what you did if you want to edit this +-- !variant_after_errors -- +20 false true {"recovered":true} +21 false \N "not-json" +22 false \N "{\\"typed\\":\\"string\\"}" diff --git a/regression-test/data/paimon_write/test_paimon_write_variant_nested.out b/regression-test/data/paimon_write/test_paimon_write_variant_nested.out new file mode 100644 index 00000000000000..e779da3b1426b2 --- /dev/null +++ b/regression-test/data/paimon_write/test_paimon_write_variant_nested.out @@ -0,0 +1,16 @@ +-- This file is automatically generated. You should know what you did if you want to edit this +-- !variant_nested_values -- +array-object null \N 7 2 null \N struct-object first second + +-- !variant_nested_containers -- +2 false 0 false 0 false +3 true \N true \N true + +-- !variant_deep_value -- +1 depth-1 deep-ok + +-- !variant_deep_null -- +true + +-- !variant_deep_count -- +2 diff --git a/regression-test/data/paimon_write/test_paimon_write_variant_shredding.out b/regression-test/data/paimon_write/test_paimon_write_variant_shredding.out new file mode 100644 index 00000000000000..9bfdb63cb0e664 --- /dev/null +++ b/regression-test/data/paimon_write/test_paimon_write_variant_shredding.out @@ -0,0 +1,22 @@ +-- This file is automatically generated. You should know what you did if you want to edit this +-- !variant_explicit_shredding -- +1 27 Beijing true alice 10 \N kept false +2 28 \N \N \N \N \N \N false +3 29 \N true \N \N \N \N false +4 \N \N \N \N \N \N \N false +5 \N \N \N \N \N \N \N false +6 \N \N \N \N \N \N \N false +7 \N \N \N \N \N \N \N true +8 \N \N \N bob 30 nested-kept root-kept false + +-- !variant_mixed_layout -- +100 100 old \N \N +101 \N \N residual \N +200 200 new \N kept +201 201 \N \N \N + +-- !variant_inferred_shredding -- +300 30 alice \N \N first +301 31 bob \N \N second +400 \N \N Hangzhou true third +401 \N \N Shanghai false fourth diff --git a/regression-test/data/paimon_write/test_paimon_write_variant_table_modes.out b/regression-test/data/paimon_write/test_paimon_write_variant_table_modes.out new file mode 100644 index 00000000000000..55e403ae05ecfe --- /dev/null +++ b/regression-test/data/paimon_write/test_paimon_write_variant_table_modes.out @@ -0,0 +1,15 @@ +-- This file is automatically generated. You should know what you did if you want to edit this +-- !variant_pk -- +1 v3 100 3 +2 stable 2 1 + +-- !variant_dynamic_bucket -- +32 496 0 + +-- !variant_schema_evolution -- +1 before \N \N +2 after-add added \N +3 after-normal-column continued default-note + +-- !non_variant -- +1 {"plain":"string"} diff --git a/regression-test/suites/paimon_write/test_paimon_write_schema_change.groovy b/regression-test/suites/paimon_write/test_paimon_write_schema_change.groovy index 97951822fcdb52..9882ec73560824 100644 --- a/regression-test/suites/paimon_write/test_paimon_write_schema_change.groovy +++ b/regression-test/suites/paimon_write/test_paimon_write_schema_change.groovy @@ -637,7 +637,7 @@ suite("test_paimon_write_schema_change", "p0,external,paimon") { """, "ORDER BY id") - // Paimon 1.3.1 has no partition-key evolution in SchemaChange. + // Paimon 1.4.2 has no partition-key evolution in SchemaChange. // Doris therefore rejects ADD, DROP and REPLACE before catalog mutation. test { sql """ diff --git a/regression-test/suites/paimon_write/test_paimon_write_transaction.groovy b/regression-test/suites/paimon_write/test_paimon_write_transaction.groovy index b4ce0a942f3ae5..bfdeffa376bbe6 100644 --- a/regression-test/suites/paimon_write/test_paimon_write_transaction.groovy +++ b/regression-test/suites/paimon_write/test_paimon_write_transaction.groovy @@ -51,7 +51,10 @@ suite("test_paimon_write_transaction", "p0,external,paimon") { CREATE TABLE paimon.${dbName}.t_overwrite_part ( id INT, name STRING, region STRING ) USING paimon - PARTITIONED BY (region); + PARTITIONED BY (region) + TBLPROPERTIES ( + 'dynamic-partition-overwrite' = 'true' + ); DROP TABLE IF EXISTS paimon.${dbName}.t_overwrite_part_case; CREATE TABLE paimon.${dbName}.t_overwrite_part_case ( diff --git a/regression-test/suites/paimon_write/test_paimon_write_variant.groovy b/regression-test/suites/paimon_write/test_paimon_write_variant.groovy new file mode 100644 index 00000000000000..097d8c7aa352f1 --- /dev/null +++ b/regression-test/suites/paimon_write/test_paimon_write_variant.groovy @@ -0,0 +1,156 @@ +// 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. + +suite("test_paimon_write_variant", "p0,external,paimon") { + String enabled = context.config.otherConfigs.get("enablePaimonTest") + if (enabled == null || !enabled.equalsIgnoreCase("true")) { + logger.info("disable paimon test.") + return + } + + String externalEnvIp = context.config.otherConfigs.get("externalEnvIp") + String minioPort = context.config.otherConfigs.get("iceberg_minio_port") + String catalogName = "test_pw_variant_catalog" + String dbName = "test_pw_variant_db" + + spark_paimon_multi """ + CREATE DATABASE IF NOT EXISTS paimon.${dbName}; + DROP TABLE IF EXISTS paimon.${dbName}.t_variant_basic; + CREATE TABLE paimon.${dbName}.t_variant_basic ( + id INT, + payload VARIANT, + secondary VARIANT + ) USING paimon + TBLPROPERTIES ('file.format' = 'parquet'); + """ + + sql """DROP CATALOG IF EXISTS ${catalogName}""" + sql """ + CREATE CATALOG ${catalogName} PROPERTIES ( + 'type' = 'paimon', + 'paimon.catalog.type' = 'filesystem', + 'warehouse' = 's3://warehouse/wh', + 's3.endpoint' = 'http://${externalEnvIp}:${minioPort}', + 's3.access_key' = 'admin', + 's3.secret_key' = 'password', + 's3.path.style.access' = 'true' + ) + """ + sql """SWITCH ${catalogName}""" + sql """USE ${dbName}""" + + try { + // Paimon Variant writes are deliberately V2-only. + sql """SET enable_variant_v2 = false""" + test { + sql """INSERT INTO t_variant_basic VALUES + (0, parse_to_variant('{"disabled":true}'), NULL)""" + exception "set enable_variant_v2=true" + } + sql """SET enable_variant_v2 = true""" + sql """SET force_jni_scanner = true""" + + // JSON containers, JSON null and SQL NULL are different logical values. + sql """ + INSERT INTO t_variant_basic VALUES + (1, parse_to_variant(CONCAT( + '{"object":{"name":"doris","address":{"city":"Hangzhou"}},"array":[1,true,null,"x"],"emptyObject":{},"emptyArray":[],"explicitNull":null,"escaped":"line', + CHAR(92), 'nquote', CHAR(92), CHAR(34), '","unicode":"中文😀"}')), + parse_to_variant('{"second":2}')), + (2, parse_to_variant('{}'), parse_to_variant('[]')), + (3, parse_to_variant('[]'), parse_to_variant('{}')), + (4, parse_to_variant('null'), parse_to_variant('null')), + (5, CAST(NULL AS VARIANT), CAST(NULL AS VARIANT)) + """ + + // Typed scalar values exercise every V2 primitive family used by Doris. + sql """ + INSERT INTO t_variant_basic VALUES + (10, CAST(TRUE AS VARIANT), CAST(FALSE AS VARIANT)), + (11, CAST(CAST(-128 AS TINYINT) AS VARIANT), + CAST(CAST(32767 AS SMALLINT) AS VARIANT)), + (12, CAST(CAST(-2147483648 AS INT) AS VARIANT), + CAST(CAST(9223372036854775807 AS BIGINT) AS VARIANT)), + (13, CAST(CAST(1.25 AS FLOAT) AS VARIANT), + CAST(CAST(-2.5 AS DOUBLE) AS VARIANT)), + (14, CAST(CAST(123456.789 AS DECIMAL(12, 3)) AS VARIANT), + CAST(CAST(-0.000001 AS DECIMAL(18, 6)) AS VARIANT)), + (15, CAST(CAST('plain-string' AS VARCHAR(32)) AS VARIANT), + CAST(CAST('中文😀' AS STRING) AS VARIANT)), + (16, CAST(DATE '2024-02-29' AS VARIANT), + CAST(CAST('2024-02-29 12:34:56.123456' AS DATETIMEV2(6)) AS VARIANT)), + (17, CAST(REPEAT('long-value-', 4096) AS VARIANT), + parse_to_variant('{"batch":"large-string"}')) + """ + + // Every VALUES row must reach Variant coercion before the inline table chooses a common + // type. In particular, the integer in the first row must not become the string "1". + sql """ + INSERT INTO t_variant_basic VALUES + (20, 1, 'row-one'), + (21, 'row-two', 2) + """ + order_qt_variant_heterogeneous """ + SELECT id, payload, secondary + FROM t_variant_basic + WHERE id IN (20, 21) + ORDER BY id + """ + + order_qt_variant_object """ + SELECT + CAST(payload['object']['name'] AS STRING), + CAST(payload['object']['address']['city'] AS STRING), + CAST(payload['array'][1] AS INT), + CAST(payload['array'][2] AS BOOLEAN), + CAST(payload['emptyObject'] AS STRING), + CAST(payload['emptyArray'] AS STRING), + CAST(payload['unicode'] AS STRING), + CAST(secondary['second'] AS INT) + FROM t_variant_basic + WHERE id = 1 + """ + + order_qt_variant_nulls """ + SELECT id, payload, payload IS NULL, payload['missing'] IS NULL + FROM t_variant_basic + WHERE id IN (4, 5) + ORDER BY id + """ + + order_qt_variant_scalars """ + SELECT id, payload, secondary + FROM t_variant_basic + WHERE id BETWEEN 10 AND 16 + ORDER BY id + """ + + qt_variant_long_string """ + SELECT LENGTH(CAST(payload AS STRING)), CAST(secondary['batch'] AS STRING) + FROM t_variant_basic + WHERE id = 17 + """ + + // Refresh metadata and verify that all Doris-written rows remain readable through the + // Paimon JNI Variant reader. + sql """REFRESH TABLE t_variant_basic""" + qt_variant_row_count """SELECT COUNT(*) FROM t_variant_basic""" + } finally { + sql """SET force_jni_scanner = false""" + sql """DROP CATALOG IF EXISTS ${catalogName}""" + } +} diff --git a/regression-test/suites/paimon_write/test_paimon_write_variant_dml.groovy b/regression-test/suites/paimon_write/test_paimon_write_variant_dml.groovy new file mode 100644 index 00000000000000..d3808c93dec3ce --- /dev/null +++ b/regression-test/suites/paimon_write/test_paimon_write_variant_dml.groovy @@ -0,0 +1,185 @@ +// 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. + +suite("test_paimon_write_variant_dml", "p0,external,paimon") { + String enabled = context.config.otherConfigs.get("enablePaimonTest") + if (enabled == null || !enabled.equalsIgnoreCase("true")) { + logger.info("disable paimon test.") + return + } + + String externalEnvIp = context.config.otherConfigs.get("externalEnvIp") + String minioPort = context.config.otherConfigs.get("iceberg_minio_port") + String catalogName = "test_pw_variant_dml_catalog" + String dbName = "test_pw_variant_dml_db" + + spark_paimon_multi """ + CREATE DATABASE IF NOT EXISTS paimon.${dbName}; + + DROP TABLE IF EXISTS paimon.${dbName}.t_variant_dml; + CREATE TABLE paimon.${dbName}.t_variant_dml ( + id INT, + payload VARIANT, + note STRING NOT NULL DEFAULT 'default-note', + pt STRING + ) USING paimon + PARTITIONED BY (pt) + TBLPROPERTIES ('file.format' = 'parquet'); + + DROP TABLE IF EXISTS paimon.${dbName}.t_variant_overwrite; + CREATE TABLE paimon.${dbName}.t_variant_overwrite ( + id INT, + payload VARIANT, + pt STRING + ) USING paimon + PARTITIONED BY (pt) + TBLPROPERTIES ('file.format' = 'parquet'); + """ + + sql """DROP CATALOG IF EXISTS ${catalogName}""" + sql """ + CREATE CATALOG ${catalogName} PROPERTIES ( + 'type' = 'paimon', + 'paimon.catalog.type' = 'filesystem', + 'warehouse' = 's3://warehouse/wh', + 's3.endpoint' = 'http://${externalEnvIp}:${minioPort}', + 's3.access_key' = 'admin', + 's3.secret_key' = 'password', + 's3.path.style.access' = 'true' + ) + """ + sql """SWITCH ${catalogName}""" + sql """USE ${dbName}""" + sql """SET enable_variant_v2 = true""" + sql """SET force_jni_scanner = true""" + + try { + // INSERT SELECT preserves the V2 value and metadata buffers without routing + // through an internal OLAP Variant column, whose storage format is legacy V1. + sql """ + INSERT INTO t_variant_dml (id, payload, pt) + SELECT 1, parse_to_variant('{"source":"direct","n":1}'), 'p1' + UNION ALL + SELECT 2, parse_to_variant('["direct",2]'), 'p1' + UNION ALL + SELECT 3, parse_to_variant('null'), 'p2' + UNION ALL + SELECT 4, CAST(NULL AS VARIANT), 'p2' + """ + + // Reordered columns, partial columns and writer-side defaults. + sql """ + INSERT INTO t_variant_dml (pt, note, payload, id) VALUES + ('p3', 'reordered', parse_to_variant('{"mode":"reordered"}'), 10) + """ + sql """ + INSERT INTO t_variant_dml (pt, payload, id) VALUES + ('p3', parse_to_variant('{"mode":"default"}'), 11) + """ + sql """ + INSERT INTO t_variant_dml (pt, id) VALUES ('p3', 12) + """ + + // CTE, UNION ALL, expression-generated Variant and an empty input. + sql """ + INSERT INTO t_variant_dml + WITH source AS ( + SELECT 20 AS id, parse_to_variant('{"mode":"cte"}') AS payload, + 'cte-note' AS note, 'p4' AS pt + ) + SELECT id, payload, note, pt FROM source + """ + sql """ + INSERT INTO t_variant_dml + SELECT 21, parse_to_variant('{"mode":"union-a"}'), 'union', 'p4' + UNION ALL + SELECT 22, CAST(CAST(22 AS BIGINT) AS VARIANT), 'union', 'p4' + """ + sql """ + INSERT INTO t_variant_dml + SELECT 30 + number, + parse_to_variant(CONCAT('{"generated":', number, '}')), + 'generated', + 'p5' + FROM numbers("number" = "8") + """ + sql """ + INSERT INTO t_variant_dml + SELECT 100, parse_to_variant('{"unused":true}'), 'empty', 'p0' + WHERE 1 = 0 + """ + + // Static and dynamic partition writes. + sql """ + INSERT INTO t_variant_dml PARTITION (pt = 'static') + VALUES (50, parse_to_variant('{"partition":"static"}'), 'static-note') + """ + sql """ + INSERT INTO t_variant_dml VALUES + (51, parse_to_variant('{"partition":"dynamic-a"}'), 'dynamic', 'dynamic-a'), + (52, parse_to_variant('{"partition":"dynamic-b"}'), 'dynamic', 'dynamic-b') + """ + + order_qt_variant_dml_rows """ + SELECT id, payload, note, pt, payload IS NULL + FROM t_variant_dml + WHERE id IN (1, 2, 3, 4, 10, 11, 12, 20, 21, 22, 30, 37, 50, 51, 52) + ORDER BY id + """ + + // Static-partition overwrite exercises the overwrite writer with Variant V2 rows. + sql """ + INSERT INTO t_variant_overwrite VALUES + (1, parse_to_variant('{"state":"old-east"}'), 'east'), + (2, parse_to_variant('{"state":"old-west"}'), 'west') + """ + sql """ + INSERT OVERWRITE TABLE t_variant_overwrite + PARTITION (pt = 'east') + VALUES (10, parse_to_variant('{"state":"new-east"}')) + """ + order_qt_variant_partition_overwrite """ + SELECT id, payload, pt + FROM t_variant_overwrite + ORDER BY id + """ + + // A full-table overwrite on a partitioned table must remove untouched old partitions. + sql """ + INSERT OVERWRITE TABLE t_variant_overwrite VALUES + (20, parse_to_variant('{"state":"full-a"}'), 'all'), + (21, parse_to_variant('{"state":"full-b"}'), 'all') + """ + order_qt_variant_full_overwrite """ + SELECT id, payload, pt + FROM t_variant_overwrite + ORDER BY id + """ + + // Empty full-table overwrite must still commit an empty snapshot. + sql """ + INSERT OVERWRITE TABLE t_variant_overwrite + SELECT 30, parse_to_variant('{"state":"unused"}'), 'empty' + WHERE 1 = 0 + """ + qt_variant_empty_overwrite """SELECT COUNT(*) FROM t_variant_overwrite""" + + } finally { + sql """SET force_jni_scanner = false""" + sql """DROP CATALOG IF EXISTS ${catalogName}""" + } +} diff --git a/regression-test/suites/paimon_write/test_paimon_write_variant_errors.groovy b/regression-test/suites/paimon_write/test_paimon_write_variant_errors.groovy new file mode 100644 index 00000000000000..fd762cae16ac01 --- /dev/null +++ b/regression-test/suites/paimon_write/test_paimon_write_variant_errors.groovy @@ -0,0 +1,133 @@ +// 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. + +suite("test_paimon_write_variant_errors", "p0,external,paimon") { + String enabled = context.config.otherConfigs.get("enablePaimonTest") + if (enabled == null || !enabled.equalsIgnoreCase("true")) { + logger.info("disable paimon test.") + return + } + + String externalEnvIp = context.config.otherConfigs.get("externalEnvIp") + String minioPort = context.config.otherConfigs.get("iceberg_minio_port") + String catalogName = "test_pw_variant_errors_catalog" + String dbName = "test_pw_variant_errors_db" + + spark_paimon_multi """ + CREATE DATABASE IF NOT EXISTS paimon.${dbName}; + + DROP TABLE IF EXISTS paimon.${dbName}.t_variant_error; + CREATE TABLE paimon.${dbName}.t_variant_error ( + id INT, + payload VARIANT + ) USING paimon + TBLPROPERTIES ('file.format' = 'parquet'); + + DROP TABLE IF EXISTS paimon.${dbName}.t_variant_nested_error; + CREATE TABLE paimon.${dbName}.t_variant_nested_error ( + id INT, + payloads ARRAY + ) USING paimon + TBLPROPERTIES ('file.format' = 'parquet'); + """ + + sql """DROP CATALOG IF EXISTS ${catalogName}""" + sql """ + CREATE CATALOG ${catalogName} PROPERTIES ( + 'type' = 'paimon', + 'paimon.catalog.type' = 'filesystem', + 'warehouse' = 's3://warehouse/wh', + 's3.endpoint' = 'http://${externalEnvIp}:${minioPort}', + 's3.access_key' = 'admin', + 's3.secret_key' = 'password', + 's3.path.style.access' = 'true' + ) + """ + sql """SWITCH ${catalogName}""" + sql """USE ${dbName}""" + sql """CREATE DATABASE IF NOT EXISTS internal.${dbName}""" + + try { + // Both top-level and nested targets fail during analysis when V2 is disabled. + sql """SET enable_variant_v2 = false""" + test { + sql """INSERT INTO t_variant_error VALUES + (1, parse_to_variant('{"disabled":"top"}'))""" + exception "set enable_variant_v2=true" + } + test { + sql """INSERT INTO t_variant_nested_error VALUES + (1, CAST(NULL AS ARRAY))""" + exception "set enable_variant_v2=true" + } + + // Capture V1 output types without executing a V1 load, then verify that a + // later Paimon write rejects them in analysis before reaching the BE. + sql """DROP VIEW IF EXISTS internal.${dbName}.v_variant_v1""" + sql """DROP VIEW IF EXISTS internal.${dbName}.v_variant_v1_nested""" + sql """ + CREATE VIEW internal.${dbName}.v_variant_v1 AS + SELECT 10 AS id, parse_to_variant('{"legacy":true}') AS payload + """ + sql """ + CREATE VIEW internal.${dbName}.v_variant_v1_nested AS + SELECT 11 AS id, + CAST(NULL AS ARRAY) AS payloads + """ + + sql """SET enable_variant_v2 = true""" + sql """SET force_jni_scanner = true""" + test { + sql """ + INSERT INTO t_variant_error + SELECT id, payload + FROM internal.${dbName}.v_variant_v1 + """ + exception "Variant V1" + } + test { + sql """ + INSERT INTO t_variant_nested_error + SELECT id, payloads + FROM internal.${dbName}.v_variant_v1_nested + """ + exception "Variant V1" + } + + // Valid V2 writes still work after analysis failures in the same session. + sql """ + INSERT INTO t_variant_error VALUES + (20, parse_to_variant('{"recovered":true}')), + (21, try_parse_to_variant('not-json')), + (22, CAST(CAST('{"typed":"string"}' AS STRING) AS VARIANT)) + """ + // Invalid JSON is preserved as a Variant string unless the global + // throw-on-invalid-JSON option is enabled. + order_qt_variant_after_errors """ + SELECT id, payload IS NULL, + CAST(payload['recovered'] AS BOOLEAN), + payload + FROM t_variant_error + ORDER BY id + """ + } finally { + sql """SET force_jni_scanner = false""" + sql """DROP VIEW IF EXISTS internal.${dbName}.v_variant_v1""" + sql """DROP VIEW IF EXISTS internal.${dbName}.v_variant_v1_nested""" + sql """DROP CATALOG IF EXISTS ${catalogName}""" + } +} diff --git a/regression-test/suites/paimon_write/test_paimon_write_variant_nested.groovy b/regression-test/suites/paimon_write/test_paimon_write_variant_nested.groovy new file mode 100644 index 00000000000000..86a72be8839a47 --- /dev/null +++ b/regression-test/suites/paimon_write/test_paimon_write_variant_nested.groovy @@ -0,0 +1,195 @@ +// 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. + +suite("test_paimon_write_variant_nested", "p0,external,paimon") { + String enabled = context.config.otherConfigs.get("enablePaimonTest") + if (enabled == null || !enabled.equalsIgnoreCase("true")) { + logger.info("disable paimon test.") + return + } + + String externalEnvIp = context.config.otherConfigs.get("externalEnvIp") + String minioPort = context.config.otherConfigs.get("iceberg_minio_port") + String catalogName = "test_pw_variant_nested_catalog" + String dbName = "test_pw_variant_nested_db" + + spark_paimon_multi """ + CREATE DATABASE IF NOT EXISTS paimon.${dbName}; + + DROP TABLE IF EXISTS paimon.${dbName}.t_variant_nested; + CREATE TABLE paimon.${dbName}.t_variant_nested ( + id INT, + variants ARRAY, + variant_map MAP, + variant_struct STRUCT, + first_payload VARIANT, + second_payload VARIANT + ) USING paimon + TBLPROPERTIES ('file.format' = 'parquet'); + + DROP TABLE IF EXISTS paimon.${dbName}.t_variant_deep; + CREATE TABLE paimon.${dbName}.t_variant_deep ( + id INT, + deep STRUCT< + level1:ARRAY< + MAP> + > + > + ) USING paimon + TBLPROPERTIES ('file.format' = 'parquet'); + """ + + sql """DROP CATALOG IF EXISTS ${catalogName}""" + sql """ + CREATE CATALOG ${catalogName} PROPERTIES ( + 'type' = 'paimon', + 'paimon.catalog.type' = 'filesystem', + 'warehouse' = 's3://warehouse/wh', + 's3.endpoint' = 'http://${externalEnvIp}:${minioPort}', + 's3.access_key' = 'admin', + 's3.secret_key' = 'password', + 's3.path.style.access' = 'true' + ) + """ + sql """SWITCH ${catalogName}""" + sql """USE ${dbName}""" + sql """SET enable_variant_v2 = true""" + sql """SET force_jni_scanner = true""" + + try { + // ARRAY, MAP, STRUCT and multiple Variant columns in one Arrow batch. + sql """ + INSERT INTO t_variant_nested VALUES + ( + 1, + array( + parse_to_variant('{"kind":"array-object","n":1}'), + parse_to_variant('null'), + CAST(NULL AS VARIANT), + CAST(CAST(7 AS INT) AS VARIANT) + ), + map( + 'object', parse_to_variant('{"kind":"map-object","n":2}'), + 'json_null', parse_to_variant('null'), + 'sql_null', CAST(NULL AS VARIANT) + ), + named_struct( + 'label', 'struct-value', + 'payload', parse_to_variant('{"kind":"struct-object","n":3}') + ), + parse_to_variant('{"column":"first"}'), + parse_to_variant('["second",2]') + ), + ( + 2, + array(), + map(), + named_struct('label', 'empty', 'payload', parse_to_variant('{}')), + parse_to_variant('[]'), + CAST(NULL AS VARIANT) + ), + (3, NULL, NULL, NULL, NULL, NULL) + """ + + order_qt_variant_nested_values """ + SELECT + CAST(variants[1]['kind'] AS STRING), + variants[2], + variants[3], + CAST(variants[4] AS INT), + CAST(variant_map['object']['n'] AS INT), + variant_map['json_null'], + variant_map['sql_null'], + CAST(variant_struct.payload['kind'] AS STRING), + CAST(first_payload['column'] AS STRING), + CAST(second_payload[1] AS STRING) + FROM t_variant_nested + WHERE id = 1 + """ + + order_qt_variant_nested_containers """ + SELECT id, + variants IS NULL, SIZE(variants), + variant_map IS NULL, SIZE(variant_map), + variant_struct IS NULL + FROM t_variant_nested + WHERE id IN (2, 3) + ORDER BY id + """ + + // Deep nesting is P0: STRUCT -> ARRAY -> MAP -> STRUCT -> VARIANT. + sql """ + INSERT INTO t_variant_deep VALUES + ( + 1, + named_struct( + 'level1', + array( + map( + 'outer', + named_struct( + 'note', 'depth-1', + 'payload', parse_to_variant( + '{"level2":{"level3":{"level4":{"value":"deep-ok"}}}}') + ) + ) + ) + ) + ), + ( + 2, + named_struct( + 'level1', + array( + map( + 'null-leaf', + named_struct( + 'note', 'depth-null', + 'payload', CAST(NULL AS VARIANT) + ) + ) + ) + ) + ) + """ + + order_qt_variant_deep_value """ + SELECT id, + deep.level1[1]['outer'].note, + CAST(deep.level1[1]['outer'].payload['level2']['level3']['level4']['value'] + AS STRING) + FROM t_variant_deep + WHERE id = 1 + """ + + order_qt_variant_deep_null """ + SELECT deep.level1[1]['null-leaf'].payload IS NULL + FROM t_variant_deep + WHERE id = 2 + """ + + // Refreshing metadata must not affect nested Variant reads. + sql """REFRESH TABLE t_variant_deep""" + qt_variant_deep_count """SELECT COUNT(*) FROM t_variant_deep""" + } finally { + sql """SET force_jni_scanner = false""" + sql """DROP CATALOG IF EXISTS ${catalogName}""" + } +} diff --git a/regression-test/suites/paimon_write/test_paimon_write_variant_shredding.groovy b/regression-test/suites/paimon_write/test_paimon_write_variant_shredding.groovy new file mode 100644 index 00000000000000..ad3e6d53aa90eb --- /dev/null +++ b/regression-test/suites/paimon_write/test_paimon_write_variant_shredding.groovy @@ -0,0 +1,303 @@ +// 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. + +suite("test_paimon_write_variant_shredding", "p0,external,paimon") { + String enabled = context.config.otherConfigs.get("enablePaimonTest") + if (enabled == null || !enabled.equalsIgnoreCase("true")) { + logger.info("disable paimon test.") + return + } + + String externalEnvIp = context.config.otherConfigs.get("externalEnvIp") + String minioPort = context.config.otherConfigs.get("iceberg_minio_port") + String catalogName = "test_pw_variant_shredding_catalog" + String dbName = "test_pw_variant_shredding_db" + String shreddingSchema = + '{"type":"ROW","fields":[{"name":"payload","type":{"type":"ROW","fields":[' + + '{"name":"age","type":"INT"},' + + '{"name":"city","type":"STRING"},' + + '{"name":"active","type":"BOOLEAN"},' + + '{"name":"profile","type":{"type":"ROW","fields":[' + + '{"name":"name","type":"STRING"},' + + '{"name":"scores","type":{"type":"ARRAY","element":"INT"}}' + + ']}}]}}]}' + + // TODO: Use variant.shreddingSchema after Paimon passes the global option to its Parquet + // builder. In 1.4.2 the builder still requires the fallback spelling used below. + spark_paimon_multi """ + CREATE DATABASE IF NOT EXISTS paimon.${dbName}; + + DROP TABLE IF EXISTS paimon.${dbName}.t_variant_shredded; + CREATE TABLE paimon.${dbName}.t_variant_shredded ( + id INT, + payload VARIANT + ) USING paimon + TBLPROPERTIES ( + 'file.format' = 'parquet', + 'write-only' = 'true', + 'parquet.variant.shreddingSchema' = '${shreddingSchema}', + 'variant.inferShreddingSchema' = 'true' + ); + + DROP TABLE IF EXISTS paimon.${dbName}.t_variant_inferred; + CREATE TABLE paimon.${dbName}.t_variant_inferred ( + id INT, + payload VARIANT + ) USING paimon + TBLPROPERTIES ( + 'file.format' = 'parquet', + 'write-only' = 'true', + 'variant.inferShreddingSchema' = 'true' + ); + + DROP TABLE IF EXISTS paimon.${dbName}.t_variant_mixed; + CREATE TABLE paimon.${dbName}.t_variant_mixed ( + id INT, + payload VARIANT + ) USING paimon + TBLPROPERTIES ( + 'file.format' = 'parquet', + 'write-only' = 'true' + ); + """ + + def createDorisCatalog = { + sql """DROP CATALOG IF EXISTS ${catalogName}""" + sql """ + CREATE CATALOG ${catalogName} PROPERTIES ( + 'type' = 'paimon', + 'paimon.catalog.type' = 'filesystem', + 'warehouse' = 's3://warehouse/wh', + 's3.endpoint' = 'http://${externalEnvIp}:${minioPort}', + 's3.access_key' = 'admin', + 's3.secret_key' = 'password', + 's3.path.style.access' = 'true' + ) + """ + sql """SWITCH ${catalogName}""" + sql """USE ${dbName}""" + sql """SET enable_variant_v2 = true""" + sql """SET force_jni_scanner = true""" + } + + def sparkValues = { rows -> + rows.collect { row -> + row.collect { value -> value == null ? null : value.toString() } + } + } + + String filesTableSuffix = '$files' + def dataFiles = { String tableName -> + String filesQuery = """ + SELECT file_path + FROM paimon.${dbName}.`${tableName}${filesTableSuffix}` + ORDER BY file_path + """ + spark_paimon(filesQuery).collect { row -> row[0].toString() } + } + + // Read the data file as ordinary Parquet, bypassing Paimon's logical Variant reader. This + // proves that shredding produced typed_value columns rather than merely round-tripping the + // original value/metadata pair. + def rawParquetSource = { String path -> + return """S3( + "uri" = "${path}", + "s3.endpoint" = "http://${externalEnvIp}:${minioPort}", + "s3.access_key" = "admin", + "s3.secret_key" = "password", + "s3.region" = "us-east-1", + "use_path_style" = "true", + "format" = "parquet" + )""" + } + def rawPayloadType = { String path -> + def columns = sql """DESC FUNCTION ${rawParquetSource(path)}""" + def payloadColumn = columns.find { it[0].toString().equalsIgnoreCase("payload") } + assertTrue(payloadColumn != null, "No payload column in Paimon data file ${path}") + return payloadColumn[1].toString().toLowerCase() + } + + createDorisCatalog() + try { + // Cover typed fields, residual object fields, type mismatch fallback, nested ROW/ARRAY, + // root scalars, empty objects, Variant null, and SQL null. + sql """ + INSERT INTO t_variant_shredded VALUES + (1, parse_to_variant('{"age":27,"city":"Beijing","active":true,"profile":{"name":"alice","scores":[10,20]},"other":"kept"}')), + (2, parse_to_variant('{"age":28}')), + (3, parse_to_variant('{"age":"29","active":"true"}')), + (4, parse_to_variant('"scalar"')), + (5, parse_to_variant('{}')), + (6, parse_to_variant('null')), + (7, CAST(NULL AS VARIANT)), + (8, parse_to_variant('{"profile":{"name":"bob","scores":[30],"extra":"nested-kept"},"other":"root-kept"}')) + """ + + // Doris's Paimon reader must unshred typed and residual components back into one logical + // Variant value. + order_qt_variant_explicit_shredding """ + SELECT id, + CAST(payload['age'] AS STRING), + CAST(payload['city'] AS STRING), + CAST(payload['active'] AS BOOLEAN), + CAST(payload['profile']['name'] AS STRING), + CAST(payload['profile']['scores'][1] AS INT), + CAST(payload['profile']['extra'] AS STRING), + CAST(payload['other'] AS STRING), + payload IS NULL + FROM t_variant_shredded + ORDER BY id + """ + + def shreddedFiles = dataFiles("t_variant_shredded") + assertTrue(!shreddedFiles.isEmpty()) + def physicalRows = [] + shreddedFiles.each { filePath -> + String payloadType = rawPayloadType(filePath) + assertTrue(payloadType.contains("metadata:text")) + assertTrue(payloadType.contains("value:text")) + assertTrue(payloadType.contains("typed_value:struct")) + assertTrue(payloadType.contains("age:struct")) + assertTrue(payloadType.contains("profile:struct")) + // The explicit schema wins over inference; residual-only fields must not be promoted. + assertFalse(payloadType.contains("other:struct")) + + physicalRows.addAll(sql(""" + SELECT id, + payload.typed_value.age.typed_value, + CAST(payload.typed_value.age.value IS NOT NULL AS INT), + payload.typed_value.city.typed_value, + CAST(payload.typed_value.active.typed_value AS INT), + payload.typed_value.profile.typed_value.name.typed_value, + payload.typed_value.profile.typed_value.scores.typed_value[1].typed_value, + CAST(payload.value IS NOT NULL AS INT), + CAST(payload.metadata IS NOT NULL AS INT) + FROM ${rawParquetSource(filePath)} + WHERE id IN (1, 3) + """)) + } + physicalRows.sort { left, right -> + Integer.parseInt(left[0].toString()) <=> Integer.parseInt(right[0].toString()) + } + assertEquals([ + ["1", "27", "0", "Beijing", "1", "alice", "10", "1", "1"], + ["3", null, "1", null, null, null, null, "0", "1"] + ], sparkValues(physicalRows)) + + // First create ordinary value/metadata files, then enable shredding for the same table. + // Paimon's reader detects the physical schema per file and must read both layouts together. + sql """ + INSERT INTO t_variant_mixed VALUES + (100, parse_to_variant('{"age":100,"city":"old"}')), + (101, parse_to_variant('{"legacy":"residual"}')) + """ + def unshreddedFiles = dataFiles("t_variant_mixed") + assertTrue(!unshreddedFiles.isEmpty()) + unshreddedFiles.each { filePath -> + String payloadType = rawPayloadType(filePath) + assertTrue(payloadType.contains("value:text")) + assertTrue(payloadType.contains("metadata:text")) + assertFalse(payloadType.contains("typed_value")) + } + + spark_paimon """ + ALTER TABLE paimon.${dbName}.t_variant_mixed + SET TBLPROPERTIES ('parquet.variant.shreddingSchema' = '${shreddingSchema}') + """ + // Reload the serialized Paimon table used by the JNI writer so the next write observes + // the new file-format option. + createDorisCatalog() + sql """ + INSERT INTO t_variant_mixed VALUES + (200, parse_to_variant('{"age":200,"city":"new","extra":"kept"}')), + (201, parse_to_variant('{"age":"201"}')) + """ + + def mixedFiles = dataFiles("t_variant_mixed") + def newlyShreddedFiles = mixedFiles.findAll { !unshreddedFiles.contains(it) } + assertTrue(!newlyShreddedFiles.isEmpty()) + newlyShreddedFiles.each { filePath -> + assertTrue(rawPayloadType(filePath).contains("typed_value:struct")) + } + + order_qt_variant_mixed_layout """ + SELECT id, + CAST(payload['age'] AS STRING), + CAST(payload['city'] AS STRING), + CAST(payload['legacy'] AS STRING), + CAST(payload['extra'] AS STRING) + FROM t_variant_mixed + ORDER BY id + """ + + // Paimon 1.4 can infer one shredding schema per file writer. Doris still sends the same + // logical value/metadata pair; the SDK buffers the rows, chooses typed fields, and writes + // typed_value without a caller-provided schema. + sql """ + INSERT INTO t_variant_inferred VALUES + (300, parse_to_variant('{"age":30,"profile":{"name":"alice"},"extra":"first"}')), + (301, parse_to_variant('{"age":31,"profile":{"name":"bob"},"extra":"second"}')) + """ + def firstInferredFiles = dataFiles("t_variant_inferred") + assertTrue(!firstInferredFiles.isEmpty()) + firstInferredFiles.each { filePath -> + String payloadType = rawPayloadType(filePath) + assertTrue(payloadType.contains("metadata:text")) + assertTrue(payloadType.contains("value:text")) + assertTrue(payloadType.contains("typed_value:struct")) + assertTrue(payloadType.contains("age:struct")) + assertTrue(payloadType.contains("profile:struct")) + assertFalse(payloadType.contains("active:struct")) + assertFalse(payloadType.contains("city:struct")) + } + + // A later Doris statement opens a new Paimon file writer and may infer a different schema. + // Keep write-only enabled so compaction cannot hide the per-file schema difference. + sql """ + INSERT INTO t_variant_inferred VALUES + (400, parse_to_variant('{"city":"Hangzhou","active":true,"extra":"third"}')), + (401, parse_to_variant('{"city":"Shanghai","active":false,"extra":"fourth"}')) + """ + def allInferredFiles = dataFiles("t_variant_inferred") + def secondInferredFiles = allInferredFiles.findAll { !firstInferredFiles.contains(it) } + assertTrue(!secondInferredFiles.isEmpty()) + secondInferredFiles.each { filePath -> + String payloadType = rawPayloadType(filePath) + assertTrue(payloadType.contains("typed_value:struct")) + assertTrue(payloadType.contains("active:struct")) + assertTrue(payloadType.contains("city:struct")) + assertFalse(payloadType.contains("age:struct")) + assertFalse(payloadType.contains("profile:struct")) + } + + // Unshredding is a reader responsibility. It must use each file's physical schema and + // combine typed fields with residual values into one logical Variant column. + order_qt_variant_inferred_shredding """ + SELECT id, + CAST(payload['age'] AS INT), + CAST(payload['profile']['name'] AS STRING), + CAST(payload['city'] AS STRING), + CAST(payload['active'] AS BOOLEAN), + CAST(payload['extra'] AS STRING) + FROM t_variant_inferred + ORDER BY id + """ + } finally { + sql """SET force_jni_scanner = false""" + sql """DROP CATALOG IF EXISTS ${catalogName}""" + } +} diff --git a/regression-test/suites/paimon_write/test_paimon_write_variant_table_modes.groovy b/regression-test/suites/paimon_write/test_paimon_write_variant_table_modes.groovy new file mode 100644 index 00000000000000..df53da8c337b75 --- /dev/null +++ b/regression-test/suites/paimon_write/test_paimon_write_variant_table_modes.groovy @@ -0,0 +1,165 @@ +// 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. + +suite("test_paimon_write_variant_table_modes", "p0,external,paimon") { + String enabled = context.config.otherConfigs.get("enablePaimonTest") + if (enabled == null || !enabled.equalsIgnoreCase("true")) { + logger.info("disable paimon test.") + return + } + + String externalEnvIp = context.config.otherConfigs.get("externalEnvIp") + String minioPort = context.config.otherConfigs.get("iceberg_minio_port") + String catalogName = "test_pw_variant_modes_catalog" + String dbName = "test_pw_variant_modes_db" + + spark_paimon_multi """ + CREATE DATABASE IF NOT EXISTS paimon.${dbName}; + + DROP TABLE IF EXISTS paimon.${dbName}.t_variant_pk; + CREATE TABLE paimon.${dbName}.t_variant_pk ( + id INT, + payload VARIANT, + version BIGINT + ) USING paimon + TBLPROPERTIES ( + 'primary-key' = 'id', + 'bucket' = '2', + 'bucket-key' = 'id', + 'file.format' = 'parquet' + ); + + DROP TABLE IF EXISTS paimon.${dbName}.t_variant_dynamic_bucket; + CREATE TABLE paimon.${dbName}.t_variant_dynamic_bucket ( + id INT, + payload VARIANT + ) USING paimon + TBLPROPERTIES ( + 'primary-key' = 'id', + 'bucket' = '-1', + 'file.format' = 'parquet' + ); + + DROP TABLE IF EXISTS paimon.${dbName}.t_variant_schema; + CREATE TABLE paimon.${dbName}.t_variant_schema ( + id INT, + name STRING + ) USING paimon + TBLPROPERTIES ('file.format' = 'parquet'); + + DROP TABLE IF EXISTS paimon.${dbName}.t_non_variant; + CREATE TABLE paimon.${dbName}.t_non_variant ( + id INT, + payload STRING + ) USING paimon; + + DROP TABLE IF EXISTS paimon.${dbName}.t_variant_required; + CREATE TABLE paimon.${dbName}.t_variant_required ( + id INT, + payload VARIANT NOT NULL + ) USING paimon + TBLPROPERTIES ('file.format' = 'parquet'); + """ + + sql """DROP CATALOG IF EXISTS ${catalogName}""" + sql """ + CREATE CATALOG ${catalogName} PROPERTIES ( + 'type' = 'paimon', + 'paimon.catalog.type' = 'filesystem', + 'warehouse' = 's3://warehouse/wh', + 's3.endpoint' = 'http://${externalEnvIp}:${minioPort}', + 's3.access_key' = 'admin', + 's3.secret_key' = 'password', + 's3.path.style.access' = 'true' + ) + """ + sql """SWITCH ${catalogName}""" + sql """USE ${dbName}""" + sql """SET enable_variant_v2 = true""" + sql """SET force_jni_scanner = true""" + + try { + // Fixed-bucket primary-key table: later rows replace the same key. + sql """ + INSERT INTO t_variant_pk VALUES + (1, parse_to_variant('{"state":"v1","n":1}'), 1), + (2, parse_to_variant('{"state":"stable","n":2}'), 1), + (1, parse_to_variant('{"state":"v2","n":10}'), 2) + """ + sql """ + INSERT INTO t_variant_pk VALUES + (1, parse_to_variant('{"state":"v3","n":100}'), 3) + """ + order_qt_variant_pk """ + SELECT id, + CAST(payload['state'] AS STRING), + CAST(payload['n'] AS INT), + version + FROM t_variant_pk + ORDER BY id + """ + + // Dynamic bucket routing with Variant values. + sql """ + INSERT INTO t_variant_dynamic_bucket + SELECT number, + parse_to_variant(CONCAT('{"bucket":"dynamic","id":', number, '}')) + FROM numbers("number" = "32") + """ + order_qt_variant_dynamic_bucket """ + SELECT COUNT(*), + SUM(CAST(payload['id'] AS INT)), + SUM(CASE WHEN id <> CAST(payload['id'] AS INT) + THEN 1 ELSE 0 END) + FROM t_variant_dynamic_bucket + """ + + // Schema evolution: add Variant, write it, then add a normal column and continue writing. + sql """INSERT INTO t_variant_schema VALUES (1, 'before')""" + sql """ALTER TABLE t_variant_schema ADD COLUMN payload VARIANT NULL AFTER name""" + sql """ + INSERT INTO t_variant_schema (payload, name, id) VALUES + (parse_to_variant('{"schema":"added"}'), 'after-add', 2) + """ + sql """ALTER TABLE t_variant_schema ADD COLUMN note STRING NULL DEFAULT 'default-note'""" + sql """ + INSERT INTO t_variant_schema (id, name, payload) VALUES + (3, 'after-normal-column', parse_to_variant('{"schema":"continued"}')) + """ + sql """REFRESH TABLE t_variant_schema""" + order_qt_variant_schema_evolution """ + SELECT id, name, + CAST(payload['schema'] AS STRING), + note + FROM t_variant_schema + ORDER BY id + """ + + // Non-Variant Paimon writes remain unchanged while the session enables V2. + sql """INSERT INTO t_non_variant VALUES (1, '{"plain":"string"}')""" + order_qt_non_variant """SELECT id, payload FROM t_non_variant ORDER BY id""" + + // Paimon's real NOT NULL schema is enforced by the SDK. + test { + sql """INSERT INTO t_variant_required VALUES (1, CAST(NULL AS VARIANT))""" + exception "Cannot write null to non-null column(payload)" + } + } finally { + sql """SET force_jni_scanner = false""" + sql """DROP CATALOG IF EXISTS ${catalogName}""" + } +} diff --git a/regression-test/suites/variant_p0/test_variant_equality_contexts.groovy b/regression-test/suites/variant_p0/test_variant_equality_contexts.groovy index 050aaa566670df..7c74c725080945 100644 --- a/regression-test/suites/variant_p0/test_variant_equality_contexts.groovy +++ b/regression-test/suites/variant_p0/test_variant_equality_contexts.groovy @@ -165,7 +165,7 @@ suite("test_variant_equality_contexts", "p0,nonConcurrent") { test { sql "SELECT CAST(CAST(1.23 AS DECIMAL(76, 2)) AS VARIANT)" - exception "to Variant V2 is not supported" + exception "cannot cast DECIMALV3(76, 2) to variant" } order_qt_intersect_encoded_canonical_numeric """