| /* |
| * Licensed to the Apache Software Foundation (ASF) under one |
| * or more contributor license agreements. See the NOTICE file |
| * distributed with this work for additional information |
| * regarding copyright ownership. The ASF licenses this file |
| * to you under the Apache License, Version 2.0 (the |
| * "License"); you may not use this file except in compliance |
| * with the License. You may obtain a copy of the License at |
| * |
| * http://www.apache.org/licenses/LICENSE-2.0 |
| * |
| * Unless required by applicable law or agreed to in writing, |
| * software distributed under the License is distributed on an |
| * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY |
| * KIND, either express or implied. See the License for the |
| * specific language governing permissions and limitations |
| * under the License. |
| */ |
| |
| #include "paimon/core/schema/schema_validation.h" |
| |
| #include <algorithm> |
| #include <cassert> |
| #include <functional> |
| #include <map> |
| #include <optional> |
| #include <set> |
| #include <sstream> |
| #include <unordered_map> |
| #include <unordered_set> |
| #include <utility> |
| |
| #include "arrow/type.h" |
| #include "fmt/format.h" |
| #include "fmt/ranges.h" |
| #include "paimon/common/data/blob_utils.h" |
| #include "paimon/common/data/shredding/map_shared_shredding_utils.h" |
| #include "paimon/common/data/variant/variant_type_utils.h" |
| #include "paimon/common/table/special_fields.h" |
| #include "paimon/common/types/data_field.h" |
| #include "paimon/common/utils/arrow/vector_utils.h" |
| #include "paimon/common/utils/checked_cast.h" |
| #include "paimon/common/utils/object_utils.h" |
| #include "paimon/common/utils/preconditions.h" |
| #include "paimon/common/utils/string_utils.h" |
| #include "paimon/core/core_options.h" |
| #include "paimon/core/options/changelog_producer.h" |
| #include "paimon/core/options/expire_config.h" |
| #include "paimon/core/options/map_storage_layout.h" |
| #include "paimon/core/options/merge_engine.h" |
| #include "paimon/core/schema/arrow_schema_validator.h" |
| #include "paimon/core/schema/table_schema.h" |
| #include "paimon/core/table/bucket_mode.h" |
| #include "paimon/defs.h" |
| #include "paimon/format/file_format.h" |
| #include "paimon/result.h" |
| |
| namespace paimon { |
| namespace { |
| |
| bool ContainsBlobField(const std::shared_ptr<arrow::Field>& field) { |
| if (BlobUtils::IsBlobField(field)) { |
| return true; |
| } |
| const std::shared_ptr<arrow::DataType>& type = field->type(); |
| if (type->id() == arrow::Type::STRUCT) { |
| for (const auto& child : type->fields()) { |
| if (ContainsBlobField(child)) { |
| return true; |
| } |
| } |
| } else if (type->id() == arrow::Type::LIST) { |
| return ContainsBlobField(type->fields().front()); |
| } else if (type->id() == arrow::Type::MAP) { |
| const auto& map_type = checked_cast<const arrow::MapType&>(*type); |
| return ContainsBlobField(map_type.key_field()) || ContainsBlobField(map_type.item_field()); |
| } |
| return false; |
| } |
| |
| Status ValidateSharedShreddingCompression(const std::string& option_key, |
| const std::string& compression) { |
| std::string normalized = StringUtils::ToLowerCase(compression); |
| if (normalized != "none" && normalized != "lz4" && normalized != "zstd") { |
| return Status::Invalid(fmt::format( |
| "MAP shared-shredding only supports none/lz4/zstd compression, but {} is {}.", |
| option_key, compression)); |
| } |
| return Status::OK(); |
| } |
| |
| Status ValidateSharedShreddingFileFormat(const std::string& option_key, |
| const std::string& file_format) { |
| std::string normalized = StringUtils::ToLowerCase(file_format); |
| if (normalized != "parquet" && normalized != "orc") { |
| return Status::Invalid(fmt::format( |
| "MAP shared-shredding only supports parquet/orc file formats, but {} is {}.", |
| option_key, file_format)); |
| } |
| return Status::OK(); |
| } |
| |
| Status ValidateVectorFileFormat(const std::string& option_key, const std::string& file_format) { |
| if (!StringUtils::EqualsIgnoreCase(file_format, "parquet")) { |
| return Status::Invalid( |
| fmt::format("VECTOR currently only supports parquet data files, but {} is {}.", |
| option_key, file_format)); |
| } |
| return Status::OK(); |
| } |
| |
| Status ValidatePerLevelOption( |
| const std::map<std::string, std::string>& options, const std::string& option_key, |
| const std::function<Status(const std::string&, const std::string&)>& validator) { |
| auto it = options.find(option_key); |
| if (it == options.end() || it->second.empty()) { |
| return Status::OK(); |
| } |
| auto entries = StringUtils::Split(it->second, std::string(",")); |
| for (const std::string& entry : entries) { |
| auto level_and_value = StringUtils::Split(entry, std::string(":")); |
| if (level_and_value.size() == 2) { |
| PAIMON_RETURN_NOT_OK( |
| validator(option_key + "." + level_and_value[0], level_and_value[1])); |
| } |
| } |
| return Status::OK(); |
| } |
| |
| } // namespace |
| |
| bool SchemaValidation::IsComplexType(const std::shared_ptr<arrow::Field>& field) { |
| arrow::Type::type arrow_type_id = field->type()->id(); |
| return (arrow_type_id == arrow::Type::TIMESTAMP || arrow_type_id == arrow::Type::DECIMAL128 || |
| BlobUtils::IsBlobField(field)); |
| } |
| |
| Status SchemaValidation::ValidateTableSchema(const TableSchema& schema) { |
| const auto& field_names = schema.FieldNames(); |
| PAIMON_RETURN_NOT_OK(ValidateNoDuplicateField(schema.BucketKeys(), "bucket key")); |
| PAIMON_RETURN_NOT_OK(ValidateNoDuplicateField(schema.PrimaryKeys(), "primary key")); |
| PAIMON_RETURN_NOT_OK(ValidateNoDuplicateField(schema.PartitionKeys(), "partition key")); |
| PAIMON_RETURN_NOT_OK( |
| Preconditions::CheckState(ObjectUtils::ContainsAll(field_names, schema.PartitionKeys()), |
| "Table column {} should include all partition fields {}", |
| field_names, schema.PartitionKeys())); |
| PAIMON_RETURN_NOT_OK( |
| Preconditions::CheckState(ObjectUtils::ContainsAll(field_names, schema.PrimaryKeys()), |
| "Table column {} should include all primary key constraint {}", |
| field_names, schema.PrimaryKeys())); |
| |
| PAIMON_RETURN_NOT_OK( |
| ValidateOnlyContainPrimitiveType(schema.Fields(), schema.PrimaryKeys(), "primary key")); |
| PAIMON_RETURN_NOT_OK( |
| ValidateOnlyContainPrimitiveType(schema.Fields(), schema.PartitionKeys(), "partition")); |
| // TODO(lisizhuo.lsz): C++ Paimon do not support timestamp & decimal & float & double type in |
| // partition keys for now. |
| PAIMON_RETURN_NOT_OK(ValidateNotContainSpecificType(schema.Fields(), schema.PartitionKeys())); |
| |
| PAIMON_ASSIGN_OR_RAISE(CoreOptions options, CoreOptions::FromMap(schema.Options())); |
| PAIMON_RETURN_NOT_OK(ValidateBucket(schema, options)); |
| // PAIMON_RETURN_NOT_OK(ValidateDefaultValues(schema)); |
| // PAIMON_RETURN_NOT_OK(ValidateStartupMode(options)); |
| PAIMON_RETURN_NOT_OK(ValidateFieldsPrefix(schema, options)); |
| PAIMON_RETURN_NOT_OK(ValidateSequenceField(schema, options)); |
| PAIMON_RETURN_NOT_OK(ValidateSequenceGroup(schema, options)); |
| PAIMON_RETURN_NOT_OK(ValidateChangelogProducer(schema, options)); |
| |
| PAIMON_RETURN_NOT_OK(Preconditions::CheckState( |
| options.GetExpireConfig().GetSnapshotRetainMin() > 0, |
| std::string(Options::SNAPSHOT_NUM_RETAINED_MIN) + " should be at least 1")); |
| PAIMON_RETURN_NOT_OK(Preconditions::CheckState( |
| options.GetExpireConfig().GetSnapshotRetainMin() <= |
| options.GetExpireConfig().GetSnapshotRetainMax(), |
| std::string(Options::SNAPSHOT_NUM_RETAINED_MIN) + " should not be larger than " + |
| std::string(Options::SNAPSHOT_NUM_RETAINED_MAX))); |
| |
| // TODO(yonghao.fyh): check changelog num retain |
| // TODO(yonghao.fyh): support file format validate data fields |
| for (const auto& field_name : field_names) { |
| if (SpecialFields::IsSystemField(field_name)) { |
| return Status::Invalid( |
| fmt::format("field name '{}' in schema cannot be special field.", field_name)); |
| } |
| } |
| // TODO(yonghao.fyh): check streaming read overwrite |
| // TODO(yonghao.fyh): check 'partition.expiration-time' |
| // TODO(yonghao.fyh): check 'rowkind.field' |
| if (options.DeletionVectorsEnabled()) { |
| PAIMON_RETURN_NOT_OK(ValidateForDeletionVectors(options)); |
| } |
| |
| PAIMON_RETURN_NOT_OK(ValidateRowTracking(schema, options)); |
| PAIMON_RETURN_NOT_OK(ValidateBlobFields(schema, options)); |
| PAIMON_RETURN_NOT_OK(ValidateMosaicDataFields(schema, options)); |
| PAIMON_RETURN_NOT_OK(ValidateMapStorageLayout(schema, options)); |
| PAIMON_RETURN_NOT_OK(ValidateVectorFields(schema, options)); |
| return Status::OK(); |
| } |
| |
| Status SchemaValidation::ValidateNoDuplicateField(const std::vector<std::string>& field_names, |
| const std::string& error_message_intro) { |
| auto duplicate_field_names = ObjectUtils::DuplicateItems(field_names); |
| PAIMON_RETURN_NOT_OK(Preconditions::CheckState( |
| duplicate_field_names.empty(), |
| fmt::format("{} [{}] must not contain duplicate fields. Found: [{}]", error_message_intro, |
| fmt::join(field_names, ", "), fmt::join(duplicate_field_names, ", ")))); |
| return Status::OK(); |
| } |
| |
| Status SchemaValidation::ValidateOnlyContainPrimitiveType( |
| const std::vector<DataField>& fields, const std::vector<std::string>& field_names, |
| const std::string& error_message_intro) { |
| if (field_names.empty()) { |
| return Status::OK(); |
| } |
| std::unordered_map<std::string, std::shared_ptr<arrow::DataType>> fields_map; |
| for (const auto& field : fields) { |
| fields_map[field.Name()] = field.Type(); |
| } |
| for (const auto& field_name : field_names) { |
| auto it = fields_map.find(field_name); |
| if (it != fields_map.end()) { |
| auto data_type = it->second; |
| if (ArrowSchemaValidator::IsNestedType(data_type)) { |
| return Status::Invalid(fmt::format("The type {} in {} field {} is unsupported", |
| data_type->ToString(), error_message_intro, |
| it->first)); |
| } |
| } else { |
| assert(false); |
| return Status::Invalid( |
| fmt::format("unexpected error, field {} not found in fields map", field_name)); |
| } |
| } |
| return Status::OK(); |
| } |
| |
| Status SchemaValidation::ValidateNotContainSpecificType( |
| const std::vector<DataField>& fields, const std::vector<std::string>& field_names) { |
| if (field_names.empty()) { |
| return Status::OK(); |
| } |
| std::unordered_map<std::string, std::shared_ptr<arrow::Field>> fields_map; |
| for (const auto& field : fields) { |
| fields_map[field.Name()] = field.ArrowField(); |
| } |
| for (const auto& field_name : field_names) { |
| auto it = fields_map.find(field_name); |
| if (it != fields_map.end()) { |
| auto field = it->second; |
| if (IsComplexType(field)) { |
| return Status::Invalid( |
| fmt::format("partition field {} cannot be TIMESTAMP/DECIMAL/BLOB", field_name)); |
| } |
| if (field->type()->id() == arrow::Type::FLOAT || |
| field->type()->id() == arrow::Type::DOUBLE) { |
| return Status::Invalid( |
| fmt::format("partition field {} cannot be FLOAT/DOUBLE", field_name)); |
| } |
| } else { |
| assert(false); |
| return Status::Invalid(fmt::format( |
| "unexpected error, partition field {} not found in schema", field_name)); |
| } |
| } |
| return Status::OK(); |
| } |
| |
| bool SchemaValidation::IsPostponeBucketTable(const TableSchema& schema, int32_t bucket) { |
| return !schema.PrimaryKeys().empty() && bucket == BucketModeDefine::POSTPONE_BUCKET; |
| } |
| |
| Status SchemaValidation::ValidateBucket(const TableSchema& schema, const CoreOptions& options) { |
| int32_t bucket = options.GetBucket(); |
| if (bucket == -1) { |
| if (options.ToMap().count(Options::BUCKET_KEY)) { |
| return Status::Invalid( |
| fmt::format("Cannot define '{}' with bucket -1, please specify a bucket number.", |
| Options::BUCKET_KEY)); |
| } |
| if (schema.PrimaryKeys().empty() && |
| options.ToMap().count("full-compaction.delta-commits")) { |
| return Status::Invalid( |
| "AppendOnlyTable of unware or dynamic bucket does not support " |
| "'full-compaction.delta-commits'"); |
| } |
| } else if (bucket < 1 && !IsPostponeBucketTable(schema, bucket)) { |
| return Status::Invalid("The number of buckets needs to be greater than 0."); |
| } else { |
| if (schema.CrossPartitionUpdate()) { |
| return Status::Invalid(fmt::format( |
| "You should use dynamic bucket (bucket = -1) mode in cross partition update case " |
| "(Primary key constraint '{}' not include all partition fields '{}').", |
| fmt::join(schema.PrimaryKeys(), ", "), fmt::join(schema.PartitionKeys(), ", "))); |
| } |
| if (schema.PrimaryKeys().empty() && schema.BucketKeys().empty()) { |
| return Status::Invalid("You should define a 'bucket-key' for bucketed append mode."); |
| } |
| if (!schema.BucketKeys().empty()) { |
| std::vector<std::string> bucket_keys = schema.BucketKeys(); |
| std::vector<std::string> nested_fields; |
| |
| for (const auto& field : schema.Fields()) { |
| if (std::find(bucket_keys.begin(), bucket_keys.end(), field.Name()) != |
| bucket_keys.end() && |
| ArrowSchemaValidator::IsNestedType(field.Type())) { |
| nested_fields.push_back(field.Name()); |
| } |
| } |
| |
| if (!nested_fields.empty()) { |
| return Status::Invalid(fmt::format( |
| "Nested type cannot be in bucket-key, in your table these keys are: {}", |
| fmt::join(nested_fields, ", "))); |
| } |
| } |
| } |
| return Status::OK(); |
| } |
| |
| Status SchemaValidation::ValidateChangelogProducer(const TableSchema& schema, |
| const CoreOptions& options) { |
| ChangelogProducer changelog_producer = options.GetChangelogProducer(); |
| if (schema.PrimaryKeys().empty() && changelog_producer != ChangelogProducer::NONE) { |
| return Status::Invalid( |
| fmt::format("Can not set {} on table without primary keys, please define primary keys.", |
| Options::CHANGELOG_PRODUCER)); |
| } |
| |
| bool row_deduplicate = options.ChangelogRowDeduplicate(); |
| const std::vector<std::string>& ignore_fields = |
| options.GetChangelogRowDeduplicateIgnoreFields(); |
| PAIMON_RETURN_NOT_OK(Preconditions::CheckState( |
| ignore_fields.empty() || row_deduplicate, "'{}' is only valid when '{}' is true.", |
| Options::CHANGELOG_PRODUCER_ROW_DEDUPLICATE_IGNORE_FIELDS, |
| Options::CHANGELOG_PRODUCER_ROW_DEDUPLICATE)); |
| PAIMON_RETURN_NOT_OK(Preconditions::CheckState( |
| ObjectUtils::ContainsAll(schema.FieldNames(), ignore_fields), |
| "Fields {} configured in '{}' can not be found in table schema.", ignore_fields, |
| Options::CHANGELOG_PRODUCER_ROW_DEDUPLICATE_IGNORE_FIELDS)); |
| PAIMON_RETURN_NOT_OK(Preconditions::CheckState( |
| !row_deduplicate || changelog_producer == ChangelogProducer::LOOKUP || |
| changelog_producer == ChangelogProducer::FULL_COMPACTION, |
| "'{}' is only valid for 'lookup' or 'full-compaction' changelog producer.", |
| Options::CHANGELOG_PRODUCER_ROW_DEDUPLICATE)); |
| return Preconditions::CheckState( |
| options.GetMergeEngine() != MergeEngine::FIRST_ROW || |
| changelog_producer == ChangelogProducer::NONE || |
| changelog_producer == ChangelogProducer::LOOKUP, |
| "Only support 'none' and 'lookup' changelog-producer on FIRST_ROW merge engine"); |
| } |
| |
| Status SchemaValidation::ValidateForDeletionVectors(const CoreOptions& options) { |
| PAIMON_RETURN_NOT_OK(Preconditions::CheckState( |
| options.GetChangelogProducer() == ChangelogProducer::NONE || |
| options.GetChangelogProducer() == ChangelogProducer::INPUT || |
| options.GetChangelogProducer() == ChangelogProducer::LOOKUP, |
| "Deletion vectors mode is only supported for NONE/INPUT/LOOKUP changelog producer now.")); |
| return Preconditions::CheckState( |
| options.GetMergeEngine() != MergeEngine::FIRST_ROW, |
| "First row merge engine does not need deletion vectors because there is " |
| "no deletion of old data in this merge engine."); |
| } |
| |
| Status SchemaValidation::ValidateSequenceGroup(const TableSchema& schema, |
| const CoreOptions& options) { |
| std::unordered_map<std::string, std::set<std::string>> fields2_group; |
| auto sequence_groups_map = options.GetFieldsSequenceGroups(); |
| const std::vector<std::string>& field_names = schema.FieldNames(); |
| for (const auto& [k, v] : sequence_groups_map) { |
| std::vector<std::string> sequence_field_names = |
| StringUtils::Split(k, Options::FIELDS_SEPARATOR); |
| for (const auto& sequence_field_name : sequence_field_names) { |
| if (std::find(field_names.begin(), field_names.end(), sequence_field_name) == |
| field_names.end()) { |
| return Status::Invalid( |
| fmt::format("The sequence field group: {} can not be found in table schema.", |
| sequence_field_name)); |
| } |
| } |
| |
| for (const auto& field : StringUtils::Split(v, Options::FIELDS_SEPARATOR)) { |
| if (std::find(field_names.begin(), field_names.end(), field) == field_names.end()) { |
| return Status::Invalid( |
| fmt::format("Field {} can not be found in table schema.", field)); |
| } |
| |
| if (fields2_group.count(field)) { |
| std::vector<std::vector<std::string>> sequence_groups; |
| sequence_groups.emplace_back(fields2_group[field].begin(), |
| fields2_group[field].end()); |
| sequence_groups.push_back(sequence_field_names); |
| |
| std::ostringstream sequence_groups_msg; |
| for (const auto& group : sequence_groups) { |
| sequence_groups_msg << "{"; |
| for (const auto& group_field : group) { |
| sequence_groups_msg << group_field << " "; |
| } |
| sequence_groups_msg << "} "; |
| } |
| return Status::Invalid( |
| fmt::format("Field {} is defined repeatedly by multiple groups: {}.", field, |
| sequence_groups_msg.str())); |
| } |
| fields2_group[field].insert(sequence_field_names.begin(), sequence_field_names.end()); |
| } |
| } |
| |
| std::set<std::string> illegal_group; |
| for (const auto& group : fields2_group) { |
| for (const auto& field : group.second) { |
| PAIMON_ASSIGN_OR_RAISE(std::optional<std::string> agg_func, |
| options.GetFieldAggFunc(field)); |
| if (agg_func) { |
| illegal_group.insert(field); |
| } |
| } |
| } |
| |
| if (!illegal_group.empty()) { |
| std::ostringstream illegal_group_msg; |
| illegal_group_msg << "Should not define aggregation function on sequence group: "; |
| for (const auto& field : illegal_group) { |
| illegal_group_msg << field << " "; |
| } |
| return Status::Invalid(illegal_group_msg.str()); |
| } |
| return Status::OK(); |
| } |
| |
| Status SchemaValidation::ValidateSequenceField(const TableSchema& schema, |
| const CoreOptions& options) { |
| std::vector<std::string> sequence_field = options.GetSequenceField(); |
| if (!sequence_field.empty()) { |
| // Create field count map |
| std::unordered_map<std::string, int> field_count; |
| for (const auto& field : sequence_field) { |
| field_count[field]++; |
| } |
| |
| const auto& field_names = schema.FieldNames(); |
| for (const auto& field : sequence_field) { |
| PAIMON_RETURN_NOT_OK(Preconditions::CheckState( |
| std::find(field_names.begin(), field_names.end(), field) != field_names.end(), |
| fmt::format("Sequence field: '{}' cannot be found in table schema.", field))); |
| |
| PAIMON_ASSIGN_OR_RAISE(std::optional<std::string> agg_func, |
| options.GetFieldAggFunc(field)); |
| PAIMON_RETURN_NOT_OK(Preconditions::CheckState( |
| agg_func == std::nullopt, |
| fmt::format("Should not define aggregation on sequence field: '{}'.", field))); |
| |
| PAIMON_RETURN_NOT_OK(Preconditions::CheckState( |
| field_count[field] == 1, "Sequence field '" + field + "' is defined repeatedly.")); |
| } |
| |
| // Check for FIRST_ROW merge engine |
| if (options.GetMergeEngine() == MergeEngine::FIRST_ROW) { |
| return Status::Invalid( |
| "Do not support using sequence field on FIRST_ROW merge engine."); |
| } |
| |
| // Check for cross partition update |
| if (schema.CrossPartitionUpdate()) { |
| return Status::Invalid(fmt::format( |
| "You cannot use sequence.field in cross partition update case (Primary " |
| "key constraint '{}' not including all partition fields '{}').", |
| fmt::join(schema.PrimaryKeys(), ", "), fmt::join(schema.PartitionKeys(), ", "))); |
| } |
| } |
| return Status::OK(); |
| } |
| |
| Status SchemaValidation::ValidateFieldsPrefix(const TableSchema& schema, |
| const CoreOptions& options) { |
| const auto& field_names = schema.FieldNames(); |
| const auto& options_map = options.ToMap(); |
| for (const auto& [k, v] : options_map) { |
| if (StringUtils::StartsWith(k, Options::FIELDS_PREFIX)) { |
| std::vector<std::string> cols = StringUtils::Split(k, "."); |
| if (cols.size() < 2) { |
| return Status::Invalid("invalid options key " + k); |
| } |
| std::vector<std::string> fields = |
| StringUtils::Split(cols[1], Options::FIELDS_SEPARATOR); |
| for (const auto& field : fields) { |
| PAIMON_RETURN_NOT_OK(Preconditions::CheckState( |
| Options::DEFAULT_AGG_FUNCTION == field || |
| std::find(field_names.begin(), field_names.end(), field) != |
| field_names.end(), |
| "Field " + field + " can not be found in table schema.")); |
| } |
| } |
| } |
| return Status::OK(); |
| } |
| |
| Status SchemaValidation::ValidateRowTracking(const TableSchema& table_schema, |
| const CoreOptions& options) { |
| bool row_tracking_enabled = options.RowTrackingEnabled(); |
| if (row_tracking_enabled) { |
| PAIMON_RETURN_NOT_OK(Preconditions::CheckState( |
| options.GetBucket() == -1, |
| "Cannot define {} for row tracking table, it only support bucket = -1", |
| Options::BUCKET)); |
| PAIMON_RETURN_NOT_OK( |
| Preconditions::CheckState(table_schema.PrimaryKeys().empty(), |
| "Cannot define primary key for row tracking table")); |
| } |
| if (options.DataEvolutionEnabled()) { |
| PAIMON_RETURN_NOT_OK(Preconditions::CheckState( |
| row_tracking_enabled, "Data evolution config must enabled with row-tracking.enabled")); |
| } |
| |
| std::vector<std::string> blob_names; |
| for (const auto& field : table_schema.Fields()) { |
| if (BlobUtils::IsBlobField(field.ArrowField())) { |
| blob_names.push_back(field.Name()); |
| } |
| } |
| if (!blob_names.empty()) { |
| // Validate blob fields cannot be partition keys |
| for (const auto& blob_field_name : blob_names) { |
| if (std::find(table_schema.PartitionKeys().begin(), table_schema.PartitionKeys().end(), |
| blob_field_name) != table_schema.PartitionKeys().end()) { |
| return Status::Invalid( |
| fmt::format("Blob field {} cannot be a partition key.", blob_field_name)); |
| } |
| } |
| |
| // Validate data evolution must be enabled when blob-field is configured |
| PAIMON_RETURN_NOT_OK(Preconditions::CheckState( |
| options.DataEvolutionEnabled(), |
| "Data evolution config must be enabled for table with BLOB type column.")); |
| PAIMON_RETURN_NOT_OK(Preconditions::CheckState( |
| table_schema.Fields().size() > blob_names.size(), |
| "Table with BLOB type column must have other normal columns.")); |
| } |
| return Status::OK(); |
| } |
| |
| Status SchemaValidation::ValidateBlobFields(const TableSchema& schema, const CoreOptions& options) { |
| const auto& configured_blob_names = options.GetBlobFields(); |
| const auto& blob_descriptor_names = options.GetBlobDescriptorFields(); |
| const auto& blob_view_names = options.GetBlobViewFields(); |
| if (configured_blob_names.empty() && blob_descriptor_names.empty() && blob_view_names.empty()) { |
| return Status::OK(); |
| } |
| |
| auto validate_blob_fields = [&](const std::vector<std::string>& field_names, |
| const std::string& option_key) -> Status { |
| if (field_names.empty()) { |
| return Status::OK(); |
| } |
| PAIMON_RETURN_NOT_OK(ValidateNoDuplicateField(field_names, option_key)); |
| PAIMON_ASSIGN_OR_RAISE(std::vector<DataField> blob_fields, schema.GetFields(field_names)); |
| for (const auto& blob_field : blob_fields) { |
| if (!BlobUtils::IsBlobField(blob_field.ArrowField())) { |
| return Status::Invalid( |
| fmt::format("Field '{}' in '{}' must be a BLOB field in table schema.", |
| blob_field.Name(), option_key)); |
| } |
| } |
| return Status::OK(); |
| }; |
| |
| PAIMON_RETURN_NOT_OK(validate_blob_fields(configured_blob_names, Options::BLOB_FIELD)); |
| PAIMON_RETURN_NOT_OK( |
| validate_blob_fields(blob_descriptor_names, Options::BLOB_DESCRIPTOR_FIELD)); |
| PAIMON_RETURN_NOT_OK(validate_blob_fields(blob_view_names, Options::BLOB_VIEW_FIELD)); |
| |
| std::set<std::string> blob_descriptor_name_set(blob_descriptor_names.begin(), |
| blob_descriptor_names.end()); |
| for (const auto& blob_view_name : blob_view_names) { |
| if (blob_descriptor_name_set.count(blob_view_name) > 0) { |
| return Status::Invalid(fmt::format("Field '{}' in '{}' can not also be in '{}'.", |
| blob_view_name, Options::BLOB_VIEW_FIELD, |
| Options::BLOB_DESCRIPTOR_FIELD)); |
| } |
| } |
| return Status::OK(); |
| } |
| |
| Status SchemaValidation::ValidateMosaicDataField(const std::shared_ptr<arrow::Field>& field) { |
| if (VariantTypeUtils::IsVariantField(field)) { |
| return Status::Invalid("Mosaic file format does not support type VARIANT"); |
| } |
| if (BlobUtils::IsBlobField(field)) { |
| return Status::Invalid("Mosaic file format does not support type BLOB"); |
| } |
| |
| const std::shared_ptr<arrow::DataType>& type = field->type(); |
| switch (type->id()) { |
| case arrow::Type::BOOL: |
| case arrow::Type::INT8: |
| case arrow::Type::INT16: |
| case arrow::Type::INT32: |
| case arrow::Type::INT64: |
| case arrow::Type::FLOAT: |
| case arrow::Type::DOUBLE: |
| case arrow::Type::DATE32: |
| case arrow::Type::STRING: |
| case arrow::Type::BINARY: |
| case arrow::Type::TIME32: |
| case arrow::Type::DECIMAL128: |
| return Status::OK(); |
| case arrow::Type::TIMESTAMP: { |
| const auto& timestamp_type = checked_cast<const arrow::TimestampType&>(*type); |
| if (timestamp_type.unit() == arrow::TimeUnit::SECOND) { |
| return Status::Invalid("Mosaic file format does not support TIMESTAMP(0)"); |
| } |
| return Status::OK(); |
| } |
| case arrow::Type::LIST: |
| return ValidateMosaicDataField(type->field(0)); |
| case arrow::Type::MAP: { |
| const auto& map_type = checked_cast<const arrow::MapType&>(*type); |
| PAIMON_RETURN_NOT_OK(ValidateMosaicDataField(map_type.key_field())); |
| return ValidateMosaicDataField(map_type.item_field()); |
| } |
| case arrow::Type::FIXED_SIZE_LIST: |
| return Status::Invalid("Mosaic file format does not support type VECTOR"); |
| case arrow::Type::STRUCT: |
| return Status::Invalid("Mosaic file format does not support type ROW"); |
| default: |
| break; |
| } |
| return Status::Invalid( |
| fmt::format("Mosaic file format does not support type {}", type->ToString())); |
| } |
| |
| Status SchemaValidation::ValidateMosaicDataFields(const TableSchema& schema, |
| const CoreOptions& options) { |
| if (StringUtils::ToLowerCase(options.GetFileFormat()->Identifier()) != "mosaic") { |
| return Status::OK(); |
| } |
| |
| const std::vector<std::string> inline_blob_fields = options.GetBlobInlineFields(); |
| const std::set<std::string> inline_blob_field_set(inline_blob_fields.begin(), |
| inline_blob_fields.end()); |
| // Match Java SchemaValidation by validating only fields stored in the normal data file. C++ |
| // permits BLOB only as a top-level field; descriptor and view fields are inline, so Mosaic |
| // must reject them here. |
| for (const DataField& field : schema.Fields()) { |
| if (BlobUtils::IsBlobField(field.ArrowField()) && |
| inline_blob_field_set.count(field.Name()) == 0) { |
| continue; |
| } |
| PAIMON_RETURN_NOT_OK(ValidateMosaicDataField(field.ArrowField())); |
| } |
| return Status::OK(); |
| } |
| |
| Status SchemaValidation::ValidateMapStorageLayout(const TableSchema& schema, |
| const CoreOptions& options) { |
| // Extract all field names that have map.storage-layout configured from options |
| const std::string layout_suffix = std::string(".") + std::string(Options::MAP_STORAGE_LAYOUT); |
| const auto& options_map = options.ToMap(); |
| |
| std::unordered_map<std::string, std::shared_ptr<arrow::DataType>> schema_fields; |
| for (const auto& field : schema.Fields()) { |
| schema_fields[field.Name()] = field.Type(); |
| } |
| |
| std::string fields_prefix_str = std::string(Options::FIELDS_PREFIX); |
| bool has_shared_shredding = false; |
| for (const auto& [key, value] : options_map) { |
| if (!StringUtils::StartsWith(key, fields_prefix_str)) { |
| continue; |
| } |
| if (!StringUtils::EndsWith(key, layout_suffix)) { |
| continue; |
| } |
| // key = "fields.<field_name>.map.storage-layout" |
| // Extract field_name: skip "fields." prefix and ".map.storage-layout" suffix |
| std::string field_name = |
| key.substr(fields_prefix_str.size() + 1, |
| key.size() - fields_prefix_str.size() - 1 - layout_suffix.size()); |
| |
| // Check field exists in schema |
| auto it = schema_fields.find(field_name); |
| if (it == schema_fields.end()) { |
| return Status::Invalid( |
| fmt::format("Column '{}' is configured with map.storage-layout " |
| "but does not exist in table schema.", |
| field_name)); |
| } |
| |
| // Any column configured with map.storage-layout must be a MAP type |
| const auto& field_type = it->second; |
| if (field_type->id() != arrow::Type::MAP) { |
| return Status::Invalid( |
| fmt::format("Column '{}' is configured with map.storage-layout " |
| "but its type is not MAP.", |
| field_name)); |
| } |
| |
| PAIMON_ASSIGN_OR_RAISE(MapStorageLayout layout, options.GetMapStorageLayout(field_name)); |
| if (layout != MapStorageLayout::SHARED_SHREDDING) { |
| continue; |
| } |
| has_shared_shredding = true; |
| for (const auto& field : schema.Fields()) { |
| if (VariantTypeUtils::ContainsVariantField(field.ArrowField())) { |
| return Status::Invalid( |
| "MAP shared-shredding currently cannot be used with Variant fields."); |
| } |
| } |
| // Column configured with shared-shredding must be MAP<STRING, T> |
| if (!MapSharedShreddingUtils::IsShreddingKeyMap(field_type)) { |
| return Status::Invalid( |
| fmt::format("Column '{}' is configured with map.storage-layout=shared-shredding " |
| "but its type is not MAP<STRING NOT NULL, T>.", |
| field_name)); |
| } |
| auto map_type = checked_pointer_cast<arrow::MapType>(field_type); |
| if (map_type->key_field()->nullable()) { |
| return Status::Invalid( |
| fmt::format("Column '{}' is configured with map.storage-layout=shared-shredding " |
| "but its map key type is nullable.", |
| field_name)); |
| } |
| if (ContainsBlobField(map_type->item_field())) { |
| return Status::Invalid("MAP shared-shredding currently cannot contain BLOB fields."); |
| } |
| if (VectorUtils::ContainsVectorField(map_type->item_field())) { |
| return Status::Invalid("MAP shared-shredding currently cannot contain VECTOR fields."); |
| } |
| // Validate max-columns config |
| PAIMON_RETURN_NOT_OK(options.GetMapSharedShreddingMaxColumns(field_name)); |
| // Validate placement policy config |
| PAIMON_RETURN_NOT_OK(options.GetMapSharedShreddingColumnPlacementPolicy(field_name)); |
| } |
| if (!has_shared_shredding) { |
| return Status::OK(); |
| } |
| |
| if (IsPostponeBucketTable(schema, options.GetBucket())) { |
| return Status::Invalid( |
| "MAP shared-shredding currently does not support postpone bucket mode."); |
| } |
| |
| PAIMON_RETURN_NOT_OK(ValidateSharedShreddingFileFormat(Options::FILE_FORMAT, |
| options.GetFileFormat()->Identifier())); |
| PAIMON_RETURN_NOT_OK(ValidatePerLevelOption(options_map, Options::FILE_FORMAT_PER_LEVEL, |
| ValidateSharedShreddingFileFormat)); |
| std::shared_ptr<FileFormat> changelog_format = options.GetChangelogFileFormat(); |
| if (changelog_format) { |
| PAIMON_RETURN_NOT_OK(ValidateSharedShreddingFileFormat(Options::CHANGELOG_FILE_FORMAT, |
| changelog_format->Identifier())); |
| } |
| PAIMON_RETURN_NOT_OK(ValidateSharedShreddingCompression(Options::FILE_COMPRESSION, |
| options.GetFileCompression())); |
| PAIMON_RETURN_NOT_OK(ValidatePerLevelOption(options_map, Options::FILE_COMPRESSION_PER_LEVEL, |
| ValidateSharedShreddingCompression)); |
| std::optional<std::string> changelog_compression = options.GetChangelogFileCompression(); |
| if (changelog_compression) { |
| PAIMON_RETURN_NOT_OK(ValidateSharedShreddingCompression(Options::CHANGELOG_FILE_COMPRESSION, |
| changelog_compression.value())); |
| } |
| |
| return Status::OK(); |
| } |
| |
| Status SchemaValidation::ValidateVectorFields(const TableSchema& schema, |
| const CoreOptions& options) { |
| bool has_vector = false; |
| for (const auto& field : schema.Fields()) { |
| if (VectorUtils::ContainsVectorField(field.ArrowField())) { |
| has_vector = true; |
| break; |
| } |
| } |
| if (!has_vector) { |
| return Status::OK(); |
| } |
| if (!schema.PrimaryKeys().empty()) { |
| return Status::NotImplemented( |
| "VECTOR fields in primary-key tables are not implemented yet."); |
| } |
| if (options.DataEvolutionEnabled()) { |
| return Status::NotImplemented( |
| "VECTOR fields in data-evolution tables are not implemented yet."); |
| } |
| PAIMON_RETURN_NOT_OK( |
| ValidateVectorFileFormat(Options::FILE_FORMAT, options.GetFileFormat()->Identifier())); |
| return ValidatePerLevelOption(options.ToMap(), Options::FILE_FORMAT_PER_LEVEL, |
| ValidateVectorFileFormat); |
| } |
| |
| } // namespace paimon |