Skip to content
Open
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
Original file line number Diff line number Diff line change
Expand Up @@ -828,10 +828,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);
}
}

Expand Down
187 changes: 138 additions & 49 deletions src/Storages/ObjectStorage/DataLakes/Iceberg/MetadataGenerator.cpp
Original file line number Diff line number Diff line change
@@ -1,5 +1,4 @@
#include <Core/Settings.h>
#include <DataTypes/DataTypeNullable.h>
#include <IO/ReadBufferFromString.h>
#include <IO/ReadHelpers.h>
#include <Interpreters/Context.h>
Expand Down Expand Up @@ -357,23 +356,34 @@ 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<Int32>(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<String>(Iceberg::f_name) != column_name)
continue;
return field->getValue<bool>(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<Int32>(Iceberg::f_last_column_id);
auto expected_type = Iceberg::getIcebergType(type, unused_field_id);
if (field->getValue<bool>(Iceberg::f_required) != expected_type.second
|| !icebergTypesEqualIgnoringIds(field->get(Iceberg::f_type), expected_type.first))
return false;
}

if (first)
return i == 0;
if (!after_column.empty() && after_column != column_name)
return i > 0 && fields->getObject(i - 1)->getValue<String>(Iceberg::f_name) == after_column;
return true;
}
return false;
}
Expand Down Expand Up @@ -722,7 +732,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 @@ -748,81 +758,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<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);
}

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<Int32>(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<String>(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<String>(Iceberg::f_name) != column_name)
continue;

if (current_field->getValue<bool>(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<String>(),
context,
context->getSettingsRef()[Setting::allow_experimental_geo_types_in_iceberg]);
if (!current_field->getValue<bool>(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<String>());
}
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 (!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<bool>(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<bool>(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<UInt32>(schema_fields->size());
for (UInt32 i = 0; i < schema_fields->size(); ++i)
{
if (schema_fields->getObject(i)->getValue<String>(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<String>(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)
Expand Down
10 changes: 6 additions & 4 deletions src/Storages/ObjectStorage/DataLakes/Iceberg/MetadataGenerator.h
Original file line number Diff line number Diff line change
Expand Up @@ -41,18 +41,19 @@ class MetadataGenerator
std::optional<Int64> 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
Expand All @@ -62,7 +63,8 @@ class MetadataGenerator
bool isAddColumnApplied(const String & column_name, DataTypePtr type) 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;
Expand Down
6 changes: 3 additions & 3 deletions src/Storages/ObjectStorage/DataLakes/Iceberg/Mutations.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -115,7 +115,7 @@ static bool alterAlreadyApplied(const MetadataGenerator & generator, const Alter
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;
}
Expand Down 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 All @@ -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;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<bool>(f_required));
}

/**
* Iceberg allows only three types of primitive type conversion:
* int -> long
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<String, Int64> traverseSchema(Poco::JSON::Array::Ptr schema);

void registerSnapshotWithSchemaId(Int64 snapshot_id, Int32 schema_id);
Expand Down
Loading
Loading