Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
34 changes: 32 additions & 2 deletions src/Storages/ObjectStorage/DataLakes/Iceberg/MetadataGenerator.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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");
Expand All @@ -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<String>(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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,7 @@ class MetadataGenerator
std::optional<Int64> 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
Expand Down
2 changes: 1 addition & 1 deletion src/Storages/ObjectStorage/DataLakes/Iceberg/Mutations.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -410,6 +410,73 @@ void expectModifyRejected(

}

TEST(IcebergMetadataGenerator, AddColumnFirstPlacesFieldAtIndexZero)
{
auto metadata = makeMetadataWithGap();
MetadataGenerator gen(metadata);

gen.generateAddColumnMetadata("z", makeNullable(std::make_shared<DataTypeInt64>()), /* first */ true);

auto current_schema_id = metadata->getValue<Int32>(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<Int32>(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<String>(f_name), "z");
EXPECT_EQ(fields->getObject(1)->getValue<String>(f_name), "x");
EXPECT_EQ(fields->getObject(2)->getValue<String>(f_name), "y");
}


TEST(IcebergMetadataGenerator, AddColumnAfterPlacesFieldAfterNamedColumn)
{
auto metadata = makeMetadataWithGap();
MetadataGenerator gen(metadata);

gen.generateAddColumnMetadata("z", makeNullable(std::make_shared<DataTypeInt64>()), /* first */ false, /* after_column */ "x");

auto current_schema_id = metadata->getValue<Int32>(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<Int32>(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<String>(f_name), "x");
EXPECT_EQ(fields->getObject(1)->getValue<String>(f_name), "z");
EXPECT_EQ(fields->getObject(2)->getValue<String>(f_name), "y");
}


TEST(IcebergMetadataGenerator, AddColumnAfterNonexistentColumnThrows)
{
auto metadata = makeMetadataWithGap();
MetadataGenerator gen(metadata);

EXPECT_THROW(
gen.generateAddColumnMetadata("z", makeNullable(std::make_shared<DataTypeInt64>()), /* first */ false, /* after_column */ "nonexistent"),
DB::Exception);
}


TEST(IcebergMetadataGenerator, ModifyColumnAppliedRecognisesTypeAlreadyInSchema)
{
/// The schema already says `long`, so a MODIFY to Int64 has taken effect.
Expand Down
Loading