diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.cpp index 945303783fe7..bd8818847d4c 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.cpp @@ -828,10 +828,35 @@ void IcebergMetadata::checkAlterIsPossible(const AlterCommands & commands) ErrorCodes::NOT_IMPLEMENTED, "Removing column property '{}' from column '{}' is not supported by Iceberg storage", command.to_remove, command.column_name); - if (command.type == AlterCommand::Type::MODIFY_COLUMN && !command.data_type) + if (command.type == AlterCommand::Type::ADD_COLUMN || command.type == AlterCommand::Type::MODIFY_COLUMN) + { + /// The Iceberg schema records only the type and the field order, so any other clause would be dropped. + const char * unsupported_property = nullptr; + if (command.comment) + unsupported_property = "comment"; + else if (command.default_expression) + unsupported_property = "default expression"; + else if (command.codec) + unsupported_property = "codec"; + else if (command.ttl) + unsupported_property = "TTL"; + else if (!command.settings_changes.empty() || !command.settings_resets.empty()) + unsupported_property = "settings"; + + if (unsupported_property) + throw Exception( + ErrorCodes::NOT_IMPLEMENTED, + "{} the {} of column '{}' is not supported by Iceberg storage", + command.type == AlterCommand::Type::ADD_COLUMN ? "Setting" : "Changing", + unsupported_property, + command.column_name); + } + + if (command.type == AlterCommand::Type::MODIFY_COLUMN && !command.data_type + && !command.first && command.after_column.empty()) throw Exception( ErrorCodes::NOT_IMPLEMENTED, - "Modifying column '{}' without changing its type is not supported by Iceberg storage", command.column_name); + "Modifying column '{}' without changing its type or position is not supported by Iceberg storage", command.column_name); } } diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/MetadataGenerator.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/MetadataGenerator.cpp index 7c232a26b8fe..c5a4762d7f34 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/MetadataGenerator.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/MetadataGenerator.cpp @@ -1,5 +1,4 @@ #include -#include #include #include #include @@ -12,6 +11,7 @@ #include #include +#include #include #include @@ -161,46 +161,43 @@ bool checkValidSchemaEvolution(Poco::Dynamic::Var old_type, Poco::Dynamic::Var n } -/// Recursively drop the field ids Iceberg assigns to nested elements of a complex type. -/// `getIcebergType` allocates them from a running counter, so regenerating the same -/// ClickHouse type with a different counter start yields a different - but structurally -/// identical - descriptor. Removing the ids makes such descriptors comparable. -void stripNestedFieldIds(Poco::JSON::Object::Ptr type_object) +String removeASCIIWhitespace(const String & s) { - for (const auto & id_field : {Iceberg::f_id, Iceberg::f_element_id, Iceberg::f_key_id, Iceberg::f_value_id}) - type_object->remove(id_field); - - for (const auto & nested_field : {Iceberg::f_element, Iceberg::f_key, Iceberg::f_value, Iceberg::f_type}) + String result; + result.reserve(s.size()); + for (char c : s) { - if (!type_object->has(nested_field)) - continue; - auto nested = type_object->get(nested_field); - if (nested.isString()) - continue; - if (auto nested_object = nested.extract()) - stripNestedFieldIds(nested_object); + if (!isWhitespaceASCII(c)) + result.push_back(c); } + return result; +} - if (type_object->has(Iceberg::f_fields)) - { - auto fields = type_object->getArray(Iceberg::f_fields); - for (UInt32 i = 0; i < fields->size(); ++i) - { - if (auto field = fields->getObject(i)) - stripNestedFieldIds(field); - } - } +/// A missing key matches only a missing key. +bool optionalBoolsEqual(const Poco::JSON::Object::Ptr & first, const Poco::JSON::Object::Ptr & second, const char * key) +{ + const bool first_has = first->has(key); + if (first_has != second->has(key)) + return false; + return !first_has || first->getValue(key) == second->getValue(key); } -/// Like `icebergTypesEqual`, but ignores the field ids embedded in complex types. -/// Used to recognize a type that a previous attempt already wrote, where the ids -/// were allocated from a lower `last-column-id` than the one we would use now. +/// Whether two Iceberg type descriptors denote the same type. +/// Only the properties that define the type are compared. Field ids are ignored because +/// `getIcebergType` allocates them from a running counter, so a previous attempt may have written +/// the very same type with lower ids. Other keys (`doc`, defaults, the `required` that +/// `getIcebergType` puts on a list) are ignored because other writers emit or omit them freely. bool icebergTypesEqualIgnoringIds(Poco::Dynamic::Var old_type, Poco::Dynamic::Var new_type) { - if (old_type.isString() && new_type.isString()) - return old_type.extract() == new_type.extract(); - if (old_type.isString() || new_type.isString()) + { + if (!old_type.isString() || !new_type.isString()) + return false; + /// Writers differ in spacing, e.g. `decimal(10,2)` and `decimal(10, 2)`. + return removeASCIIWhitespace(old_type.extract()) == removeASCIIWhitespace(new_type.extract()); + } + + if (old_type.type() != typeid(Poco::JSON::Object::Ptr) || new_type.type() != typeid(Poco::JSON::Object::Ptr)) return false; auto old_object = old_type.extract(); @@ -208,16 +205,52 @@ bool icebergTypesEqualIgnoringIds(Poco::Dynamic::Var old_type, Poco::Dynamic::Va if (!old_object || !new_object) return false; - auto old_stripped = deepCopy(old_object); - auto new_stripped = deepCopy(new_object); - stripNestedFieldIds(old_stripped); - stripNestedFieldIds(new_stripped); + if (old_object->has("precision") || new_object->has("precision")) + { + return old_object->has("precision") && new_object->has("precision") + && old_object->getValue("precision") == new_object->getValue("precision") + && old_object->getValue("scale") == new_object->getValue("scale"); + } + + if (!old_object->has(Iceberg::f_type) || !new_object->has(Iceberg::f_type)) + return false; + const auto kind = old_object->getValue(Iceberg::f_type); + if (kind != new_object->getValue(Iceberg::f_type)) + return false; + + if (kind == Iceberg::f_list) + { + return optionalBoolsEqual(old_object, new_object, Iceberg::f_element_required) + && icebergTypesEqualIgnoringIds(old_object->get(Iceberg::f_element), new_object->get(Iceberg::f_element)); + } + + if (kind == Iceberg::f_map) + { + return optionalBoolsEqual(old_object, new_object, Iceberg::f_value_required) + && icebergTypesEqualIgnoringIds(old_object->get(Iceberg::f_key), new_object->get(Iceberg::f_key)) + && icebergTypesEqualIgnoringIds(old_object->get(Iceberg::f_value), new_object->get(Iceberg::f_value)); + } + + if (kind == Iceberg::f_struct) + { + auto old_fields = old_object->getArray(Iceberg::f_fields); + auto new_fields = new_object->getArray(Iceberg::f_fields); + if (!old_fields || !new_fields || old_fields->size() != new_fields->size()) + return false; + for (UInt32 i = 0; i < old_fields->size(); ++i) + { + auto old_field = old_fields->getObject(i); + auto new_field = new_fields->getObject(i); + if (!old_field || !new_field + || old_field->getValue(Iceberg::f_name) != new_field->getValue(Iceberg::f_name) + || !optionalBoolsEqual(old_field, new_field, Iceberg::f_required) + || !icebergTypesEqualIgnoringIds(old_field->get(Iceberg::f_type), new_field->get(Iceberg::f_type))) + return false; + } + return true; + } - std::ostringstream oss_old; // STYLE_CHECK_ALLOW_STD_STRING_STREAM - std::ostringstream oss_new; // STYLE_CHECK_ALLOW_STD_STRING_STREAM - old_stripped->stringify(oss_old); - new_stripped->stringify(oss_new); - return oss_old.str() == oss_new.str(); + return false; } /// The Iceberg spec marks `snapshots`, `metadata-log` and `snapshot-log` as optional, so table @@ -246,6 +279,18 @@ Int32 getNextSchemaId(Poco::JSON::Object::Ptr metadata_object) return max_id + 1; } +/// Whether the field at `index` of `fields` satisfies the `FIRST` / `AFTER after_column` clause of an ALTER. +/// `AFTER` the column itself does not move it, so any position satisfies it. +bool isFieldAtRequestedPosition( + const Poco::JSON::Array::Ptr & fields, UInt32 index, const String & column_name, bool first, const String & after_column) +{ + if (first) + return index == 0; + if (!after_column.empty() && after_column != column_name) + return index > 0 && fields->getObject(index - 1)->getValue(Iceberg::f_name) == after_column; + return true; +} + } MetadataGenerator::MetadataGenerator(Poco::JSON::Object::Ptr metadata_object_) @@ -300,7 +345,7 @@ Poco::JSON::Object::Ptr MetadataGenerator::getCurrentSchema() const return current_schema; } -bool MetadataGenerator::isAddColumnApplied(const String & column_name, DataTypePtr type) const +bool MetadataGenerator::isAddColumnApplied(const String & column_name, DataTypePtr type, bool first, const String & after_column) const { auto current_schema = findCurrentSchema(); if (!current_schema) @@ -318,7 +363,8 @@ bool MetadataGenerator::isAddColumnApplied(const String & column_name, DataTypeP /// The stored descriptor was produced from a lower `last-column-id` than the one we /// just used, so the ids of nested elements differ even for the very same type. return field->getValue(Iceberg::f_required) == expected_type.second - && icebergTypesEqualIgnoringIds(field->get(Iceberg::f_type), expected_type.first); + && icebergTypesEqualIgnoringIds(field->get(Iceberg::f_type), expected_type.first) + && isFieldAtRequestedPosition(fields, i, column_name, first, after_column); } return false; } @@ -357,23 +403,30 @@ bool MetadataGenerator::isRenameColumnApplied(const String & column_name, const return found_new_name; } -bool MetadataGenerator::isModifyColumnApplied(const String & column_name, DataTypePtr type) const +bool MetadataGenerator::isModifyColumnApplied(const String & column_name, DataTypePtr type, bool first, const String & after_column) const { auto current_schema = findCurrentSchema(); if (!current_schema) return false; - Int32 unused_field_id = metadata_object->getValue(Iceberg::f_last_column_id); - auto expected_type = Iceberg::getIcebergType(type, unused_field_id); - auto fields = current_schema->getArray(Iceberg::f_fields); for (UInt32 i = 0; i < fields->size(); ++i) { auto field = fields->getObject(i); if (field->getValue(Iceberg::f_name) != column_name) continue; - return field->getValue(Iceberg::f_required) == expected_type.second - && icebergTypesEqualIgnoringIds(field->get(Iceberg::f_type), expected_type.first); + + /// A position-only `MODIFY COLUMN c FIRST` carries no type. + if (type) + { + Int32 unused_field_id = metadata_object->getValue(Iceberg::f_last_column_id); + auto expected_type = Iceberg::getIcebergType(type, unused_field_id); + if (field->getValue(Iceberg::f_required) != expected_type.second + || !icebergTypesEqualIgnoringIds(field->get(Iceberg::f_type), expected_type.first)) + return false; + } + + return isFieldAtRequestedPosition(fields, i, column_name, first, after_column); } return false; } @@ -722,7 +775,7 @@ void MetadataGenerator::generateDropColumnMetadata(const String & column_name) metadata_object->getArray(Iceberg::f_schemas)->add(current_schema); } -void MetadataGenerator::generateAddColumnMetadata(const String & column_name, DataTypePtr type) +void MetadataGenerator::generateAddColumnMetadata(const String & column_name, DataTypePtr type, bool first, const String & after_column) { if (!type->isNullable()) throw Exception(ErrorCodes::BAD_ARGUMENTS, "Iceberg spec doesn't allow to add non-nullable columns"); @@ -748,81 +801,160 @@ void MetadataGenerator::generateAddColumnMetadata(const String & column_name, Da metadata_object->set(Iceberg::f_last_column_id, last_column_id + 1); - current_schema->getArray(Iceberg::f_fields)->add(new_field); + if (first || !after_column.empty()) + { + Poco::JSON::Array::Ptr new_fields = new Poco::JSON::Array; + if (first) + { + new_fields->add(new_field); + for (UInt32 i = 0; i < existing_fields->size(); ++i) + new_fields->add(existing_fields->get(i)); + } + else + { + bool inserted = false; + for (UInt32 i = 0; i < existing_fields->size(); ++i) + { + new_fields->add(existing_fields->get(i)); + if (existing_fields->getObject(i)->getValue(Iceberg::f_name) == after_column) + { + new_fields->add(new_field); + inserted = true; + } + } + if (!inserted) + throw Exception(ErrorCodes::BAD_ARGUMENTS, "Column {} not found for AFTER positioning", after_column); + } + current_schema->set(Iceberg::f_fields, new_fields); + } + else + { + existing_fields->add(new_field); + } + current_schema->set(Iceberg::f_schema_id, next_schema_id); metadata_object->set(Iceberg::f_current_schema_id, next_schema_id); metadata_object->getArray(Iceberg::f_schemas)->add(current_schema); } -bool MetadataGenerator::generateModifyColumnMetadata(const String & column_name, DataTypePtr type, ContextPtr context) +bool MetadataGenerator::generateModifyColumnMetadata(const String & column_name, DataTypePtr type, ContextPtr context, bool first, const String & after_column) { auto current_schema = getCurrentSchema(); auto last_column_id = metadata_object->getValue(Iceberg::f_last_column_id); - auto new_type = Iceberg::getIcebergType(type, last_column_id); auto schema_fields = current_schema->getArray(Iceberg::f_fields); - for (UInt32 i = 0; i < schema_fields->size(); ++i) + /// `AFTER` the column itself keeps its position, as in `ColumnsDescription::modifyColumnOrder`. + bool needs_reposition = first || (!after_column.empty() && after_column != column_name); + bool type_changed = false; + + if (type) { - auto current_field = schema_fields->getObject(i); - if (current_field->getValue(Iceberg::f_name) == column_name) + auto new_type = Iceberg::getIcebergType(type, last_column_id); + + for (UInt32 i = 0; i < schema_fields->size(); ++i) { + auto current_field = schema_fields->getObject(i); + if (current_field->getValue(Iceberg::f_name) != column_name) + continue; + if (current_field->getValue(Iceberg::f_required) == new_type.second && icebergTypesEqualIgnoringIds(current_field->get(Iceberg::f_type), new_type.first)) { - auto existing_iceberg_type = current_field->get(Iceberg::f_type); - if (existing_iceberg_type.isString()) - { - auto reconstructed_ch_type = Iceberg::IcebergSchemaProcessor::getSimpleType( - existing_iceberg_type.extract(), - context, - context->getSettingsRef()[Setting::allow_experimental_geo_types_in_iceberg]); - if (!current_field->getValue(Iceberg::f_required) && reconstructed_ch_type->canBeInsideNullable()) - reconstructed_ch_type = makeNullable(reconstructed_ch_type); - - if (reconstructed_ch_type->equals(*type)) - return false; + Iceberg::IcebergSchemaProcessor schema_processor( + context, context->getSettingsRef()[Setting::allow_experimental_geo_types_in_iceberg]); + auto current_ch_type = schema_processor.getClickHouseFieldType(current_field, context); + if (!current_ch_type->equals(*type)) throw Exception( ErrorCodes::BAD_ARGUMENTS, "Cannot MODIFY COLUMN '{}' from {} to {}: both map to the same Iceberg type '{}' " "so the change cannot be recorded in the Iceberg schema", column_name, - reconstructed_ch_type->getName(), + current_ch_type->getName(), type->getName(), - existing_iceberg_type.extract()); - } + current_field->get(Iceberg::f_type).toString()); - throw Exception( - ErrorCodes::BAD_ARGUMENTS, - "Cannot MODIFY COLUMN '{}': the requested and existing types both map to the same " - "Iceberg complex type, and the change cannot be recorded in the Iceberg schema", - column_name); + if (!needs_reposition) + return false; } + else + { + if (!checkValidSchemaEvolution(current_field->get(Iceberg::f_type), new_type.first)) + throw Exception(ErrorCodes::BAD_ARGUMENTS, "Iceberg spec doesn't allow schema evolution to type {}", type->getPrettyName()); + + if (!current_field->getValue(Iceberg::f_required) && !type->isNullable()) + throw Exception(ErrorCodes::BAD_ARGUMENTS, "Iceberg spec doesn't allow change type from nullable to non-nullable {}", type->getPrettyName()); - if (!checkValidSchemaEvolution(current_field->get(Iceberg::f_type), new_type.first)) - throw Exception(ErrorCodes::BAD_ARGUMENTS, "Iceberg spec doesn't allow schema evolution to type {}", type->getPrettyName()); + type_changed = true; + } + break; + } + } - if (!current_field->getValue(Iceberg::f_required) && !type->isNullable()) - throw Exception(ErrorCodes::BAD_ARGUMENTS, "Iceberg spec doesn't allow change type from nullable to non-nullable {}", type->getPrettyName()); + UInt32 target_index = static_cast(schema_fields->size()); + for (UInt32 i = 0; i < schema_fields->size(); ++i) + { + if (schema_fields->getObject(i)->getValue(Iceberg::f_name) == column_name) + { + target_index = i; + break; + } + } + if (target_index == schema_fields->size()) + throw Exception(ErrorCodes::BAD_ARGUMENTS, "Column {} not found in schema", column_name); - const auto next_schema_id = getNextSchemaId(metadata_object); + if (!type_changed && !needs_reposition) + return false; - current_schema = deepCopy(current_schema); - schema_fields = current_schema->getArray(Iceberg::f_fields); - current_field = schema_fields->getObject(i); + const auto next_schema_id = getNextSchemaId(metadata_object); + current_schema = deepCopy(current_schema); + schema_fields = current_schema->getArray(Iceberg::f_fields); + auto target_field = schema_fields->getObject(target_index); - current_field->set(Iceberg::f_type, new_type.first); - current_field->set(Iceberg::f_required, new_type.second); + if (type_changed) + { + auto new_type = Iceberg::getIcebergType(type, last_column_id); + target_field->set(Iceberg::f_type, new_type.first); + target_field->set(Iceberg::f_required, new_type.second); + } - metadata_object->set(Iceberg::f_current_schema_id, next_schema_id); - current_schema->set(Iceberg::f_schema_id, next_schema_id); - metadata_object->getArray(Iceberg::f_schemas)->add(current_schema); - return true; + if (needs_reposition) + { + Poco::JSON::Array::Ptr new_fields = new Poco::JSON::Array; + if (first) + { + new_fields->add(schema_fields->get(target_index)); + for (UInt32 i = 0; i < schema_fields->size(); ++i) + { + if (i != target_index) + new_fields->add(schema_fields->get(i)); + } } + else + { + bool inserted = false; + for (UInt32 i = 0; i < schema_fields->size(); ++i) + { + if (i == target_index) + continue; + new_fields->add(schema_fields->get(i)); + if (schema_fields->getObject(i)->getValue(Iceberg::f_name) == after_column) + { + new_fields->add(schema_fields->get(target_index)); + inserted = true; + } + } + if (!inserted) + throw Exception(ErrorCodes::BAD_ARGUMENTS, "Column {} not found for AFTER positioning", after_column); + } + current_schema->set(Iceberg::f_fields, new_fields); } - throw Exception(ErrorCodes::BAD_ARGUMENTS, "Column {} not found in schema", column_name); + metadata_object->set(Iceberg::f_current_schema_id, next_schema_id); + current_schema->set(Iceberg::f_schema_id, next_schema_id); + metadata_object->getArray(Iceberg::f_schemas)->add(current_schema); + return true; } void MetadataGenerator::generateRenameColumnMetadata(const String & column_name, const String & new_column_name) diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/MetadataGenerator.h b/src/Storages/ObjectStorage/DataLakes/Iceberg/MetadataGenerator.h index 33ff405a65d7..01d41ca2749a 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/MetadataGenerator.h +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/MetadataGenerator.h @@ -41,28 +41,31 @@ class MetadataGenerator std::optional user_defined_timestamp = std::nullopt, bool is_truncate = false); + void generateAddColumnMetadata(const String & column_name, DataTypePtr type, bool first = false, const String & after_column = {}); /// Create a manifest-only rewrite snapshot (`replace` operation) carrying `total-*` counters forward so `OPTIMIZE ... MANIFEST` is idempotent. NextMetadataResult generateManifestOnlySnapshot( FileNamesGenerator & generator, const Iceberg::IcebergPathFromMetadata & metadata_file_path, Int64 parent_snapshot_id); - void generateAddColumnMetadata(const String & column_name, DataTypePtr type); void generateDropColumnMetadata(const String & column_name); - /// Returns false when the column already has the requested type (no metadata change). + /// Returns false when neither the type nor the position changed (true no-op). + /// Throws when the requested type differs from the current one but maps to the same Iceberg type. /// `context` supplies the settings used to map the stored Iceberg type back to a ClickHouse /// type (the timestamptz timezone and whether geo types are allowed). - bool generateModifyColumnMetadata(const String & column_name, DataTypePtr type, ContextPtr context); + bool generateModifyColumnMetadata(const String & column_name, DataTypePtr type, ContextPtr context, bool first = false, const String & after_column = {}); void generateRenameColumnMetadata(const String & column_name, const String & new_column_name); /// A commit attempt can land in the catalog even when the client observes a failure /// (the Iceberg "commit state unknown" case, e.g. a proxy returning 5xx after the catalog /// applied the update). These predicates let a retry detect that the requested change is /// already present instead of applying it a second time and failing. - bool isAddColumnApplied(const String & column_name, DataTypePtr type) const; + /// `first` and `after_column` are checked against the field order, as for `isModifyColumnApplied`. + bool isAddColumnApplied(const String & column_name, DataTypePtr type, bool first = false, const String & after_column = {}) const; bool isDropColumnApplied(const String & column_name) const; bool isRenameColumnApplied(const String & column_name, const String & new_column_name) const; - bool isModifyColumnApplied(const String & column_name, DataTypePtr type) const; + /// `type` may be null for a position-only `MODIFY COLUMN`; `first` and `after_column` are checked against the field order. + bool isModifyColumnApplied(const String & column_name, DataTypePtr type, bool first = false, const String & after_column = {}) const; private: Poco::JSON::Object::Ptr metadata_object; diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/Mutations.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/Mutations.cpp index 2933f50f9296..9478acb06b04 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/Mutations.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/Mutations.cpp @@ -109,13 +109,13 @@ static bool alterAlreadyApplied(const MetadataGenerator & generator, const Alter switch (command.type) { case AlterCommand::Type::ADD_COLUMN: - return generator.isAddColumnApplied(command.column_name, command.data_type); + return generator.isAddColumnApplied(command.column_name, command.data_type, command.first, command.after_column); case AlterCommand::Type::DROP_COLUMN: return generator.isDropColumnApplied(command.column_name); case AlterCommand::Type::RENAME_COLUMN: return generator.isRenameColumnApplied(command.column_name, command.rename_to); case AlterCommand::Type::MODIFY_COLUMN: - return generator.isModifyColumnApplied(command.column_name, command.data_type); + return generator.isModifyColumnApplied(command.column_name, command.data_type, command.first, command.after_column); default: return false; } @@ -900,7 +900,7 @@ void alter( switch (params[0].type) { case AlterCommand::Type::ADD_COLUMN: - metadata_json_generator.generateAddColumnMetadata(params[0].column_name, params[0].data_type); + metadata_json_generator.generateAddColumnMetadata(params[0].column_name, params[0].data_type, params[0].first, params[0].after_column); break; case AlterCommand::Type::DROP_COLUMN: if (params[0].clear) @@ -909,7 +909,7 @@ void alter( break; case AlterCommand::Type::MODIFY_COLUMN: { - if (!metadata_json_generator.generateModifyColumnMetadata(params[0].column_name, params[0].data_type, context)) + if (!metadata_json_generator.generateModifyColumnMetadata(params[0].column_name, params[0].data_type, context, params[0].first, params[0].after_column)) { succeeded = true; } diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/SchemaProcessor.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/SchemaProcessor.cpp index f6c989b13ad3..a785229d911b 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/SchemaProcessor.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/SchemaProcessor.cpp @@ -563,6 +563,11 @@ DataTypePtr IcebergSchemaProcessor::getFieldType( throw Exception(ErrorCodes::BAD_ARGUMENTS, "Unexpected 'type' field: {}", type.toString()); } +DataTypePtr IcebergSchemaProcessor::getClickHouseFieldType(const Poco::JSON::Object::Ptr & field, ContextPtr context_) +{ + return getFieldType(field, f_type, context_, field->getValue(f_required)); +} + /** * Iceberg allows only three types of primitive type conversion: * int -> long diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/SchemaProcessor.h b/src/Storages/ObjectStorage/DataLakes/Iceberg/SchemaProcessor.h index 8533614fe799..2c29579fc9b0 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/SchemaProcessor.h +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/SchemaProcessor.h @@ -97,6 +97,9 @@ class IcebergSchemaProcessor : private WithContext static DataTypePtr getSimpleType(const String & type_name, ContextPtr context_, bool allow_geo_parser = true); + /// The ClickHouse type the reader derives for a top-level schema field. + DataTypePtr getClickHouseFieldType(const Poco::JSON::Object::Ptr & field, ContextPtr context_); + static std::unordered_map traverseSchema(Poco::JSON::Array::Ptr schema); void registerSnapshotWithSchemaId(Int64 snapshot_id, Int32 schema_id); diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/tests/gtest_iceberg_metadata_generator.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/tests/gtest_iceberg_metadata_generator.cpp index 1f2c5bec0cb3..d25e19e2f42f 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/tests/gtest_iceberg_metadata_generator.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/tests/gtest_iceberg_metadata_generator.cpp @@ -6,6 +6,7 @@ #include #include +#include #include #include #include @@ -367,6 +368,26 @@ Poco::JSON::Object::Ptr makeMetadataWithField( return metadata; } +/// Ordered field names from the schema `current-schema-id` points at. +std::vector getCurrentFieldNames(const Poco::JSON::Object::Ptr & metadata) +{ + auto current_schema_id = metadata->getValue(f_current_schema_id); + auto schemas = metadata->getArray(f_schemas); + for (UInt32 i = 0; i < schemas->size(); ++i) + { + auto schema = schemas->getObject(i); + if (schema->getValue(f_schema_id) != current_schema_id) + continue; + auto fields = schema->getArray(f_fields); + std::vector names; + names.reserve(fields->size()); + for (UInt32 j = 0; j < fields->size(); ++j) + names.push_back(fields->getObject(j)->getValue(f_name)); + return names; + } + return {}; +} + /// The Iceberg type recorded for `name` in the schema `current-schema-id` points at. Poco::Dynamic::Var findCurrentFieldType(const Poco::JSON::Object::Ptr & metadata, const String & name) { @@ -389,13 +410,17 @@ Poco::Dynamic::Var findCurrentFieldType(const Poco::JSON::Object::Ptr & metadata } void expectModifyRejected( - const Poco::JSON::Object::Ptr & metadata, const String & column, const DataTypePtr & requested_type) + const Poco::JSON::Object::Ptr & metadata, + const String & column, + const DataTypePtr & requested_type, + bool first = false, + const String & after_column = {}) { const auto before = readSchemaState(metadata); MetadataGenerator gen(metadata); try { - gen.generateModifyColumnMetadata(column, requested_type, getContext().context); + gen.generateModifyColumnMetadata(column, requested_type, getContext().context, first, after_column); FAIL() << "MODIFY COLUMN " << column << " " << requested_type->getName() << " cannot be recorded in Iceberg and must be rejected"; } @@ -410,6 +435,113 @@ void expectModifyRejected( } +TEST(IcebergMetadataGenerator, AddColumnFirstPlacesFieldAtIndexZero) +{ + auto metadata = makeMetadataWithGap(); + MetadataGenerator gen(metadata); + + gen.generateAddColumnMetadata("z", makeNullable(std::make_shared()), /* first */ true); + + auto current_schema_id = metadata->getValue(f_current_schema_id); + auto schemas = metadata->getArray(f_schemas); + Poco::JSON::Object::Ptr current_schema; + for (UInt32 i = 0; i < schemas->size(); ++i) + { + if (schemas->getObject(i)->getValue(f_schema_id) == current_schema_id) + { + current_schema = schemas->getObject(i); + break; + } + } + ASSERT_NE(current_schema.get(), nullptr); + + auto fields = current_schema->getArray(f_fields); + ASSERT_GE(fields->size(), 1u); + EXPECT_EQ(fields->getObject(0)->getValue(f_name), "z"); + EXPECT_EQ(fields->getObject(1)->getValue(f_name), "x"); + EXPECT_EQ(fields->getObject(2)->getValue(f_name), "y"); +} + + +TEST(IcebergMetadataGenerator, AddColumnAfterPlacesFieldAfterNamedColumn) +{ + auto metadata = makeMetadataWithGap(); + MetadataGenerator gen(metadata); + + gen.generateAddColumnMetadata("z", makeNullable(std::make_shared()), /* first */ false, /* after_column */ "x"); + + auto current_schema_id = metadata->getValue(f_current_schema_id); + auto schemas = metadata->getArray(f_schemas); + Poco::JSON::Object::Ptr current_schema; + for (UInt32 i = 0; i < schemas->size(); ++i) + { + if (schemas->getObject(i)->getValue(f_schema_id) == current_schema_id) + { + current_schema = schemas->getObject(i); + break; + } + } + ASSERT_NE(current_schema.get(), nullptr); + + auto fields = current_schema->getArray(f_fields); + ASSERT_EQ(fields->size(), 3u); + EXPECT_EQ(fields->getObject(0)->getValue(f_name), "x"); + EXPECT_EQ(fields->getObject(1)->getValue(f_name), "z"); + EXPECT_EQ(fields->getObject(2)->getValue(f_name), "y"); +} + + +TEST(IcebergMetadataGenerator, AddColumnAfterNonexistentColumnThrows) +{ + auto metadata = makeMetadataWithGap(); + MetadataGenerator gen(metadata); + + EXPECT_THROW( + gen.generateAddColumnMetadata("z", makeNullable(std::make_shared()), /* first */ false, /* after_column */ "nonexistent"), + DB::Exception); +} + + +TEST(IcebergMetadataGenerator, AddColumnAppliedRequiresFirstPosition) +{ + /// A concurrent add of the same column and type landed at the end: a retry of `ADD COLUMN z ... FIRST` + /// must not report success. + auto type = makeNullable(std::make_shared()); + { + auto metadata = makeMetadataWithGap(); + MetadataGenerator gen(metadata); + gen.generateAddColumnMetadata("z", type); + EXPECT_FALSE(gen.isAddColumnApplied("z", type, /* first */ true)); + EXPECT_TRUE(gen.isAddColumnApplied("z", type)); + } + { + auto metadata = makeMetadataWithGap(); + MetadataGenerator gen(metadata); + gen.generateAddColumnMetadata("z", type, /* first */ true); + EXPECT_TRUE(gen.isAddColumnApplied("z", type, /* first */ true)); + } +} + + +TEST(IcebergMetadataGenerator, AddColumnAppliedRequiresAfterPosition) +{ + auto type = makeNullable(std::make_shared()); + { + auto metadata = makeMetadataWithGap(); + MetadataGenerator gen(metadata); + gen.generateAddColumnMetadata("z", type); + EXPECT_FALSE(gen.isAddColumnApplied("z", type, /* first */ false, /* after_column */ "x")); + EXPECT_TRUE(gen.isAddColumnApplied("z", type, /* first */ false, /* after_column */ "y")); + } + { + auto metadata = makeMetadataWithGap(); + MetadataGenerator gen(metadata); + gen.generateAddColumnMetadata("z", type, /* first */ false, /* after_column */ "x"); + EXPECT_TRUE(gen.isAddColumnApplied("z", type, /* first */ false, /* after_column */ "x")); + } +} + + TEST(IcebergMetadataGenerator, ModifyColumnAppliedRecognisesTypeAlreadyInSchema) { /// The schema already says `long`, so a MODIFY to Int64 has taken effect. @@ -491,4 +623,317 @@ TEST(IcebergMetadataGenerator, ModifyColumnWideningRecordsTheNewTypeInANewSchema EXPECT_EQ(stored_type.extract(), "long"); } + +TEST(IcebergMetadataGenerator, ModifyColumnFirstMovesFieldToIndexZero) +{ + auto metadata = makeMetadataWithGap(); + const auto before = readSchemaState(metadata); + MetadataGenerator gen(metadata); + + EXPECT_TRUE(gen.generateModifyColumnMetadata( + "y", makeNullable(std::make_shared()), getContext().context, + /* first */ true)); + + const auto after = readSchemaState(metadata); + EXPECT_EQ(after.schema_count, before.schema_count + 1); + + auto names = getCurrentFieldNames(metadata); + ASSERT_EQ(names.size(), 2u); + EXPECT_EQ(names[0], "y"); + EXPECT_EQ(names[1], "x"); +} + + +TEST(IcebergMetadataGenerator, ModifyColumnAfterMovesFieldAfterNamedColumn) +{ + /// Build a 3-column schema: a, b, c. + auto metadata = Poco::JSON::Object::Ptr(new Poco::JSON::Object); + metadata->set(f_format_version, 2); + metadata->set(f_current_schema_id, 0); + metadata->set(f_last_column_id, 3); + + auto schema = Poco::JSON::Object::Ptr(new Poco::JSON::Object); + schema->set(f_schema_id, 0); + schema->set(f_type, "struct"); + auto fields = Poco::JSON::Array::Ptr(new Poco::JSON::Array); + for (Int32 i = 0; i < 3; ++i) + { + auto field = Poco::JSON::Object::Ptr(new Poco::JSON::Object); + field->set(f_id, i + 1); + field->set(f_name, String(1, static_cast('a' + i))); + field->set(f_required, false); + field->set(f_type, "long"); + fields->add(field); + } + schema->set(f_fields, fields); + auto schemas = Poco::JSON::Array::Ptr(new Poco::JSON::Array); + schemas->add(schema); + metadata->set(f_schemas, schemas); + + MetadataGenerator gen(metadata); + + EXPECT_TRUE(gen.generateModifyColumnMetadata( + "a", makeNullable(std::make_shared()), getContext().context, + /* first */ false, /* after_column */ "b")); + + auto names = getCurrentFieldNames(metadata); + ASSERT_EQ(names.size(), 3u); + EXPECT_EQ(names[0], "b"); + EXPECT_EQ(names[1], "a"); + EXPECT_EQ(names[2], "c"); +} + + +TEST(IcebergMetadataGenerator, ModifyColumnTypeAndFirstDoesBoth) +{ + auto metadata = makeMetadataWithGap(); + MetadataGenerator gen(metadata); + + EXPECT_TRUE(gen.generateModifyColumnMetadata( + "x", std::make_shared(), getContext().context, + /* first */ true)); + + auto names = getCurrentFieldNames(metadata); + ASSERT_EQ(names.size(), 2u); + EXPECT_EQ(names[0], "x"); + EXPECT_EQ(names[1], "y"); + + auto stored_type = findCurrentFieldType(metadata, "x"); + ASSERT_TRUE(stored_type.isString()); + EXPECT_EQ(stored_type.extract(), "long"); +} + + +TEST(IcebergMetadataGenerator, ModifyColumnAfterNonexistentColumnThrows) +{ + auto metadata = makeMetadataWithGap(); + MetadataGenerator gen(metadata); + + EXPECT_THROW( + gen.generateModifyColumnMetadata( + "x", std::make_shared(), getContext().context, + /* first */ false, /* after_column */ "nonexistent"), + DB::Exception); +} + + +TEST(IcebergMetadataGenerator, ModifyColumnAppliedPositionOnlyFirst) +{ + /// `MODIFY COLUMN c FIRST` carries no type; the predicate must check only the position. + auto metadata = makeMetadataWithGap(); + MetadataGenerator gen(metadata); + EXPECT_TRUE(gen.isModifyColumnApplied("x", nullptr, /* first */ true)); + EXPECT_FALSE(gen.isModifyColumnApplied("y", nullptr, /* first */ true)); + EXPECT_FALSE(gen.isModifyColumnApplied("absent", nullptr, /* first */ true)); +} + +TEST(IcebergMetadataGenerator, ModifyColumnAppliedPositionOnlyAfter) +{ + auto metadata = makeMetadataWithGap(); + MetadataGenerator gen(metadata); + EXPECT_TRUE(gen.isModifyColumnApplied("y", nullptr, /* first */ false, /* after_column */ "x")); + EXPECT_FALSE(gen.isModifyColumnApplied("x", nullptr, /* first */ false, /* after_column */ "y")); +} + +TEST(IcebergMetadataGenerator, ModifyColumnAppliedRequiresPositionWithType) +{ + /// The type is already in place, but a retry must not report success until the column moved. + auto metadata = makeMetadataWithGap(); + MetadataGenerator gen(metadata); + auto nullable_string = makeNullable(std::make_shared()); + EXPECT_FALSE(gen.isModifyColumnApplied("y", nullable_string, /* first */ true)); + EXPECT_TRUE(gen.isModifyColumnApplied("y", nullable_string, /* first */ false, /* after_column */ "x")); + EXPECT_TRUE(gen.isModifyColumnApplied("x", std::make_shared(), /* first */ true)); + EXPECT_FALSE(gen.isModifyColumnApplied("x", std::make_shared(), /* first */ true)); +} + +TEST(IcebergMetadataGenerator, ModifyColumnAppliedSelfAfterIgnoresPosition) +{ + auto metadata = makeMetadataWithGap(); + MetadataGenerator gen(metadata); + EXPECT_TRUE(gen.isModifyColumnApplied("x", std::make_shared(), /* first */ false, /* after_column */ "x")); +} + +TEST(IcebergMetadataGenerator, ModifyColumnRejectsIndistinguishablePrimitiveTypeWithFirst) +{ + /// Int32 and UInt32 are both Iceberg `int`; positioning must not make the change recordable. + auto metadata = makeMetadataWithGap(); + expectModifyRejected(metadata, "x", std::make_shared(), /* first */ true); +} + +TEST(IcebergMetadataGenerator, ModifyColumnRejectsIndistinguishablePrimitiveTypeWithAfter) +{ + auto metadata = makeMetadataWithGap(); + expectModifyRejected(metadata, "x", std::make_shared(), /* first */ false, /* after_column */ "y"); +} + +TEST(IcebergMetadataGenerator, ModifyColumnSameTypeWithFirstOnlyRepositions) +{ + auto metadata = makeMetadataWithGap(); + MetadataGenerator gen(metadata); + + EXPECT_TRUE(gen.generateModifyColumnMetadata( + "y", makeNullable(std::make_shared()), getContext().context, /* first */ true)); + + auto names = getCurrentFieldNames(metadata); + ASSERT_EQ(names.size(), 2u); + EXPECT_EQ(names[0], "y"); + EXPECT_EQ(names[1], "x"); + + auto stored_type = findCurrentFieldType(metadata, "y"); + ASSERT_TRUE(stored_type.isString()); + EXPECT_EQ(stored_type.extract(), "string"); +} + +TEST(IcebergMetadataGenerator, ModifyColumnAfterItselfWithoutTypeIsNoop) +{ + auto metadata = makeMetadataWithGap(); + const auto before = readSchemaState(metadata); + + EXPECT_FALSE(MetadataGenerator(metadata).generateModifyColumnMetadata( + "y", nullptr, getContext().context, /* first */ false, /* after_column */ "y")); + expectSchemaUnchanged(metadata, before); +} + +TEST(IcebergMetadataGenerator, ModifyColumnAfterItselfAppliesTypeAndKeepsPosition) +{ + auto metadata = makeMetadataWithGap(); + + EXPECT_TRUE(MetadataGenerator(metadata).generateModifyColumnMetadata( + "x", std::make_shared(), getContext().context, /* first */ false, /* after_column */ "x")); + + auto names = getCurrentFieldNames(metadata); + ASSERT_EQ(names.size(), 2u); + EXPECT_EQ(names[0], "x"); + EXPECT_EQ(names[1], "y"); + + auto stored_type = findCurrentFieldType(metadata, "x"); + ASSERT_TRUE(stored_type.isString()); + EXPECT_EQ(stored_type.extract(), "long"); +} + +namespace +{ + +/// Metadata whose current schema holds `x int` followed by `t` of the given complex Iceberg type. +Poco::JSON::Object::Ptr makeMetadataWithIntAndComplexField(const Poco::JSON::Object::Ptr & complex_type, Int32 last_column_id) +{ + auto metadata = makeMetadataWithField("x", "int", /* required */ true, last_column_id); + auto field = Poco::JSON::Object::Ptr(new Poco::JSON::Object); + field->set(f_id, 2); + field->set(f_name, "t"); + field->set(f_required, true); + field->set(f_type, complex_type); + metadata->getArray(f_schemas)->getObject(0)->getArray(f_fields)->add(field); + return metadata; +} + +Poco::JSON::Object::Ptr makeListType(const String & element_type, Int32 element_id) +{ + auto type = Poco::JSON::Object::Ptr(new Poco::JSON::Object); + type->set(f_type, "list"); + type->set(f_element_id, element_id); + type->set(f_element, element_type); + type->set(f_element_required, true); + return type; +} + +DataTypePtr makeTupleOf(const DataTypePtr & element_type) +{ + return std::make_shared(DataTypes{element_type}, Names{"a"}); +} + +} + +TEST(IcebergMetadataGenerator, ModifyColumnRejectsIndistinguishableComplexTypeWithFirst) +{ + /// Tuple(a Int32) and Tuple(a UInt32) are one Iceberg struct; positioning must not make the change recordable. + auto metadata = makeMetadataWithIntAndComplexField(makeStructType("a", "int", 3), /* last_column_id */ 3); + expectModifyRejected(metadata, "t", makeTupleOf(std::make_shared()), /* first */ true); +} + +TEST(IcebergMetadataGenerator, ModifyColumnRejectsIndistinguishableComplexTypeWithAfter) +{ + auto metadata = makeMetadataWithIntAndComplexField(makeStructType("a", "int", 3), /* last_column_id */ 3); + expectModifyRejected(metadata, "t", makeTupleOf(std::make_shared()), /* first */ false, /* after_column */ "x"); +} + +TEST(IcebergMetadataGenerator, ModifyColumnRejectsIndistinguishableArrayTypeWithFirst) +{ + auto metadata = makeMetadataWithIntAndComplexField(makeListType("int", 3), /* last_column_id */ 3); + expectModifyRejected( + metadata, "t", std::make_shared(std::make_shared()), /* first */ true); +} + +TEST(IcebergMetadataGenerator, ModifyColumnSameComplexTypeWithFirstOnlyRepositions) +{ + auto metadata = makeMetadataWithIntAndComplexField(makeStructType("a", "int", 3), /* last_column_id */ 3); + + EXPECT_TRUE(MetadataGenerator(metadata).generateModifyColumnMetadata( + "t", makeTupleOf(std::make_shared()), getContext().context, /* first */ true)); + + auto names = getCurrentFieldNames(metadata); + ASSERT_EQ(names.size(), 2u); + EXPECT_EQ(names[0], "t"); + EXPECT_EQ(names[1], "x"); + + auto stored_type = findCurrentFieldType(metadata, "t"); + ASSERT_EQ(stored_type.type(), typeid(Poco::JSON::Object::Ptr)); + auto stored_fields = stored_type.extract()->getArray(f_fields); + ASSERT_EQ(stored_fields->size(), 1u); + EXPECT_EQ(stored_fields->getObject(0)->getValue(f_type), "int"); +} + +TEST(IcebergMetadataGenerator, ModifyColumnToComplexTypeAlreadyInSchemaAddsNoSchema) +{ + auto metadata = makeMetadataWithIntAndComplexField(makeStructType("a", "int", 3), /* last_column_id */ 3); + const auto before = readSchemaState(metadata); + + EXPECT_FALSE(MetadataGenerator(metadata).generateModifyColumnMetadata( + "t", makeTupleOf(std::make_shared()), getContext().context)); + expectSchemaUnchanged(metadata, before); +} + +TEST(IcebergMetadataGenerator, ModifyColumnSameArrayTypeWithFirstOnlyRepositions) +{ + /// A list as other writers store it: unlike `getIcebergType`, no `required` on the list itself. + auto metadata = makeMetadataWithIntAndComplexField(makeListType("int", 3), /* last_column_id */ 3); + + EXPECT_TRUE(MetadataGenerator(metadata).generateModifyColumnMetadata( + "t", std::make_shared(std::make_shared()), getContext().context, /* first */ true)); + + auto names = getCurrentFieldNames(metadata); + ASSERT_EQ(names.size(), 2u); + EXPECT_EQ(names[0], "t"); + EXPECT_EQ(names[1], "x"); + + auto stored_type = findCurrentFieldType(metadata, "t"); + ASSERT_EQ(stored_type.type(), typeid(Poco::JSON::Object::Ptr)); + auto stored_list = stored_type.extract(); + EXPECT_EQ(stored_list->getValue(f_element), "int"); + EXPECT_TRUE(stored_list->getValue(f_element_required)); + EXPECT_FALSE(stored_list->has(f_required)); +} + +TEST(IcebergMetadataGenerator, ModifyColumnSameComplexTypeWithDocOnlyRepositions) +{ + auto struct_type = makeStructType("a", "int", 3); + struct_type->getArray(f_fields)->getObject(0)->set("doc", "documented by another writer"); + auto metadata = makeMetadataWithIntAndComplexField(struct_type, /* last_column_id */ 3); + + EXPECT_TRUE(MetadataGenerator(metadata).generateModifyColumnMetadata( + "t", makeTupleOf(std::make_shared()), getContext().context, /* first */ true)); + + auto names = getCurrentFieldNames(metadata); + ASSERT_EQ(names.size(), 2u); + EXPECT_EQ(names[0], "t"); + EXPECT_EQ(names[1], "x"); + + auto stored_type = findCurrentFieldType(metadata, "t"); + ASSERT_EQ(stored_type.type(), typeid(Poco::JSON::Object::Ptr)); + auto stored_fields = stored_type.extract()->getArray(f_fields); + ASSERT_EQ(stored_fields->size(), 1u); + EXPECT_EQ(stored_fields->getObject(0)->getValue(f_type), "int"); +} + #endif diff --git a/tests/integration/test_storage_iceberg_no_spark/test_writes_modify_column_position.py b/tests/integration/test_storage_iceberg_no_spark/test_writes_modify_column_position.py new file mode 100644 index 000000000000..3e4b35afabd4 --- /dev/null +++ b/tests/integration/test_storage_iceberg_no_spark/test_writes_modify_column_position.py @@ -0,0 +1,195 @@ +import pytest + +from helpers.iceberg_utils import ( + create_iceberg_table, + get_uuid_str, +) + +INSERT_SETTINGS = {"allow_insert_into_iceberg": 1} + + +@pytest.mark.parametrize("format_version", [1, 2]) +@pytest.mark.parametrize("storage_type", ["local", "s3"]) +def test_modify_column_first(started_cluster_iceberg_no_spark, format_version, storage_type): + """MODIFY COLUMN ... FIRST moves the column to the front of the schema.""" + instance = started_cluster_iceberg_no_spark.instances["node1"] + TABLE_NAME = "test_modify_first_" + storage_type + "_" + get_uuid_str() + + create_iceberg_table( + storage_type, + instance, + TABLE_NAME, + started_cluster_iceberg_no_spark, + "(a Int32, b Nullable(String), c Nullable(Int64))", + format_version, + ) + + instance.query(f"INSERT INTO {TABLE_NAME} VALUES (1, 'x', 10);", settings=INSERT_SETTINGS) + + instance.query( + f"ALTER TABLE {TABLE_NAME} MODIFY COLUMN c Nullable(Int64) FIRST;", + settings=INSERT_SETTINGS, + ) + + col_names = instance.query( + f"SELECT name FROM system.columns WHERE database = currentDatabase() AND table = '{TABLE_NAME}' ORDER BY position" + ).strip() + assert col_names == "c\na\nb", f"Expected c,a,b but got: {col_names}" + + result = instance.query(f"SELECT * FROM {TABLE_NAME}") + assert result.strip() == "10\t1\tx" + + +@pytest.mark.parametrize("format_version", [1, 2]) +@pytest.mark.parametrize("storage_type", ["local", "s3"]) +def test_modify_column_after(started_cluster_iceberg_no_spark, format_version, storage_type): + """MODIFY COLUMN ... AFTER moves the column after the named column.""" + instance = started_cluster_iceberg_no_spark.instances["node1"] + TABLE_NAME = "test_modify_after_" + storage_type + "_" + get_uuid_str() + + create_iceberg_table( + storage_type, + instance, + TABLE_NAME, + started_cluster_iceberg_no_spark, + "(a Int32, b Nullable(String), c Nullable(Int64))", + format_version, + ) + + instance.query(f"INSERT INTO {TABLE_NAME} VALUES (1, 'x', 10);", settings=INSERT_SETTINGS) + + instance.query( + f"ALTER TABLE {TABLE_NAME} MODIFY COLUMN a Int32 AFTER b;", + settings=INSERT_SETTINGS, + ) + + col_names = instance.query( + f"SELECT name FROM system.columns WHERE database = currentDatabase() AND table = '{TABLE_NAME}' ORDER BY position" + ).strip() + assert col_names == "b\na\nc", f"Expected b,a,c but got: {col_names}" + + +@pytest.mark.parametrize("format_version", [1, 2]) +@pytest.mark.parametrize("storage_type", ["local", "s3"]) +def test_modify_column_type_and_first(started_cluster_iceberg_no_spark, format_version, storage_type): + """MODIFY COLUMN with both a type change and FIRST does both.""" + instance = started_cluster_iceberg_no_spark.instances["node1"] + TABLE_NAME = "test_modify_type_first_" + storage_type + "_" + get_uuid_str() + + create_iceberg_table( + storage_type, + instance, + TABLE_NAME, + started_cluster_iceberg_no_spark, + "(a Int32, b Nullable(String), c Nullable(Int32))", + format_version, + ) + + instance.query(f"INSERT INTO {TABLE_NAME} VALUES (1, 'x', 10);", settings=INSERT_SETTINGS) + + instance.query( + f"ALTER TABLE {TABLE_NAME} MODIFY COLUMN c Nullable(Int64) FIRST;", + settings=INSERT_SETTINGS, + ) + + col_names = instance.query( + f"SELECT name FROM system.columns WHERE database = currentDatabase() AND table = '{TABLE_NAME}' ORDER BY position" + ).strip() + assert col_names == "c\na\nb", f"Expected c,a,b but got: {col_names}" + + instance.query(f"INSERT INTO {TABLE_NAME} VALUES (3000000000, 2, 'y');", settings=INSERT_SETTINGS) + result = instance.query(f"SELECT c FROM {TABLE_NAME} ORDER BY c") + assert "3000000000" in result + + +@pytest.mark.parametrize("format_version", [1, 2]) +@pytest.mark.parametrize("storage_type", ["local", "s3"]) +def test_modify_column_unrecordable_type_with_position_rejected(started_cluster_iceberg_no_spark, format_version, storage_type): + """Positioning must not bypass the rejection of a type change Iceberg cannot record.""" + instance = started_cluster_iceberg_no_spark.instances["node1"] + TABLE_NAME = "test_modify_unrecordable_pos_" + storage_type + "_" + get_uuid_str() + + create_iceberg_table( + storage_type, + instance, + TABLE_NAME, + started_cluster_iceberg_no_spark, + "(a Int32, b Nullable(String), c Int32)", + format_version, + ) + + for clause in ["FIRST", "AFTER a"]: + error = instance.query_and_get_error( + f"ALTER TABLE {TABLE_NAME} MODIFY COLUMN c UInt32 {clause};", + settings=INSERT_SETTINGS, + ) + assert "BAD_ARGUMENTS" in error, error + + columns = instance.query( + f"SELECT name, type FROM system.columns WHERE database = currentDatabase() AND table = '{TABLE_NAME}' ORDER BY position" + ).strip() + assert columns == "a\tInt32\nb\tNullable(String)\nc\tInt32", columns + + +@pytest.mark.parametrize("format_version", [1, 2]) +@pytest.mark.parametrize("storage_type", ["local", "s3"]) +def test_modify_column_after_itself(started_cluster_iceberg_no_spark, format_version, storage_type): + """MODIFY COLUMN ... AFTER the column itself keeps its position.""" + instance = started_cluster_iceberg_no_spark.instances["node1"] + TABLE_NAME = "test_modify_after_self_" + storage_type + "_" + get_uuid_str() + + create_iceberg_table( + storage_type, + instance, + TABLE_NAME, + started_cluster_iceberg_no_spark, + "(a Int32, b Nullable(Int32), c Nullable(String))", + format_version, + ) + + instance.query(f"INSERT INTO {TABLE_NAME} VALUES (1, 10, 'x');", settings=INSERT_SETTINGS) + + instance.query(f"ALTER TABLE {TABLE_NAME} MODIFY COLUMN b Nullable(Int32) AFTER b;", settings=INSERT_SETTINGS) + instance.query(f"ALTER TABLE {TABLE_NAME} MODIFY COLUMN b Nullable(Int64) AFTER b;", settings=INSERT_SETTINGS) + + columns = instance.query( + f"SELECT name, type FROM system.columns WHERE database = currentDatabase() AND table = '{TABLE_NAME}' ORDER BY position" + ).strip() + assert columns == "a\tInt32\nb\tNullable(Int64)\nc\tNullable(String)", columns + + assert instance.query(f"SELECT * FROM {TABLE_NAME}").strip() == "1\t10\tx" + + +@pytest.mark.parametrize("format_version", [1, 2]) +@pytest.mark.parametrize("storage_type", ["local", "s3"]) +def test_modify_column_complex_type_with_position(started_cluster_iceberg_no_spark, format_version, storage_type): + """An unrecordable complex type change is rejected with FIRST/AFTER, restating the same type only moves the column.""" + instance = started_cluster_iceberg_no_spark.instances["node1"] + TABLE_NAME = "test_modify_complex_pos_" + storage_type + "_" + get_uuid_str() + + create_iceberg_table( + storage_type, + instance, + TABLE_NAME, + started_cluster_iceberg_no_spark, + "(x Int32, t Tuple(a Int32))", + format_version, + ) + + instance.query(f"INSERT INTO {TABLE_NAME} VALUES (1, tuple(10));", settings=INSERT_SETTINGS) + + for clause in ["FIRST", "AFTER x"]: + error = instance.query_and_get_error( + f"ALTER TABLE {TABLE_NAME} MODIFY COLUMN t Tuple(a UInt32) {clause};", + settings=INSERT_SETTINGS, + ) + assert "BAD_ARGUMENTS" in error, error + + instance.query(f"ALTER TABLE {TABLE_NAME} MODIFY COLUMN t Tuple(a Int32) FIRST;", settings=INSERT_SETTINGS) + + columns = instance.query( + f"SELECT name, type FROM system.columns WHERE database = currentDatabase() AND table = '{TABLE_NAME}' ORDER BY position" + ).strip() + assert columns == "t\tTuple(a Int32)\nx\tInt32", columns + + assert instance.query(f"SELECT * FROM {TABLE_NAME}").strip() == "(10)\t1" diff --git a/tests/queries/0_stateless/05054_iceberg_alter_modify_column_position_rejects_clauses.reference b/tests/queries/0_stateless/05054_iceberg_alter_modify_column_position_rejects_clauses.reference new file mode 100644 index 000000000000..aeff25700052 --- /dev/null +++ b/tests/queries/0_stateless/05054_iceberg_alter_modify_column_position_rejects_clauses.reference @@ -0,0 +1,10 @@ +MODIFY COLUMN c1 COMMENT 'x' FIRST +NOT_IMPLEMENTED +MODIFY COLUMN c1 Int DEFAULT 1 AFTER c0 +NOT_IMPLEMENTED +MODIFY COLUMN c1 Int CODEC(ZSTD) FIRST +NOT_IMPLEMENTED +MODIFY COLUMN c1 Nullable(Int) COMMENT 'x' +NOT_IMPLEMENTED +c0 Int32 +c1 Int32 diff --git a/tests/queries/0_stateless/05054_iceberg_alter_modify_column_position_rejects_clauses.sh b/tests/queries/0_stateless/05054_iceberg_alter_modify_column_position_rejects_clauses.sh new file mode 100755 index 000000000000..410c4fd3dd35 --- /dev/null +++ b/tests/queries/0_stateless/05054_iceberg_alter_modify_column_position_rejects_clauses.sh @@ -0,0 +1,35 @@ +#!/usr/bin/env bash +# Tags: no-fasttest + +# The Iceberg schema records only the type and the order of columns, so `MODIFY COLUMN` clauses +# other than the type and `FIRST` / `AFTER` must be rejected rather than silently dropped. + +CUR_DIR=$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd) +# shellcheck source=../shell_config.sh +. "$CUR_DIR"/../shell_config.sh + +TABLE="t_${CLICKHOUSE_DATABASE}_${RANDOM}" +TABLE_PATH="${USER_FILES_PATH}/${TABLE}/" + +${CLICKHOUSE_CLIENT} --query "DROP TABLE IF EXISTS ${TABLE}" +${CLICKHOUSE_CLIENT} --query " + CREATE TABLE ${TABLE} (c0 Int, c1 Int) + ENGINE = IcebergLocal('${TABLE_PATH}') +" +# To have at least one real snapshot. Otherwise alter can be noop. +${CLICKHOUSE_CLIENT} --allow_insert_into_iceberg=1 --query "INSERT INTO ${TABLE} VALUES (1, 2)" + +for alter in \ + "MODIFY COLUMN c1 COMMENT 'x' FIRST" \ + "MODIFY COLUMN c1 Int DEFAULT 1 AFTER c0" \ + "MODIFY COLUMN c1 Int CODEC(ZSTD) FIRST" \ + "MODIFY COLUMN c1 Nullable(Int) COMMENT 'x'" +do + echo "${alter}" + ${CLICKHOUSE_CLIENT} --query "ALTER TABLE ${TABLE} ${alter}" 2>&1 | grep -o -m1 "NOT_IMPLEMENTED" +done + +${CLICKHOUSE_CLIENT} --query "DESCRIBE TABLE ${TABLE}" | cut -f1,2 + +${CLICKHOUSE_CLIENT} --query "DROP TABLE IF EXISTS ${TABLE}" +rm -rf "${TABLE_PATH}" diff --git a/tests/queries/0_stateless/05055_iceberg_alter_add_column_position_rejects_clauses.reference b/tests/queries/0_stateless/05055_iceberg_alter_add_column_position_rejects_clauses.reference new file mode 100644 index 000000000000..57d8f9910454 --- /dev/null +++ b/tests/queries/0_stateless/05055_iceberg_alter_add_column_position_rejects_clauses.reference @@ -0,0 +1,7 @@ +ADD COLUMN z Nullable(Int32) DEFAULT 7 FIRST +NOT_IMPLEMENTED +ADD COLUMN z Nullable(Int32) ALIAS c0 AFTER c0 +NOT_IMPLEMENTED +ADD COLUMN z Nullable(Int32) COMMENT 'x' FIRST +NOT_IMPLEMENTED +c0 Int32 diff --git a/tests/queries/0_stateless/05055_iceberg_alter_add_column_position_rejects_clauses.sh b/tests/queries/0_stateless/05055_iceberg_alter_add_column_position_rejects_clauses.sh new file mode 100755 index 000000000000..3aca967d0e0d --- /dev/null +++ b/tests/queries/0_stateless/05055_iceberg_alter_add_column_position_rejects_clauses.sh @@ -0,0 +1,34 @@ +#!/usr/bin/env bash +# Tags: no-fasttest + +# The Iceberg schema records only the type and the order of columns, so `ADD COLUMN` clauses +# other than the type and `FIRST` / `AFTER` must be rejected rather than silently dropped. + +CUR_DIR=$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd) +# shellcheck source=../shell_config.sh +. "$CUR_DIR"/../shell_config.sh + +TABLE="t_${CLICKHOUSE_DATABASE}_${RANDOM}" +TABLE_PATH="${USER_FILES_PATH}/${TABLE}/" + +${CLICKHOUSE_CLIENT} --query "DROP TABLE IF EXISTS ${TABLE}" +${CLICKHOUSE_CLIENT} --query " + CREATE TABLE ${TABLE} (c0 Int) + ENGINE = IcebergLocal('${TABLE_PATH}') +" +# To have at least one real snapshot. Otherwise alter can be noop. +${CLICKHOUSE_CLIENT} --allow_insert_into_iceberg=1 --query "INSERT INTO ${TABLE} VALUES (1)" + +for alter in \ + "ADD COLUMN z Nullable(Int32) DEFAULT 7 FIRST" \ + "ADD COLUMN z Nullable(Int32) ALIAS c0 AFTER c0" \ + "ADD COLUMN z Nullable(Int32) COMMENT 'x' FIRST" +do + echo "${alter}" + ${CLICKHOUSE_CLIENT} --allow_insert_into_iceberg=1 --query "ALTER TABLE ${TABLE} ${alter}" 2>&1 | grep -o -m1 "NOT_IMPLEMENTED" +done + +${CLICKHOUSE_CLIENT} --query "DESCRIBE TABLE ${TABLE}" | cut -f1,2 + +${CLICKHOUSE_CLIENT} --query "DROP TABLE IF EXISTS ${TABLE}" +rm -rf "${TABLE_PATH}"