From 073b19707f898ae9c29561ea6e7b93716a4c2317 Mon Sep 17 00:00:00 2001 From: Kanthi Subramanian Date: Fri, 28 Aug 2026 18:11:21 +0200 Subject: [PATCH 1/4] Added support for adlter add column with first/last --- .../DataLakes/Iceberg/MetadataGenerator.cpp | 34 +++++++++- .../DataLakes/Iceberg/MetadataGenerator.h | 2 +- .../DataLakes/Iceberg/Mutations.cpp | 2 +- .../gtest_iceberg_metadata_generator.cpp | 67 +++++++++++++++++++ 4 files changed, 101 insertions(+), 4 deletions(-) diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/MetadataGenerator.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/MetadataGenerator.cpp index 55ee1c99baf3..e34bf5af873d 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/MetadataGenerator.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/MetadataGenerator.cpp @@ -564,7 +564,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"); @@ -590,7 +590,37 @@ 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); diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/MetadataGenerator.h b/src/Storages/ObjectStorage/DataLakes/Iceberg/MetadataGenerator.h index a20c1cfdc827..7e880616baf4 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/MetadataGenerator.h +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/MetadataGenerator.h @@ -43,7 +43,7 @@ class MetadataGenerator std::optional user_defined_timestamp = std::nullopt, bool is_truncate = false); - void generateAddColumnMetadata(const String & column_name, DataTypePtr type); + void generateAddColumnMetadata(const String & column_name, DataTypePtr type, bool first = false, const String & after_column = {}); void generateDropColumnMetadata(const String & column_name); /// Returns false when the column already has the requested type (no metadata change). /// `context` supplies the settings used to map the stored Iceberg type back to a ClickHouse diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/Mutations.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/Mutations.cpp index 2933f50f9296..7098b01bacf6 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/Mutations.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/Mutations.cpp @@ -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) 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..5f141e5b066b 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 @@ -410,6 +410,73 @@ 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_TRUE(current_schema); + + 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_TRUE(current_schema); + + 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, ModifyColumnAppliedRecognisesTypeAlreadyInSchema) { /// The schema already says `long`, so a MODIFY to Int64 has taken effect. From bc1676d709a0645f6c3bce57f245ab0ef2873bf7 Mon Sep 17 00:00:00 2001 From: Kanthi Subramanian Date: Sat, 29 Aug 2026 19:57:13 +0200 Subject: [PATCH 2/4] Fix unit tests --- .../Iceberg/tests/gtest_iceberg_metadata_generator.cpp | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) 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 5f141e5b066b..cb7d85d4fb3b 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 @@ -428,7 +428,7 @@ TEST(IcebergMetadataGenerator, AddColumnFirstPlacesFieldAtIndexZero) break; } } - ASSERT_TRUE(current_schema); + ASSERT_NE(current_schema, nullptr); auto fields = current_schema->getArray(f_fields); ASSERT_GE(fields->size(), 1u); @@ -456,7 +456,7 @@ TEST(IcebergMetadataGenerator, AddColumnAfterPlacesFieldAfterNamedColumn) break; } } - ASSERT_TRUE(current_schema); + ASSERT_NE(current_schema, nullptr); auto fields = current_schema->getArray(f_fields); ASSERT_EQ(fields->size(), 3u); From b49fcbfebca0cfaeca295b422699174683a6f5cd Mon Sep 17 00:00:00 2001 From: Kanthi Subramanian Date: Mon, 31 Aug 2026 02:10:18 +0200 Subject: [PATCH 3/4] Fixed compiler error --- .../Iceberg/tests/gtest_iceberg_metadata_generator.cpp | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) 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 cb7d85d4fb3b..a017a8b4cac7 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 @@ -428,7 +428,7 @@ TEST(IcebergMetadataGenerator, AddColumnFirstPlacesFieldAtIndexZero) break; } } - ASSERT_NE(current_schema, nullptr); + ASSERT_NE(current_schema.get(), nullptr); auto fields = current_schema->getArray(f_fields); ASSERT_GE(fields->size(), 1u); @@ -456,7 +456,7 @@ TEST(IcebergMetadataGenerator, AddColumnAfterPlacesFieldAfterNamedColumn) break; } } - ASSERT_NE(current_schema, nullptr); + ASSERT_NE(current_schema.get(), nullptr); auto fields = current_schema->getArray(f_fields); ASSERT_EQ(fields->size(), 3u); From fa600ba33231c621758fb67ba2b1039ac2529ed3 Mon Sep 17 00:00:00 2001 From: Kanthi Subramanian Date: Tue, 1 Sep 2026 03:32:41 +0200 Subject: [PATCH 4/4] Added support for modify column first/after iceberg --- .../DataLakes/Iceberg/IcebergMetadata.cpp | 5 +- .../DataLakes/Iceberg/MetadataGenerator.cpp | 147 +++++++++++++----- .../DataLakes/Iceberg/MetadataGenerator.h | 4 +- .../DataLakes/Iceberg/Mutations.cpp | 2 +- .../gtest_iceberg_metadata_generator.cpp | 113 ++++++++++++++ .../test_writes_modify_column_position.py | 102 ++++++++++++ 6 files changed, 326 insertions(+), 47 deletions(-) create mode 100644 tests/integration/test_storage_iceberg_no_spark/test_writes_modify_column_position.py diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.cpp index 24930b88462a..e4c6a9f001a5 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.cpp @@ -790,10 +790,11 @@ 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::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 e34bf5af873d..3cd07058506a 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/MetadataGenerator.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/MetadataGenerator.cpp @@ -626,75 +626,138 @@ void MetadataGenerator::generateAddColumnMetadata(const String & column_name, Da 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) + bool needs_reposition = first || !after_column.empty(); + 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()) + if (!needs_reposition) { - 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); + 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; - if (reconstructed_ch_type->equals(*type)) - return false; + 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(), + type->getName(), + existing_iceberg_type.extract()); + } 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(), - type->getName(), - existing_iceberg_type.extract()); + "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); } - - 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); } + 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 (!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 (!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()); + type_changed = true; + } + break; + } + } - const auto next_schema_id = getNextSchemaId(metadata_object); + 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); - current_schema = deepCopy(current_schema); - schema_fields = current_schema->getArray(Iceberg::f_fields); - current_field = schema_fields->getObject(i); + if (!type_changed && !needs_reposition) + return false; - current_field->set(Iceberg::f_type, new_type.first); - current_field->set(Iceberg::f_required, new_type.second); + 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); - 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 (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); + } + + 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 7e880616baf4..1120edd60d27 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/MetadataGenerator.h +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/MetadataGenerator.h @@ -45,10 +45,10 @@ class MetadataGenerator void generateAddColumnMetadata(const String & column_name, DataTypePtr type, bool first = false, const String & after_column = {}); 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). /// `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 diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/Mutations.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/Mutations.cpp index 7098b01bacf6..521676d25286 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/Mutations.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/Mutations.cpp @@ -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/tests/gtest_iceberg_metadata_generator.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/tests/gtest_iceberg_metadata_generator.cpp index a017a8b4cac7..ed960304c298 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 @@ -367,6 +367,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) { @@ -558,4 +578,97 @@ 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); +} + #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..0789c342f248 --- /dev/null +++ b/tests/integration/test_storage_iceberg_no_spark/test_writes_modify_column_position.py @@ -0,0 +1,102 @@ +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