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..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 @@ -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_NE(current_schema, 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, 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, ModifyColumnAppliedRecognisesTypeAlreadyInSchema) { /// The schema already says `long`, so a MODIFY to Int64 has taken effect.