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 @@ -49,6 +49,14 @@ ChunkPartitioner::ChunkPartitioner(

FunctionOverloadResolverPtr transform;

/// Iceberg V3: multi-argument transforms use `source-ids` instead of `source-id`.
/// Writing with multi-arg transforms is not supported yet (hash semantics not finalized upstream).
if (partition_specification_field->has(Iceberg::f_source_ids))
throw Exception(
ErrorCodes::BAD_ARGUMENTS,
"Multi-argument partition transforms (source-ids) are not supported for writes. "
"Multi-argument transform evaluation is not yet implemented");

auto source_id = partition_specification_field->getValue<Int32>(Iceberg::f_source_id);
auto column_name = id_to_column[source_id];

Expand Down
1 change: 1 addition & 0 deletions src/Storages/ObjectStorage/DataLakes/Iceberg/Constant.h
Original file line number Diff line number Diff line change
Expand Up @@ -117,6 +117,7 @@ DEFINE_ICEBERG_FIELD_ALIAS(order_id, order-id);
DEFINE_ICEBERG_FIELD_ALIAS(default_sort_order_id, default-sort-order-id);
DEFINE_ICEBERG_FIELD_ALIAS(sort_orders, sort-orders);
DEFINE_ICEBERG_FIELD_ALIAS(source_id, source-id);
DEFINE_ICEBERG_FIELD_ALIAS(source_ids, source-ids);
DEFINE_ICEBERG_FIELD_ALIAS(partition_transform, transform);
DEFINE_ICEBERG_FIELD_ALIAS(partition_name, name);
DEFINE_ICEBERG_FIELD_ALIAS(default_spec_id, default-spec-id);
Expand Down
10 changes: 7 additions & 3 deletions src/Storages/ObjectStorage/DataLakes/Iceberg/ManifestFile.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@

#include <Common/logger_useful.h>
#include <fmt/format.h>
#include <fmt/ranges.h>


namespace DB::ErrorCodes
Expand Down Expand Up @@ -99,8 +100,10 @@ void requireDirectReferencedDataFileForPuffinDeletionVector(

static std::strong_ordering operator<=>(const PartitionSpecsEntry & lhs, const PartitionSpecsEntry & rhs)
{
return std::tie(lhs.source_id, lhs.transform_name, lhs.partition_name)
<=> std::tie(rhs.source_id, rhs.transform_name, rhs.partition_name);
if (auto cmp = lhs.source_ids <=> rhs.source_ids; cmp != std::strong_ordering::equal)
return cmp;
return std::tie(lhs.transform_name, lhs.partition_name)
<=> std::tie(rhs.transform_name, rhs.partition_name);
}

template <typename A>
Expand Down Expand Up @@ -138,7 +141,8 @@ static String dumpPartitionSpecification(const PartitionSpecification & partitio
{
const auto & entry = partition_specification[i];
answer += fmt::format(
"(source id: {}, transform name: {}, partition name: {})", entry.source_id, entry.transform_name, entry.partition_name);
"(source ids: [{}], transform name: {}, partition name: {})",
fmt::join(entry.source_ids, ", "), entry.transform_name, entry.partition_name);
if (i != partition_specification.size() - 1)
answer += ", ";
}
Expand Down
7 changes: 6 additions & 1 deletion src/Storages/ObjectStorage/DataLakes/Iceberg/ManifestFile.h
Original file line number Diff line number Diff line change
Expand Up @@ -61,9 +61,14 @@ String FileContentTypeToString(FileContentType type);

struct PartitionSpecsEntry
{
Int32 source_id;
/// For single-argument transforms (V1/V2 and single-arg V3) this holds one element.
/// For multi-argument V3 transforms (e.g. bucket over multiple columns) it holds multiple.
std::vector<Int32> source_ids;
String transform_name;
String partition_name;

/// Convenience: true when the transform references more than one source column (V3 multi-arg).
bool isMultiArg() const { return source_ids.size() > 1; }
};
using PartitionSpecification = std::vector<PartitionSpecsEntry>;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -182,17 +182,55 @@ std::shared_ptr<ManifestFileIterator> ManifestFileIterator::create(
{
auto partition_specification_field = partition_specification->getObject(static_cast<UInt32>(i));

auto source_id = partition_specification_field->getValue<Int32>(f_source_id);
/// Iceberg V3 spec: partition fields use either singular `source-id` (single-arg transforms)
/// or `source-ids` (multi-arg transforms, e.g. bucket over multiple columns). They are mutually exclusive.
bool has_source_id = partition_specification_field->has(f_source_id);
bool has_source_ids = partition_specification_field->has(f_source_ids);

std::vector<Int32> source_ids;
if (has_source_id && has_source_ids)
{
throw Exception(
ErrorCodes::ICEBERG_SPECIFICATION_VIOLATION,
"Partition field in manifest '{}' has both 'source-id' and 'source-ids' — they are mutually exclusive per the Iceberg spec",
path_to_manifest_file_);
}
else if (has_source_ids)
{
auto source_ids_array = partition_specification_field->getArray(f_source_ids);
for (UInt32 idx = 0; idx < source_ids_array->size(); ++idx)
source_ids.push_back(source_ids_array->getElement<Int32>(idx));
}
else if (has_source_id)
{
source_ids.push_back(partition_specification_field->getValue<Int32>(f_source_id));
}
else
{
throw Exception(
ErrorCodes::ICEBERG_SPECIFICATION_VIOLATION,
"Partition field in manifest '{}' has neither 'source-id' nor 'source-ids'",
path_to_manifest_file_);
}

auto transform_name = partition_specification_field->getValue<String>(f_partition_transform);
auto partition_name = partition_specification_field->getValue<String>(f_partition_name);
partition_spec_vec.emplace_back(PartitionSpecsEntry{std::move(source_ids), transform_name, partition_name});

/// Multi-argument transforms (V3): we cannot evaluate the transform, so skip pruning for this field.
/// Per the Iceberg V3 spec: "all v3 readers are required to read tables with unknown transforms,
/// ignoring the unsupported partition fields when filtering."
if (partition_spec_vec.back().isMultiArg())
continue;

auto source_id = partition_spec_vec.back().source_ids[0];
/// NOTE: tricky part to support RENAME column in partition key. Instead of some name
/// we use column internal number as it's name.
auto numeric_column_name = DB::backQuote(DB::toString(source_id));
std::optional<DB::NameAndTypePair> manifest_file_column_characteristics
= schema_processor.tryGetFieldCharacteristics(manifest_schema_id, source_id);
if (!manifest_file_column_characteristics.has_value())
continue;
auto transform_name = partition_specification_field->getValue<String>(f_partition_transform);
auto partition_name = partition_specification_field->getValue<String>(f_partition_name);
partition_spec_vec.emplace_back(source_id, transform_name, partition_name);
auto partition_ast = getASTFromTransform(transform_name, numeric_column_name, context_->getSettingsRef()[Setting::iceberg_partition_timezone]);
/// Unsupported partition key expression
if (partition_ast == nullptr)
Expand Down
26 changes: 26 additions & 0 deletions src/Storages/ObjectStorage/DataLakes/Iceberg/MetadataGenerator.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -526,6 +526,19 @@ void MetadataGenerator::generateDropColumnMetadata(const String & column_name)
ErrorCodes::BAD_ARGUMENTS,
"Cannot drop column '{}' (field id {}): it is referenced by the active sort order",
column_name, dropped_field_id);
/// Also check multi-arg V3 sort fields
if (sf->has(Iceberg::f_source_ids))
{
auto ids = sf->getArray(Iceberg::f_source_ids);
for (UInt32 k = 0; k < ids->size(); ++k)
{
if (ids->getElement<Int32>(k) == dropped_field_id)
throw Exception(
ErrorCodes::BAD_ARGUMENTS,
"Cannot drop column '{}' (field id {}): it is referenced by the active sort order (multi-arg transform)",
column_name, dropped_field_id);
}
}
}
break;
}
Expand Down Expand Up @@ -553,6 +566,19 @@ void MetadataGenerator::generateDropColumnMetadata(const String & column_name)
ErrorCodes::BAD_ARGUMENTS,
"Cannot drop column '{}' (field id {}): it is referenced by the active partition spec",
column_name, dropped_field_id);
/// Also check multi-arg V3 partition fields
if (pf->has(Iceberg::f_source_ids))
{
auto ids = pf->getArray(Iceberg::f_source_ids);
for (UInt32 k = 0; k < ids->size(); ++k)
{
if (ids->getElement<Int32>(k) == dropped_field_id)
throw Exception(
ErrorCodes::BAD_ARGUMENTS,
"Cannot drop column '{}' (field id {}): it is referenced by the active partition spec (multi-arg transform)",
column_name, dropped_field_id);
}
}
}
break;
}
Expand Down
93 changes: 85 additions & 8 deletions src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@
#include <string>
#include <unordered_set>
#include <config.h>
#include <fmt/ranges.h>
#include <Core/ColumnsWithTypeAndName.h>
#include <Core/Settings.h>
#include <Core/TypeId.h>
Expand Down Expand Up @@ -1460,6 +1461,16 @@ KeyDescription getSortingKeyDescriptionFromMetadata(Poco::JSON::Object::Ptr meta
for (UInt32 field_index = 0; field_index < fields->size(); ++field_index)
{
auto field = fields->getObject(field_index);

/// Iceberg V3: multi-argument transforms use `source-ids` instead of `source-id`.
/// We cannot evaluate multi-arg transforms, so skip these sort fields — this disables
/// the read-in-order optimization for such columns (safe: data is still correct).
if (field->has(f_source_ids))
continue;

if (!field->has(f_source_id))
continue;

auto source_id = field->getValue<Int64>(f_source_id);
auto column_name = source_id_to_column_name[source_id];
int direction = field->getValue<String>(f_direction) == "asc" ? 1 : -1;
Expand Down Expand Up @@ -1499,6 +1510,28 @@ KeyDescription getSortingKeyDescriptionFromMetadata(Poco::JSON::Object::Ptr meta
return KeyDescription::parse(order_by_str, column_description, {}, local_context, true);
}

/// Format a multi-argument partition field for display in Iceberg/Spark style, e.g. "bucket(16, a, b)".
static String formatPartitionFieldDisplayMultiArg(const String & iceberg_transform_name, const std::vector<String> & column_names)
{
std::string name = Poco::toLower(iceberg_transform_name);
String columns_joined = fmt::format("{}", fmt::join(column_names, ", "));

if (name.starts_with("bucket") && name.back() == ']')
{
auto p = name.find('[');
if (p != std::string::npos)
return "bucket(" + name.substr(p + 1, name.size() - p - 2) + ", " + columns_joined + ")";
}
if (name.starts_with("truncate") && name.back() == ']')
{
auto p = name.find('[');
if (p != std::string::npos)
return "truncate(" + name.substr(p + 1, name.size() - p - 2) + ", " + columns_joined + ")";
}
/// Fallback for unknown multi-arg transforms: show as transform(col1, col2, ...)
return name + "(" + columns_joined + ")";
}

/// Format one partition field for display in Iceberg/Spark style, e.g. "day(ts)" or "bucket(16, id)".
static String formatPartitionFieldDisplay(const String & iceberg_transform_name, const String & column_name)
{
Expand Down Expand Up @@ -1560,12 +1593,32 @@ std::optional<String> getPartitionKeyStringFromMetadata(Poco::JSON::Object::Ptr
for (UInt32 i = 0; i < fields->size(); ++i)
{
auto field = fields->getObject(i);
auto iceberg_transform_name = field->getValue<String>(f_transform);

/// Iceberg V3 multi-argument transform: `source-ids` is an array of column IDs.
if (field->has(f_source_ids))
{
auto source_ids_array = field->getArray(f_source_ids);
std::vector<String> column_names;
for (UInt32 idx = 0; idx < source_ids_array->size(); ++idx)
{
auto sid = source_ids_array->getElement<Int64>(idx);
auto it = source_id_to_column_name.find(sid);
if (it == source_id_to_column_name.end())
return std::nullopt;
column_names.push_back(it->second);
}
part_exprs.push_back(formatPartitionFieldDisplayMultiArg(iceberg_transform_name, column_names));
continue;
}

if (!field->has(f_source_id))
return std::nullopt;
auto source_id = field->getValue<Int64>(f_source_id);
auto it = source_id_to_column_name.find(source_id);
if (it == source_id_to_column_name.end())
return std::nullopt;
String column_name = it->second;
auto iceberg_transform_name = field->getValue<String>(f_transform);
part_exprs.push_back(formatPartitionFieldDisplay(iceberg_transform_name, column_name));
}
String result;
Expand Down Expand Up @@ -1600,14 +1653,38 @@ std::optional<String> getSortingKeyDisplayStringFromMetadata(Poco::JSON::Object:
for (UInt32 j = 0; j < sort_fields->size(); ++j)
{
auto field = sort_fields->getObject(j);
auto source_id = field->getValue<Int64>(f_source_id);
auto it = source_id_to_column_name.find(source_id);
if (it == source_id_to_column_name.end())
return std::nullopt;
String column_name = it->second;
String direction = field->getValue<String>(f_direction) == "asc" ? " asc" : " desc";
auto iceberg_transform_name = field->getValue<String>(f_transform);
String expr = formatPartitionFieldDisplay(iceberg_transform_name, column_name);
String direction = field->getValue<String>(f_direction) == "asc" ? " asc" : " desc";
String expr;

/// Iceberg V3 multi-argument transform
if (field->has(f_source_ids))
{
auto source_ids_array = field->getArray(f_source_ids);
std::vector<String> column_names;
for (UInt32 idx = 0; idx < source_ids_array->size(); ++idx)
{
auto sid = source_ids_array->getElement<Int64>(idx);
auto it = source_id_to_column_name.find(sid);
if (it == source_id_to_column_name.end())
return std::nullopt;
column_names.push_back(it->second);
}
expr = formatPartitionFieldDisplayMultiArg(iceberg_transform_name, column_names);
}
else if (field->has(f_source_id))
{
auto source_id = field->getValue<Int64>(f_source_id);
auto it = source_id_to_column_name.find(source_id);
if (it == source_id_to_column_name.end())
return std::nullopt;
expr = formatPartitionFieldDisplay(iceberg_transform_name, it->second);
}
else
{
return std::nullopt;
}

if (!result.empty())
result += ", ";
result += expr + direction;
Expand Down
Loading
Loading