From 4e29a3114953fe78c05ffedc39fcbe4cf61c9c31 Mon Sep 17 00:00:00 2001 From: Kanthi Subramanian Date: Wed, 30 Sep 2026 20:49:30 +0200 Subject: [PATCH 1/3] Backport of #116171 --- .../DataLakes/Iceberg/AvroSchema.h | 312 ++++++++++++++++++ .../DataLakes/Iceberg/Compaction.cpp | 17 +- .../DataLakes/Iceberg/IcebergMetadata.cpp | 8 +- .../DataLakes/Iceberg/IcebergWrites.cpp | 239 ++++++++------ .../DataLakes/Iceberg/IcebergWrites.h | 3 +- .../DataLakes/Iceberg/Mutations.cpp | 11 +- .../test_writes_v3_row_lineage.py | 100 ++++++ .../test_writes_v3_row_lineage_partitioned.py | 101 ++++++ .../test_row_lineage.py | 196 ++++++++++- .../test_row_lineage_pruning.py | 161 +++++++++ 10 files changed, 1045 insertions(+), 103 deletions(-) create mode 100644 tests/integration/test_storage_iceberg_no_spark/test_writes_v3_row_lineage.py create mode 100644 tests/integration/test_storage_iceberg_no_spark/test_writes_v3_row_lineage_partitioned.py diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/AvroSchema.h b/src/Storages/ObjectStorage/DataLakes/Iceberg/AvroSchema.h index 97c832760a11..6546247e3c91 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/AvroSchema.h +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/AvroSchema.h @@ -199,6 +199,96 @@ static constexpr const char * manifest_list_v2_schema = R"( } )"; +static constexpr const char * manifest_list_v3_schema = R"( +{ + "type": "record", + "name": "manifest_file", + "fields": [ + {"name": "manifest_path", "type": "string", "doc": "Location URI with FS scheme", "field-id": 500}, + {"name": "manifest_length", "type": "long", "doc": "Total file size in bytes", "field-id": 501}, + {"name": "partition_spec_id", "type": "int", "doc": "Spec ID used to write", "field-id": 502}, + {"name": "content", "type": "int", "doc": "Contents of the manifest: 0=data, 1=deletes", "field-id": 517}, + { + "name": "sequence_number", + "type": "long", + "doc": "Sequence number when the manifest was added", + "field-id": 515 + }, + { + "name": "min_sequence_number", + "type": "long", + "doc": "Lowest sequence number in the manifest", + "field-id": 516 + }, + {"name": "added_snapshot_id", "type": "long", "doc": "Snapshot ID that added the manifest", "field-id": 503}, + {"name": "added_files_count", "type": "int", "doc": "Added entry count", "field-id": 504}, + {"name": "existing_files_count", "type": "int", "doc": "Existing entry count", "field-id": 505}, + {"name": "deleted_files_count", "type": "int", "doc": "Deleted entry count", "field-id": 506}, + {"name": "added_rows_count", "type": "long", "doc": "Added rows count", "field-id": 512}, + {"name": "existing_rows_count", "type": "long", "doc": "Existing rows count", "field-id": 513}, + {"name": "deleted_rows_count", "type": "long", "doc": "Deleted rows count", "field-id": 514}, + { + "name": "partitions", + "type": [ + "null", + { + "type": "array", + "items": + { + "type": "record", + "name": "r508", + "fields": [ + { + "name": "contains_null", + "type": "boolean", + "doc": "True if any file has a null partition value", + "field-id": 509 + }, + { + "name": "contains_nan", + "type": ["null", "boolean"], + "doc": "True if any file has a nan partition value", + "field-id": 518 + }, + { + "name": "lower_bound", + "type": ["null", "bytes"], + "doc": "Partition lower bound for all files", + "field-id": 510 + }, + { + "name": "upper_bound", + "type": ["null", "bytes"], + "doc": "Partition upper bound for all files", + "field-id": 511 + } + ] + }, + "element-id": 508 + } + ], + "doc": "Summary for each partition", + "field-id": 507 + }, + { + "name": "key_metadata", + "type": ["null", "bytes"], + "doc": "Encryption key metadata blob", + "default": null, + "field-id": 519 + }, + { + "name": "first_row_id", + "type": ["null", "long"], + "doc": "Starting row ID to assign to new rows in ADDED data files", + "default": null, + "field-id": 520 + } + ] +} +)"; + + /// NOTE: This string is just a template for the actual schema. To use it, you must first replace "#" with the correct value. static constexpr const char * manifest_entry_v1_schema = R"( { @@ -619,4 +709,226 @@ static constexpr const char * data_file_sidecar_schema = R"( } )"; +static constexpr const char * manifest_entry_v3_schema = R"( +{ + "type": "record", + "name": "manifest_entry", + "fields": [ + {"name": "status", "type": "int", "field-id": 0}, + {"name": "snapshot_id", "type": ["null", "long"], "field-id": 1}, + {"name": "sequence_number", "type": ["null", "long"], "field-id": 3}, + {"name": "file_sequence_number", "type": ["null", "long"], "field-id": 4}, + { + "name": "data_file", + "type": + { + "type": "record", + "name": "r2", + "fields": [ + {"name": "content", "type": "int", "doc": "Type of content stored by the data file", "field-id": 134}, + {"name": "file_path", "type": "string", "doc": "Location URI with FS scheme", "field-id": 100}, + { + "name": "file_format", + "type": "string", + "doc": "File format name: avro, orc, or parquet", + "field-id": 101 + }, + { + "name": "partition", + "type": + { + "type": "record", + "name": "r102", + "fields": # + }, + "field-id": 102 + }, + {"name": "record_count", "type": "long", "doc": "Number of records in the file", "field-id": 103}, + {"name": "file_size_in_bytes", "type": "long", "doc": "Total file size in bytes", "field-id": 104}, + { + "name": "column_sizes", + "type": [ + "null", + { + "type": "array", + "items": + { + "type": "record", + "name": "k117_v118", + "fields": [ + {"name": "key", "type": "int", "field-id": 117}, + {"name": "value", "type": "long", "field-id": 118} + ] + }, + "logicalType": "map" + } + ], + "doc": "Map of column id to total size on disk", + "field-id": 108 + }, + { + "name": "value_counts", + "type": [ + "null", + { + "type": "array", + "items": + { + "type": "record", + "name": "k119_v120", + "fields": [ + {"name": "key", "type": "int", "field-id": 119}, + {"name": "value", "type": "long", "field-id": 120} + ] + }, + "logicalType": "map" + } + ], + "doc": "Map of column id to total count, including null and NaN", + "field-id": 109 + }, + { + "name": "null_value_counts", + "type": [ + "null", + { + "type": "array", + "items": + { + "type": "record", + "name": "k121_v122", + "fields": [ + {"name": "key", "type": "int", "field-id": 121}, + {"name": "value", "type": "long", "field-id": 122} + ] + }, + "logicalType": "map" + } + ], + "doc": "Map of column id to null value count", + "field-id": 110 + }, + { + "name": "nan_value_counts", + "type": [ + "null", + { + "type": "array", + "items": + { + "type": "record", + "name": "k138_v139", + "fields": [ + {"name": "key", "type": "int", "field-id": 138}, + {"name": "value", "type": "long", "field-id": 139} + ] + }, + "logicalType": "map" + } + ], + "doc": "Map of column id to number of NaN values in the column", + "field-id": 137 + }, + { + "name": "lower_bounds", + "type": [ + "null", + { + "type": "array", + "items": + { + "type": "record", + "name": "k126_v127", + "fields": [ + {"name": "key", "type": "int", "field-id": 126}, + {"name": "value", "type": "bytes", "field-id": 127} + ] + }, + "logicalType": "map" + } + ], + "doc": "Map of column id to lower bound", + "field-id": 125 + }, + { + "name": "upper_bounds", + "type": [ + "null", + { + "type": "array", + "items": + { + "type": "record", + "name": "k129_v130", + "fields": [ + {"name": "key", "type": "int", "field-id": 129}, + {"name": "value", "type": "bytes", "field-id": 130} + ] + }, + "logicalType": "map" + } + ], + "doc": "Map of column id to upper bound", + "field-id": 128 + }, + { + "name": "key_metadata", + "type": ["null", "bytes"], + "doc": "Encryption key metadata blob", + "field-id": 131 + }, + { + "name": "split_offsets", + "type": ["null", {"type": "array", "items": "long", "element-id": 133}], + "doc": "Splittable offsets", + "field-id": 132 + }, + { + "name": "equality_ids", + "type": ["null", {"type": "array", "items": "int", "element-id": 136}], + "doc": "Field ids used to determine row equality for delete files", + "field-id": 135 + }, + { + "name": "sort_order_id", + "type": ["null", "int"], + "doc": "Sort order ID", + "field-id": 140 + }, + { + "name": "first_row_id", + "type": ["null", "long"], + "doc": "The _row_id for the first row in the data file", + "default": null, + "field-id": 142 + }, + { + "name": "referenced_data_file", + "type": ["null", "string"], + "doc": "Fully qualified location (URI with FS scheme) of a data file that all deletes reference", + "default": null, + "field-id": 143 + }, + { + "name": "content_offset", + "type": ["null", "long"], + "doc": "The offset in the file where the content starts", + "default": null, + "field-id": 144 + }, + { + "name": "content_size_in_bytes", + "type": ["null", "long"], + "doc": "The length of referenced content stored in the file", + "default": null, + "field-id": 145 + } + ] + }, + "field-id": 2 + } + ] +} +)"; + } diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/Compaction.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/Compaction.cpp index 5855504b22ac..c4b3c7baf4e0 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/Compaction.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/Compaction.cpp @@ -1041,6 +1041,7 @@ static void writeMetadataFiles( Poco::JSON::Object::Ptr initial_metadata_object = plan.initial_metadata_object; std::unordered_map manifest_file_renamings; std::unordered_map manifest_file_sizes; + std::unordered_map manifest_file_row_counts; { std::unordered_map, std::unordered_set> grouped_by_manifest_files_result; @@ -1114,6 +1115,10 @@ static void writeMetadataFiles( file_byte_counts.push_back(0); } } + Int64 manifest_rows = 0; + for (const auto rows : file_row_counts) + manifest_rows += static_cast(rows); + manifest_file_row_counts[manifest_entry->patched_path] = manifest_rows; generateManifestFile( metadata_object, partition_columns, @@ -1170,8 +1175,12 @@ static void writeMetadataFiles( } } std::vector per_manifest_sizes; + std::vector per_manifest_row_counts; for (const auto & entry : renamed_manifest_entries) + { per_manifest_sizes.push_back(manifest_file_sizes[entry]); + per_manifest_row_counts.push_back(manifest_file_row_counts[entry]); + } auto buffer_manifest_list = object_storage->writeObject( StoredObject(path_resolver.resolve(renamed_manifest_list)), WriteMode::Rewrite, @@ -1189,7 +1198,13 @@ static void writeMetadataFiles( per_manifest_sizes, *buffer_manifest_list, Iceberg::FileContentType::DATA, - false); + false, + /* per_entry_content_types */ {}, + /* existing_entry_counts */ {}, + /* carry_forward_manifest_paths */ {}, + /* entry_partition_spec_ids */ {}, + /* entry_partition_summaries */ {}, + per_manifest_row_counts); buffer_manifest_list->finalize(); } diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.cpp index 945303783fe7..31b2c4da0655 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.cpp @@ -2005,7 +2005,13 @@ std::optional IcebergMetadata::commitImport try { generateManifestList( - resolver, metadata, object_storage, *secondary_storages, context, {manifest_entry_path}, new_snapshot, {manifest_lengths}, *buffer_manifest_list, Iceberg::FileContentType::DATA, true); + resolver, metadata, object_storage, *secondary_storages, context, {manifest_entry_path}, new_snapshot, {manifest_lengths}, *buffer_manifest_list, Iceberg::FileContentType::DATA, true, + /* per_entry_content_types */ {}, + /* existing_entry_counts */ {}, + /* carry_forward_manifest_paths */ {}, + /* entry_partition_spec_ids */ {}, + /* entry_partition_summaries */ {}, + /* entry_row_counts */ {total_rows}); buffer_manifest_list->finalize(); } catch (...) diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergWrites.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergWrites.cpp index 0adbf4f2c0e5..6d94f80e33f0 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergWrites.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergWrites.cpp @@ -687,8 +687,10 @@ void generateManifestFile( String schema_representation; if (version == 1) schema_representation = manifest_entry_v1_schema; - else if (version == 2 || version == 3) + else if (version == 2) schema_representation = manifest_entry_v2_schema; + else if (version == 3) + schema_representation = manifest_entry_v3_schema; else throw Exception(ErrorCodes::BAD_ARGUMENTS, "Unsupported iceberg format-version {}", version); @@ -1038,7 +1040,8 @@ void generateManifestList( const std::vector & existing_entry_counts, const std::unordered_set & carry_forward_manifest_paths, const std::vector & entry_partition_spec_ids, - const std::vector>> & entry_partition_summaries) + const std::vector>> & entry_partition_summaries, + const std::vector & entry_row_counts) { chassert( per_entry_content_types.empty() || per_entry_content_types.size() == manifest_entry_names.size(), @@ -1054,12 +1057,23 @@ void generateManifestList( existing_entry_counts.empty() || existing_entry_counts.size() == manifest_entry_names.size(), "existing_entry_counts size does not match number of manifest entries"); const bool manifest_only_rewrite = !existing_entry_counts.empty(); + if (!manifest_only_rewrite && entry_row_counts.size() != manifest_entry_names.size()) + throw Exception( + ErrorCodes::LOGICAL_ERROR, + "Iceberg manifest list needs one row count per manifest entry, got {} counts for {} entries", + entry_row_counts.size(), + manifest_entry_names.size()); + Int32 version = metadata->getValue(Iceberg::f_format_version); String schema_representation; if (version == 1) schema_representation = manifest_list_v1_schema; - else + else if (version == 2) schema_representation = manifest_list_v2_schema; + else if (version == 3) + schema_representation = manifest_list_v3_schema; + else + throw Exception(ErrorCodes::BAD_ARGUMENTS, "Unsupported iceberg format-version {}", version); // For empty manifest list (e.g. TRUNCATE), write a valid Avro container // file manually so we can embed the full schema JSON with field-ids intact, @@ -1099,6 +1113,102 @@ void generateManifestList( avro::DataFileWriter writer(std::move(adapter), schema); writer.setMetadata(Iceberg::f_format_version, std::to_string(version)); + Int64 next_first_row_id = 0; + if (version > 2) + { + if (manifest_only_rewrite) + throw Exception( + ErrorCodes::NOT_IMPLEMENTED, + "Rewriting an Iceberg manifest list is not supported for format-version 3: the row-lineage " + "'first_row_id' of a rewritten manifest is not round-tripped"); + if (!new_snapshot->has(Iceberg::f_first_row_id) || new_snapshot->isNull(Iceberg::f_first_row_id)) + throw Exception( + ErrorCodes::LOGICAL_ERROR, + "Snapshot of a format-version 3 Iceberg table has no '{}', cannot assign row ids to the added data files", + Iceberg::f_first_row_id); + next_first_row_id = new_snapshot->getValue(Iceberg::f_first_row_id); + } + + /// Copy entries from the parent snapshot's manifest list: `use_previous_snapshots` copies all, `carry_forward_manifest_paths` copies only the listed manifests. + if (use_previous_snapshots || !carry_forward_manifest_paths.empty()) + { + auto parent_snapshot_id = new_snapshot->getValue(Iceberg::f_parent_snapshot_id); + auto snapshots = metadata->getArray(Iceberg::f_snapshots); + for (size_t i = 0; i < snapshots->size(); ++i) + { + if (snapshots->getObject(static_cast(i))->getValue(Iceberg::f_metadata_snapshot_id) == parent_snapshot_id) + { + auto manifest_list = Iceberg::IcebergPathFromMetadata::deserialize( + snapshots->getObject(static_cast(i))->getValue(Iceberg::f_manifest_list)); + + auto [manifest_list_storage, resolved_manifest_list_path] = resolveObjectStorageForPath( + path_resolver.getTableLocation(), manifest_list.serialize(), object_storage, secondary_storages, context, path_resolver); + forEachAvroEntry(resolved_manifest_list_path, manifest_list_storage, context, "IcebergWrites", + [&](const avro::GenericDatum & datum) + { + const avro::GenericRecord & old_entry = datum.value(); + /// When a path filter is supplied, copy only the matching entries. + if (!carry_forward_manifest_paths.empty() + && !carry_forward_manifest_paths.contains(old_entry.field(Iceberg::f_manifest_path).value())) + return; + avro::GenericDatum new_datum(schema.root()); + avro::GenericRecord & new_entry = new_datum.value(); + new_entry.field(f_manifest_path) = old_entry.field(Iceberg::f_manifest_path); + new_entry.field(f_manifest_length) = old_entry.field(Iceberg::f_manifest_length); + new_entry.field(f_partition_spec_id) = old_entry.field(Iceberg::f_partition_spec_id); + /// iceberg-spark changed `f_added_snapshot_id` from 'null, long' to 'long' (apache/iceberg#11626); rewrite with the new schema in case we read the old type. + if (old_entry.hasField(Iceberg::f_added_snapshot_id)) + { + const avro::GenericDatum & old_added_snapshot_id_entry = old_entry.field(Iceberg::f_added_snapshot_id); + if (old_added_snapshot_id_entry.isUnion()) + { + if (old_added_snapshot_id_entry.unionBranch() == 0) /// it means add_snapshot_id is null + { + /// This only happens when we read data written by a old version of iceberg, which violates the spec of iceberg. + throw Exception( + ErrorCodes::ICEBERG_SPECIFICATION_VIOLATION, + "Manifest list {} has null value for field '{}', but it is required", + resolved_manifest_list_path, + Iceberg::f_added_snapshot_id); + } + } + new_entry.field(f_added_snapshot_id) = old_added_snapshot_id_entry.value(); + } + else + /// This only happens when we read data written by a old version of iceberg, which violates the spec of iceberg. + throw Exception( + ErrorCodes::ICEBERG_SPECIFICATION_VIOLATION, + "Manifest list {} has null value for field '{}', but it is required", + resolved_manifest_list_path, + Iceberg::f_added_snapshot_id); + auto add_field_to_datum = [&](const String & field) + { + if (old_entry.hasField(field)) + new_entry.field(field) = old_entry.field(field); + }; + add_field_to_datum(Iceberg::f_added_files_count); + add_field_to_datum(Iceberg::f_existing_files_count); + add_field_to_datum(Iceberg::f_deleted_files_count); + add_field_to_datum(Iceberg::f_partitions); + add_field_to_datum(Iceberg::f_added_rows_count); + add_field_to_datum(Iceberg::f_existing_rows_count); + add_field_to_datum(Iceberg::f_deleted_rows_count); + add_field_to_datum(Iceberg::f_key_metadata); + if (version > 1) + { + add_field_to_datum(Iceberg::f_content); + add_field_to_datum(Iceberg::f_sequence_number); + add_field_to_datum(Iceberg::f_min_sequence_number); + } + if (version > 2) + add_field_to_datum(Iceberg::f_manifest_first_row_id); + writer.write(new_datum); + }); + break; + } + } + } + for (size_t entry_idx = 0; entry_idx < manifest_entry_names.size(); ++entry_idx) { avro::GenericDatum entry_datum(schema.root()); @@ -1215,102 +1325,22 @@ void generateManifestList( if (summary->has(Iceberg::f_added_position_deletes)) entry.field(Iceberg::f_deleted_rows_count) = summary->getValue(Iceberg::f_added_position_deletes); } - - if (summary->has(Iceberg::f_added_records)) - { - set_versioned_field( - summary->getValue(Iceberg::f_added_records), - Iceberg::f_added_rows_count); - } - else - { - set_versioned_field(summary->getValue(Iceberg::f_added_position_deletes), Iceberg::f_added_rows_count); - } - set_versioned_field( + const Int64 added_rows_count = entry_row_counts[entry_idx]; + setVersionedField(entry, added_rows_count, Iceberg::f_added_rows_count); + setVersionedField( + entry, 0, Iceberg::f_existing_rows_count); - set_versioned_field(0, Iceberg::f_deleted_rows_count); + setVersionedField(entry, 0, Iceberg::f_deleted_rows_count); - writer.write(entry_datum); - } - - /// Copy entries from the parent snapshot's manifest list: `use_previous_snapshots` copies all, `carry_forward_manifest_paths` copies only the listed manifests. - if (use_previous_snapshots || !carry_forward_manifest_paths.empty()) - { - auto parent_snapshot_id = new_snapshot->getValue(Iceberg::f_parent_snapshot_id); - auto snapshots = metadata->getArray(Iceberg::f_snapshots); - for (size_t i = 0; i < snapshots->size(); ++i) + /// Only data files get row ids, so a delete manifest leaves the field null. + if (version > 2 && entry_content == Iceberg::FileContentType::DATA) { - if (snapshots->getObject(static_cast(i))->getValue(Iceberg::f_metadata_snapshot_id) == parent_snapshot_id) - { - auto manifest_list = Iceberg::IcebergPathFromMetadata::deserialize( - snapshots->getObject(static_cast(i))->getValue(Iceberg::f_manifest_list)); - - auto [manifest_list_storage, resolved_manifest_list_path] = resolveObjectStorageForPath( - path_resolver.getTableLocation(), manifest_list.serialize(), object_storage, secondary_storages, context, path_resolver); - forEachAvroEntry(resolved_manifest_list_path, manifest_list_storage, context, "IcebergWrites", - [&](const avro::GenericDatum & datum) - { - const avro::GenericRecord & old_entry = datum.value(); - /// When a path filter is supplied, copy only the matching entries. - if (!carry_forward_manifest_paths.empty() - && !carry_forward_manifest_paths.contains(old_entry.field(Iceberg::f_manifest_path).value())) - return; - avro::GenericDatum new_datum(schema.root()); - avro::GenericRecord & new_entry = new_datum.value(); - new_entry.field(f_manifest_path) = old_entry.field(Iceberg::f_manifest_path); - new_entry.field(f_manifest_length) = old_entry.field(Iceberg::f_manifest_length); - new_entry.field(f_partition_spec_id) = old_entry.field(Iceberg::f_partition_spec_id); - /// iceberg-spark changed `f_added_snapshot_id` from 'null, long' to 'long' (apache/iceberg#11626); rewrite with the new schema in case we read the old type. - if (old_entry.hasField(Iceberg::f_added_snapshot_id)) - { - const avro::GenericDatum & old_added_snapshot_id_entry = old_entry.field(Iceberg::f_added_snapshot_id); - if (old_added_snapshot_id_entry.isUnion()) - { - if (old_added_snapshot_id_entry.unionBranch() == 0) /// it means add_snapshot_id is null - { - /// This only happens when we read data written by a old version of iceberg, which violates the spec of iceberg. - throw Exception( - ErrorCodes::ICEBERG_SPECIFICATION_VIOLATION, - "Manifest list {} has null value for field '{}', but it is required", - resolved_manifest_list_path, - Iceberg::f_added_snapshot_id); - } - } - new_entry.field(f_added_snapshot_id) = old_added_snapshot_id_entry.value(); - } - else - /// This only happens when we read data written by a old version of iceberg, which violates the spec of iceberg. - throw Exception( - ErrorCodes::ICEBERG_SPECIFICATION_VIOLATION, - "Manifest list {} has null value for field '{}', but it is required", - resolved_manifest_list_path, - Iceberg::f_added_snapshot_id); - auto add_field_to_datum = [&](const String & field) - { - if (old_entry.hasField(field)) - new_entry.field(field) = old_entry.field(field); - }; - add_field_to_datum(Iceberg::f_added_files_count); - add_field_to_datum(Iceberg::f_existing_files_count); - add_field_to_datum(Iceberg::f_deleted_files_count); - add_field_to_datum(Iceberg::f_partitions); - add_field_to_datum(Iceberg::f_added_rows_count); - add_field_to_datum(Iceberg::f_existing_rows_count); - add_field_to_datum(Iceberg::f_deleted_rows_count); - add_field_to_datum(Iceberg::f_key_metadata); - /// v2 and v3 share the manifest-list schema, so these fields exist for both. - if (version > 1) - { - add_field_to_datum(Iceberg::f_content); - add_field_to_datum(Iceberg::f_sequence_number); - add_field_to_datum(Iceberg::f_min_sequence_number); - } - writer.write(new_datum); - }); - break; - } + setVersionedField(entry, next_first_row_id, Iceberg::f_manifest_first_row_id); + next_first_row_id += added_rows_count; } + + writer.write(entry_datum); } writer.close(); @@ -1597,6 +1627,7 @@ bool IcebergStorageSink::initializeMetadata() Strings manifest_entries_in_storage; std::vector manifest_entries; std::vector manifest_entry_sizes; + std::vector manifest_entry_row_counts; auto cleanup = [&] (bool retry_because_of_metadata_conflict) { @@ -1678,6 +1709,10 @@ bool IcebergStorageSink::initializeMetadata() auto manifest_entry_path = filename_generator.generateManifestEntryName(); manifest_entries_in_storage.push_back(resolver.resolve(manifest_entry_path)); manifest_entries.push_back(manifest_entry_path); + Int64 manifest_row_count = 0; + for (UInt64 data_file_row_count : writer.getDataFileRowCounts()) + manifest_row_count += static_cast(data_file_row_count); + manifest_entry_row_counts.push_back(manifest_row_count); auto buffer_manifest_entry = object_storage->writeObject( StoredObject(resolver.resolve(manifest_entry_path)), WriteMode::Rewrite, std::nullopt, DBMS_DEFAULT_BUFFER_SIZE, context->getWriteSettings()); @@ -1720,8 +1755,20 @@ bool IcebergStorageSink::initializeMetadata() try { generateManifestList( - persistent_table_components.path_resolver, metadata, object_storage, *secondary_storages, context, manifest_entries, new_snapshot, manifest_entry_sizes, *buffer_manifest_list, Iceberg::FileContentType::DATA, - /* use_previous_snapshots = */ true); + persistent_table_components.path_resolver, + metadata, object_storage, *secondary_storages, context, + manifest_entries, + new_snapshot, + manifest_entry_sizes, + *buffer_manifest_list, + Iceberg::FileContentType::DATA, + /* use_previous_snapshots = */ true, + /* per_entry_content_types = */ {}, + /* existing_entry_counts = */ {}, + /* carry_forward_manifest_paths = */ {}, + /* entry_partition_spec_ids = */ {}, + /* entry_partition_summaries = */ {}, + manifest_entry_row_counts); buffer_manifest_list->finalize(); } catch (...) diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergWrites.h b/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergWrites.h index a7c6c6a322c5..7dc951cef09f 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergWrites.h +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergWrites.h @@ -155,7 +155,8 @@ void generateManifestList( const std::vector & existing_entry_counts = {}, const std::unordered_set & carry_forward_manifest_paths = {}, const std::vector & entry_partition_spec_ids = {}, - const std::vector>> & entry_partition_summaries = {}); + const std::vector>> & entry_partition_summaries = {}, + const std::vector & entry_row_counts = {}); std::string getIcebergExportPartSidecarStoragePath(const String & data_file_storage_path); diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/Mutations.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/Mutations.cpp index 2933f50f9296..321fab80121d 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/Mutations.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/Mutations.cpp @@ -490,6 +490,7 @@ static bool writeMetadataFiles( auto manifest_entries_in_storage = std::make_shared(); std::vector manifest_entries; std::vector manifest_entry_sizes; + std::vector manifest_entry_row_counts; auto cleanup = [object_storage, &delete_filenames, &path_resolver, manifest_entries_in_storage, storage_manifest_list_name, storage_metadata_name]() { @@ -516,6 +517,7 @@ static bool writeMetadataFiles( auto manifest_entry_path = filename_generator.generateManifestEntryName(); manifest_entries_in_storage->push_back(path_resolver.resolve(manifest_entry_path)); manifest_entries.push_back(manifest_entry_path); + manifest_entry_row_counts.push_back(delete_filename.total_rows); auto buffer_manifest_entry = object_storage->writeObject( StoredObject(path_resolver.resolve(manifest_entry_path)), @@ -576,7 +578,14 @@ static bool writeMetadataFiles( new_snapshot, manifest_entry_sizes, *buffer_manifest_list, - content_type); + content_type, + /* use_previous_snapshots */ true, + /* per_entry_content_types */ {}, + /* existing_entry_counts */ {}, + /* carry_forward_manifest_paths */ {}, + /* entry_partition_spec_ids */ {}, + /* entry_partition_summaries */ {}, + manifest_entry_row_counts); buffer_manifest_list->finalize(); } catch (...) diff --git a/tests/integration/test_storage_iceberg_no_spark/test_writes_v3_row_lineage.py b/tests/integration/test_storage_iceberg_no_spark/test_writes_v3_row_lineage.py new file mode 100644 index 000000000000..eb1c7ea29a74 --- /dev/null +++ b/tests/integration/test_storage_iceberg_no_spark/test_writes_v3_row_lineage.py @@ -0,0 +1,100 @@ +import json + +import pytest + +from avro.datafile import DataFileReader +from avro.io import DatumReader + +from helpers.iceberg_utils import ( + create_iceberg_table, + get_uuid_str, +) + + +def _metadata_dir(table_name): + return f"/var/lib/clickhouse/user_files/iceberg_data/default/{table_name}/metadata" + + +def _latest_metadata(instance, table_name): + latest = instance.exec_in_container( + ["bash", "-c", f"ls -v {_metadata_dir(table_name)}/v*.metadata.json | tail -1"] + ).strip() + return json.loads(instance.exec_in_container(["cat", latest])) + + +def _read_avro(instance, remote_path, local_path): + instance.copy_file_from_container(remote_path, local_path) + with open(local_path, "rb") as handle: + reader = DataFileReader(handle, DatumReader()) + records = list(reader) + schema = json.loads(reader.meta["avro.schema"].decode("utf-8")) + reader.close() + return records, schema + + +def _field_ids(schema_fields): + return {field["name"]: field.get("field-id") for field in schema_fields} + + +def _sorted_snapshots(metadata): + return sorted(metadata["snapshots"], key=lambda s: s["sequence-number"]) + + +@pytest.mark.parametrize("storage_type", ["local"]) +def test_v3_row_lineage_written(started_cluster_iceberg_no_spark, storage_type, tmp_path): + instance = started_cluster_iceberg_no_spark.instances["node1"] + table_name = "test_v3_row_lineage_" + storage_type + "_" + get_uuid_str() + + create_iceberg_table( + storage_type, + instance, + table_name, + started_cluster_iceberg_no_spark, + "(id Int32, s String)", + format_version=3, + ) + + instance.query(f"INSERT INTO {table_name} VALUES (1, 'a'), (2, 'b')") + instance.query(f"INSERT INTO {table_name} VALUES (3, 'c'), (4, 'd'), (5, 'e')") + + assert instance.query(f"SELECT count() FROM {table_name}").strip() == "5" + + metadata = _latest_metadata(instance, table_name) + assert metadata["format-version"] == 3 + + snapshots = _sorted_snapshots(metadata) + assert len(snapshots) == 2 + + assert [s["first-row-id"] for s in snapshots] == [0, 2] + assert [s["added-rows"] for s in snapshots] == [2, 3] + assert metadata["next-row-id"] == 5 + + for snapshot in snapshots: + assert snapshot["added-rows"] == int(snapshot["summary"]["added-records"]) + + manifest_list_path = snapshots[-1]["manifest-list"] + manifest_list, manifest_list_schema = _read_avro( + instance, manifest_list_path, str(tmp_path / "manifest_list.avro") + ) + + assert _field_ids(manifest_list_schema["fields"])["first_row_id"] == 520 + + ranges = sorted( + (entry["first_row_id"], entry["added_rows_count"]) for entry in manifest_list + ) + assert ranges == [(0, 2), (2, 3)] + + next_free = 0 + for first_row_id, added_rows_count in ranges: + assert first_row_id == next_free + next_free += added_rows_count + assert next_free == metadata["next-row-id"] + + for index, entry in enumerate(manifest_list): + _, manifest_schema = _read_avro( + instance, entry["manifest_path"], str(tmp_path / f"manifest_{index}.avro") + ) + data_file_schema = next( + field for field in manifest_schema["fields"] if field["name"] == "data_file" + ) + assert _field_ids(data_file_schema["type"]["fields"])["first_row_id"] == 142 diff --git a/tests/integration/test_storage_iceberg_no_spark/test_writes_v3_row_lineage_partitioned.py b/tests/integration/test_storage_iceberg_no_spark/test_writes_v3_row_lineage_partitioned.py new file mode 100644 index 000000000000..8a8d44d1b1c7 --- /dev/null +++ b/tests/integration/test_storage_iceberg_no_spark/test_writes_v3_row_lineage_partitioned.py @@ -0,0 +1,101 @@ +import glob +import json +import os + +import avro.datafile +import avro.io +import pytest + +from helpers.iceberg_utils import ( + create_iceberg_table, + default_download_directory, + get_uuid_str, +) + +TABLE_ROOT = "/var/lib/clickhouse/user_files/iceberg_data/default" + + +def _download_table(started_cluster, storage_type, table_name): + path = f"{TABLE_ROOT}/{table_name}/" + default_download_directory(started_cluster, storage_type, path, path) + return path + + +def _read_avro(path): + with open(path, "rb") as f: + return list(avro.datafile.DataFileReader(f, avro.io.DatumReader())) + + +def _latest_metadata(table_path): + names = glob.glob(f"{table_path}/metadata/v*.metadata.json") + latest = max(names, key=lambda name: int(os.path.basename(name)[1:].split(".")[0])) + with open(latest) as f: + return json.load(f) + + +def _manifest_list_of(table_path, snapshot_id): + matches = glob.glob(f"{table_path}/metadata/snap-{snapshot_id}-*.avro") + assert len(matches) == 1, matches + return matches[0] + + +@pytest.mark.parametrize("storage_type", ["s3", "local"]) +def test_writes_v3_row_lineage_partitioned(started_cluster_iceberg_no_spark, storage_type): + instance = started_cluster_iceberg_no_spark.instances["node1"] + TABLE_NAME = ( + "test_writes_v3_row_lineage_partitioned_" + storage_type + "_" + get_uuid_str() + ) + + create_iceberg_table( + storage_type, + instance, + TABLE_NAME, + started_cluster_iceberg_no_spark, + "(part Int32, id Int32)", + format_version=3, + partition_by="part", + order_by="id", + ) + + instance.query( + f"INSERT INTO {TABLE_NAME} VALUES (1, 0), (1, 1), (1, 2), (2, 3), (2, 4)", + settings={"allow_insert_into_iceberg": 1}, + ) + instance.query( + f"INSERT INTO {TABLE_NAME} VALUES (1, 5), (2, 6), (2, 7), (2, 8)", + settings={"allow_insert_into_iceberg": 1}, + ) + + table_path = _download_table( + started_cluster_iceberg_no_spark, storage_type, TABLE_NAME + ) + metadata = _latest_metadata(table_path) + + assert metadata["format-version"] == 3 + assert metadata["next-row-id"] == 9 + + snapshots = sorted(metadata["snapshots"], key=lambda s: s["sequence-number"]) + assert [snapshot["first-row-id"] for snapshot in snapshots] == [0, 5] + assert [snapshot["added-rows"] for snapshot in snapshots] == [5, 4] + + entries = _read_avro(_manifest_list_of(table_path, metadata["current-snapshot-id"])) + assert len(entries) == 4 + + ranges_by_snapshot = {} + for entry in entries: + ranges_by_snapshot.setdefault(entry["added_snapshot_id"], []).append( + (entry["first_row_id"], entry["added_rows_count"]) + ) + + for snapshot in snapshots: + ranges = sorted(ranges_by_snapshot[snapshot["snapshot-id"]]) + assert len(ranges) == 2, ranges + next_free = snapshot["first-row-id"] + for first_row_id, added_rows_count in ranges: + assert first_row_id == next_free + next_free += added_rows_count + assert next_free == snapshot["first-row-id"] + snapshot["added-rows"] + + assert instance.query( + f"SELECT _row_id FROM {TABLE_NAME} ORDER BY _row_id FORMAT TSV" + ).split() == [str(row_id) for row_id in range(9)] diff --git a/tests/integration/test_storage_iceberg_with_spark/test_row_lineage.py b/tests/integration/test_storage_iceberg_with_spark/test_row_lineage.py index 3fd5125a93ff..eb64191a7ac0 100644 --- a/tests/integration/test_storage_iceberg_with_spark/test_row_lineage.py +++ b/tests/integration/test_storage_iceberg_with_spark/test_row_lineage.py @@ -3,6 +3,7 @@ from helpers.iceberg_utils import ( create_iceberg_table, default_upload_directory, + default_download_directory, drop_iceberg_table, get_creation_expression, get_uuid_str, @@ -10,9 +11,15 @@ def _spark_lineage(spark, table_name): - rows = spark.sql( - f"SELECT id, _row_id, _last_updated_sequence_number FROM {table_name}" - ).collect() + """Read the lineage of a table by path, which works for a table Spark never created itself. + `_row_id` and `_last_updated_sequence_number` are metadata columns: they are not part of the + schema a plain `collect` returns, so they have to be selected by name.""" + rows = ( + spark.read.format("iceberg") + .load(f"/var/lib/clickhouse/user_files/iceberg_data/default/{table_name}") + .select("id", "_row_id", "_last_updated_sequence_number") + .collect() + ) return { row["id"]: (row["_row_id"], row["_last_updated_sequence_number"]) for row in rows } @@ -49,6 +56,13 @@ def _publish(started_cluster, storage_type, table_name): ) +def _fetch(started_cluster, storage_type, table_name): + """The reverse of `_publish`: bring a table written by ClickHouse to the path Spark reads. The + download helper takes the storage path as it is, so it has to be spelled out in full.""" + path = f"/var/lib/clickhouse/user_files/iceberg_data/default/{table_name}/" + default_download_directory(started_cluster, storage_type, path, path) + + @pytest.mark.parametrize("run_on_cluster", [False, True]) @pytest.mark.parametrize("storage_type", ["s3"]) def test_row_lineage_inherited_from_manifest( @@ -276,3 +290,179 @@ def test_first_row_id_in_system_iceberg_files( ).split() == ["0", "10", "20", "30"] drop_iceberg_table(instance, TABLE_NAME) + +# The tests above have Spark write the table and ClickHouse read it. The ones below are the mirror +# image: ClickHouse writes, and Spark is the reference for what the row lineage of the result means. +INSERT_SETTINGS = {"allow_insert_into_iceberg": 1} + + +def _create_clickhouse_table(started_cluster, storage_type, table_name, schema, format_version=3, partition_by=""): + create_iceberg_table( + storage_type, + started_cluster.instances["node1"], + table_name, + started_cluster, + schema, + format_version=format_version, + partition_by=partition_by, + ) + + +@pytest.mark.parametrize("storage_type", ["s3"]) +def test_row_lineage_clickhouse(started_cluster_iceberg_with_spark, storage_type): + instance = started_cluster_iceberg_with_spark.instances["node1"] + spark = started_cluster_iceberg_with_spark.spark_session + TABLE_NAME = "test_row_lineage_written_by_clickhouse_" + storage_type + "_" + get_uuid_str() + + _create_clickhouse_table( + started_cluster_iceberg_with_spark, storage_type, TABLE_NAME, "(id Int32, s String)" + ) + + # One row per INSERT, so every row lands in its own snapshot and gets its own sequence number. + for row_key in range(40): + instance.query( + f"INSERT INTO {TABLE_NAME} VALUES ({row_key}, 'a')", settings=INSERT_SETTINGS + ) + + _fetch(started_cluster_iceberg_with_spark, storage_type, TABLE_NAME) + + spark_lineage = _spark_lineage(spark, TABLE_NAME) + + assert sorted(row_id for row_id, _ in spark_lineage.values()) == list(range(40)) + for row_key, (row_id, sequence_number) in spark_lineage.items(): + assert row_id == row_key + assert sequence_number == row_key + 1 + + assert _clickhouse_lineage(instance, TABLE_NAME) == spark_lineage + + +@pytest.mark.parametrize("storage_type", ["s3"]) +def test_row_id_clickhouse_several_files_in_one_manifest( + started_cluster_iceberg_with_spark, storage_type +): + instance = started_cluster_iceberg_with_spark.instances["node1"] + spark = started_cluster_iceberg_with_spark.spark_session + TABLE_NAME = "test_row_id_clickhouse_one_manifest_" + storage_type + "_" + get_uuid_str() + + _create_clickhouse_table( + started_cluster_iceberg_with_spark, storage_type, TABLE_NAME, "(id Int32, s String)" + ) + + # A single snapshot whose manifest lists four data files: the row ids of the second and later + # files are only right if the reader accumulates the record counts of the entries before them. + instance.query( + f"INSERT INTO {TABLE_NAME} SELECT number, 'a' FROM numbers(20)", + settings={**INSERT_SETTINGS, "iceberg_insert_max_rows_in_data_file": 5, "max_insert_threads": 1}, + ) + + _fetch(started_cluster_iceberg_with_spark, storage_type, TABLE_NAME) + + spark_lineage = _spark_lineage(spark, TABLE_NAME) + + assert sorted(row_id for row_id, _ in spark_lineage.values()) == list(range(20)) + assert all(sequence_number == 1 for _, sequence_number in spark_lineage.values()) + + assert _clickhouse_lineage(instance, TABLE_NAME) == spark_lineage + + +@pytest.mark.parametrize("storage_type", ["s3"]) +def test_row_id_clickhouse_is_not_affected_by_filter_pushdown( + started_cluster_iceberg_with_spark, storage_type +): + instance = started_cluster_iceberg_with_spark.instances["node1"] + TABLE_NAME = "test_row_id_clickhouse_pushdown_" + storage_type + "_" + get_uuid_str() + + _create_clickhouse_table( + started_cluster_iceberg_with_spark, storage_type, TABLE_NAME, "(id Int32, s String)" + ) + + instance.query( + f"INSERT INTO {TABLE_NAME} SELECT number, 'a' FROM numbers(10)", settings=INSERT_SETTINGS + ) + instance.query( + f"INSERT INTO {TABLE_NAME} SELECT number, 'b' FROM numbers(10, 10)", settings=INSERT_SETTINGS + ) + + # A row id is the position of the row in the table, not in whatever part of the file survived + # the filter, so it must not move when rows or row groups are skipped. + assert _row_ids(_clickhouse_lineage(instance, TABLE_NAME, where="WHERE id >= 15")) == { + row_key: row_key for row_key in range(15, 20) + } + + assert _row_ids(_clickhouse_lineage(instance, TABLE_NAME, where="WHERE id % 7 = 3")) == { + 3: 3, + 10: 10, + 17: 17, + } + + +@pytest.mark.parametrize("storage_type", ["s3"]) +def test_row_id_clickhouse_survives_delete(started_cluster_iceberg_with_spark, storage_type): + instance = started_cluster_iceberg_with_spark.instances["node1"] + TABLE_NAME = "test_row_id_clickhouse_after_delete_" + storage_type + "_" + get_uuid_str() + + _create_clickhouse_table( + started_cluster_iceberg_with_spark, storage_type, TABLE_NAME, "(id Int32, s String)" + ) + + instance.query( + f"INSERT INTO {TABLE_NAME} SELECT number, 'a' FROM numbers(4)", settings=INSERT_SETTINGS + ) + instance.query(f"ALTER TABLE {TABLE_NAME} DELETE WHERE id = 1", settings=INSERT_SETTINGS) + + assert int(instance.query(f"SELECT count() FROM {TABLE_NAME}")) == 3 + + # Deleting a row does not renumber the rows that stay. + assert _row_ids(_clickhouse_lineage(instance, TABLE_NAME)) == {0: 0, 2: 2, 3: 3} + + +@pytest.mark.parametrize("storage_type", ["s3"]) +def test_row_lineage_clickhouse_is_null_for_v2_table( + started_cluster_iceberg_with_spark, storage_type +): + instance = started_cluster_iceberg_with_spark.instances["node1"] + TABLE_NAME = "test_row_lineage_clickhouse_v2_null_" + storage_type + "_" + get_uuid_str() + + _create_clickhouse_table( + started_cluster_iceberg_with_spark, + storage_type, + TABLE_NAME, + "(id Int32, s String)", + format_version=2, + ) + + instance.query( + f"INSERT INTO {TABLE_NAME} SELECT number, 'a' FROM numbers(4)", settings=INSERT_SETTINGS + ) + + assert _clickhouse_lineage(instance, TABLE_NAME) == { + row_key: (None, None) for row_key in range(4) + } + + +@pytest.mark.parametrize("storage_type", ["s3"]) +def test_first_row_id_in_system_iceberg_files_clickhouse( + started_cluster_iceberg_with_spark, storage_type +): + instance = started_cluster_iceberg_with_spark.instances["node1"] + TABLE_NAME = "test_first_row_id_system_table_clickhouse_" + storage_type + "_" + get_uuid_str() + + _create_clickhouse_table( + started_cluster_iceberg_with_spark, storage_type, TABLE_NAME, "(id Int32, s String)" + ) + + for lo in range(0, 40, 10): + instance.query( + f"INSERT INTO {TABLE_NAME} SELECT number, 'a' FROM numbers({lo}, 10)", + settings=INSERT_SETTINGS, + ) + + # ClickHouse writes the manifest entries without an explicit first_row_id as well, so the values + # here are the ones resolved from the first_row_id of the manifest list entry. + assert instance.query( + f"SELECT first_row_id FROM system.iceberg_files " + f"WHERE database = currentDatabase() AND table = '{TABLE_NAME}' AND content = 'DATA' " + f"ORDER BY first_row_id FORMAT TSV" + ).split() == ["0", "10", "20", "30"] + + drop_iceberg_table(instance, TABLE_NAME) diff --git a/tests/integration/test_storage_iceberg_with_spark/test_row_lineage_pruning.py b/tests/integration/test_storage_iceberg_with_spark/test_row_lineage_pruning.py index 8ac5e032b268..8981a7f38830 100644 --- a/tests/integration/test_storage_iceberg_with_spark/test_row_lineage_pruning.py +++ b/tests/integration/test_storage_iceberg_with_spark/test_row_lineage_pruning.py @@ -20,6 +20,7 @@ from helpers.iceberg_utils import ( check_validity_and_get_prunned_files_general, + create_iceberg_table, default_upload_directory, get_creation_expression, get_uuid_str, @@ -350,3 +351,163 @@ def test_materialized_row_ids_are_not_pruned_away( ).strip() == "0\n2\n3" ) + + +# The tests above read what Spark wrote. The ones below prune over metadata ClickHouse wrote itself: +# the inherited row id range of a manifest entry is only as good as the `first_row_id` the writer put +# into the manifest list, so the same filters are replayed against a table of its own making. +INSERT_SETTINGS = {"allow_insert_into_iceberg": 1} + + +def _clickhouse_table_with_five_appends( + started_cluster, storage_type, table_name, schema="(id Int32, s String)", payload="'a'" +): + """Five inserts of ten rows: file k holds row ids [10k, 10k + 10) and sequence number k + 1.""" + instance = started_cluster.instances["node1"] + create_iceberg_table( + storage_type, + instance, + table_name, + started_cluster, + schema, + format_version=3, + ) + for lo in range(0, 50, 10): + instance.query( + f"INSERT INTO {table_name} SELECT number, {payload} FROM numbers({lo}, 10)", + settings=INSERT_SETTINGS, + ) + + +@pytest.mark.parametrize("storage_type", ["s3"]) +def test_row_id_filter_prunes_files_clickhouse(started_cluster_iceberg_with_spark, storage_type): + instance = started_cluster_iceberg_with_spark.instances["node1"] + TABLE_NAME = "test_row_id_pruning_clickhouse_" + storage_type + "_" + get_uuid_str() + + _clickhouse_table_with_five_appends( + started_cluster_iceberg_with_spark, storage_type, TABLE_NAME + ) + + assert _pruned_files(instance, TABLE_NAME, f"SELECT id FROM {TABLE_NAME} ORDER BY ALL") == 0 + + # A point lookup touches the one file whose row id range contains the value. + assert ( + _pruned_files( + instance, TABLE_NAME, f"SELECT id FROM {TABLE_NAME} WHERE _row_id = 25 ORDER BY ALL" + ) + == 4 + ) + + # A half-open range keeps the two files above it. + assert ( + _pruned_files( + instance, TABLE_NAME, f"SELECT id FROM {TABLE_NAME} WHERE _row_id >= 35 ORDER BY ALL" + ) + == 3 + ) + + assert ( + _pruned_files( + instance, TABLE_NAME, f"SELECT id FROM {TABLE_NAME} WHERE _row_id < 10 ORDER BY ALL" + ) + == 4 + ) + + # A range spanning everything prunes nothing. + assert ( + _pruned_files( + instance, TABLE_NAME, f"SELECT id FROM {TABLE_NAME} WHERE _row_id >= 0 ORDER BY ALL" + ) + == 0 + ) + + +@pytest.mark.parametrize("storage_type", ["s3"]) +def test_incremental_read_by_sequence_number_prunes_files_clickhouse( + started_cluster_iceberg_with_spark, storage_type +): + instance = started_cluster_iceberg_with_spark.instances["node1"] + TABLE_NAME = "test_sequence_number_pruning_clickhouse_" + storage_type + "_" + get_uuid_str() + + _clickhouse_table_with_five_appends( + started_cluster_iceberg_with_spark, storage_type, TABLE_NAME + ) + + assert ( + _pruned_files( + instance, + TABLE_NAME, + f"SELECT id FROM {TABLE_NAME} WHERE _last_updated_sequence_number > 3 ORDER BY ALL", + ) + == 3 + ) + + assert ( + _pruned_files( + instance, + TABLE_NAME, + f"SELECT id FROM {TABLE_NAME} WHERE _last_updated_sequence_number = 2 ORDER BY ALL", + ) + == 4 + ) + + assert ( + _pruned_files( + instance, + TABLE_NAME, + f"SELECT id FROM {TABLE_NAME} WHERE _last_updated_sequence_number > 0 ORDER BY ALL", + ) + == 0 + ) + + # An incremental consumer reads exactly the rows of the two newest files, not the whole table. + assert ( + _read_rows( + instance, + f"SELECT id FROM {TABLE_NAME} WHERE _last_updated_sequence_number > 3 FORMAT Null", + PRUNING_ENABLED, + ) + == 20 + ) + + +@pytest.mark.parametrize("storage_type", ["s3"]) +def test_row_id_pruning_is_skipped_for_orc_clickhouse( + started_cluster_iceberg_with_spark, storage_type +): + instance = started_cluster_iceberg_with_spark.instances["node1"] + TABLE_NAME = "test_row_lineage_pruning_orc_clickhouse_" + storage_type + "_" + get_uuid_str() + + create_iceberg_table( + storage_type, + instance, + TABLE_NAME, + started_cluster_iceberg_with_spark, + "(id Int32, s String)", + format_version=3, + format="ORC", + ) + for lo in range(0, 50, 10): + instance.query( + f"INSERT INTO {TABLE_NAME} SELECT number, 'a' FROM numbers({lo}, 10)", + settings=INSERT_SETTINGS, + ) + + # The ORC reader reports no physical row numbers, so `_row_id` is NULL for every row and the + # inherited range describes nothing: pruning by it would drop the rows this query asks for. + assert ( + instance.query( + f"SELECT count() FROM {TABLE_NAME} WHERE _row_id IS NULL", settings=PRUNING_ENABLED + ).strip() + == "50" + ) + + # The sequence number does not depend on row numbers, so it prunes as usual. + assert ( + _pruned_files( + instance, + TABLE_NAME, + f"SELECT id FROM {TABLE_NAME} WHERE _last_updated_sequence_number > 3 ORDER BY ALL", + ) + == 3 + ) From cfce020a19888991b4afc4ec55334dcb60d34326 Mon Sep 17 00:00:00 2001 From: Kanthi Subramanian Date: Thu, 1 Oct 2026 16:27:10 +0200 Subject: [PATCH 2/3] Fix tests --- .../DataLakes/Iceberg/Constant.h | 2 ++ .../DataLakes/Iceberg/IcebergWrites.cpp | 31 +++++++++++++++++-- 2 files changed, 30 insertions(+), 3 deletions(-) diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/Constant.h b/src/Storages/ObjectStorage/DataLakes/Iceberg/Constant.h index 2e72d799606b..be532ba5b89c 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/Constant.h +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/Constant.h @@ -33,6 +33,8 @@ DEFINE_ICEBERG_FIELD(name); DEFINE_ICEBERG_FIELD(required); DEFINE_ICEBERG_FIELD(schema); DEFINE_ICEBERG_FIELD(schemas); +/// The Avro file-header metadata key that stores the writer schema (Avro spec, "avro.schema"). +DEFINE_ICEBERG_FIELD_ALIAS(avro_schema, avro.schema); DEFINE_ICEBERG_FIELD(sequence_number); DEFINE_ICEBERG_FIELD(snapshots); DEFINE_ICEBERG_FIELD(status); diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergWrites.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergWrites.cpp index 6d94f80e33f0..dc1faa34ebbd 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergWrites.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergWrites.cpp @@ -557,8 +557,25 @@ IcebergSerializedFileStats serializeDataFileStats( static void extendSchemaForPartitions( String & schema, const std::vector & partition_columns, - const std::vector & partition_types) + const std::vector & partition_types, + const Poco::JSON::Array::Ptr & partition_spec_fields) { + /// Iceberg projects the manifest `partition` values onto the spec by field-id, so reuse + /// the persisted spec ids. A legacy v1 spec may omit them; there the spec assigns partition + /// field-ids sequentially from 1000, so fall back to that default when any id is missing. + bool spec_has_all_ids = partition_spec_fields && partition_spec_fields->size() == partition_columns.size(); + if (spec_has_all_ids) + { + for (size_t i = 0; i < partition_columns.size(); ++i) + { + if (!partition_spec_fields->getObject(static_cast(i))->has(Iceberg::f_field_id)) + { + spec_has_all_ids = false; + break; + } + } + } + Poco::JSON::Array::Ptr partition_fields = new Poco::JSON::Array; /// Types that need a `logicalType` annotation (Time/Time64 -> "time-micros") are @@ -580,7 +597,9 @@ static void extendSchemaForPartitions( for (size_t i = 0; i < partition_columns.size(); ++i) { - const Int32 field_id = static_cast(1000 + i); + const Int32 field_id = spec_has_all_ids + ? partition_spec_fields->getObject(static_cast(i))->getValue(Iceberg::f_field_id) + : static_cast(1000 + i); Poco::JSON::Object::Ptr field = new Poco::JSON::Object; field->set(Iceberg::f_field_id, field_id); @@ -694,7 +713,7 @@ void generateManifestFile( else throw Exception(ErrorCodes::BAD_ARGUMENTS, "Unsupported iceberg format-version {}", version); - extendSchemaForPartitions(schema_representation, partition_columns, partition_types); + extendSchemaForPartitions(schema_representation, partition_columns, partition_types, partition_spec->getArray(Iceberg::f_fields)); auto schema = avro::compileJsonSchemaFromString(schema_representation); const avro::NodePtr & root_schema = schema.root(); // NOLINT @@ -707,6 +726,9 @@ void generateManifestFile( auto adapter = std::make_unique(buf); avro::DataFileWriter writer(std::move(adapter), schema); + /// avro-cpp's compiled schema loses the Iceberg field-id/element-id attributes; write the + /// original id-carrying JSON as the avro.schema header so external readers can plan a scan. + writer.setMetadata(Iceberg::f_avro_schema, schema_representation); writer.setMetadata(Iceberg::f_schema, json_representation); writer.setMetadata(Iceberg::f_format_version, std::to_string(version)); @@ -1111,6 +1133,9 @@ void generateManifestList( auto adapter = std::make_unique(buf); avro::DataFileWriter writer(std::move(adapter), schema); + /// See generateManifestFile: avro-cpp drops the Iceberg field-id/element-id attributes, so + /// write the original id-carrying JSON as the avro.schema header for external readers. + writer.setMetadata(Iceberg::f_avro_schema, schema_representation); writer.setMetadata(Iceberg::f_format_version, std::to_string(version)); Int64 next_first_row_id = 0; From 6805979dd278b1dc0d3ef33d0269233f948a1387 Mon Sep 17 00:00:00 2001 From: Kanthi Subramanian Date: Thu, 1 Oct 2026 17:53:45 +0200 Subject: [PATCH 3/3] Fix integration test --- .../test_row_lineage.py | 16 ++++++++++------ 1 file changed, 10 insertions(+), 6 deletions(-) diff --git a/tests/integration/test_storage_iceberg_with_spark/test_row_lineage.py b/tests/integration/test_storage_iceberg_with_spark/test_row_lineage.py index eb64191a7ac0..e4ff2148b99c 100644 --- a/tests/integration/test_storage_iceberg_with_spark/test_row_lineage.py +++ b/tests/integration/test_storage_iceberg_with_spark/test_row_lineage.py @@ -397,9 +397,9 @@ def test_row_id_clickhouse_is_not_affected_by_filter_pushdown( @pytest.mark.parametrize("storage_type", ["s3"]) -def test_row_id_clickhouse_survives_delete(started_cluster_iceberg_with_spark, storage_type): +def test_row_id_clickhouse_v3_delete_rejected(started_cluster_iceberg_with_spark, storage_type): instance = started_cluster_iceberg_with_spark.instances["node1"] - TABLE_NAME = "test_row_id_clickhouse_after_delete_" + storage_type + "_" + get_uuid_str() + TABLE_NAME = "test_row_id_clickhouse_v3_delete_rejected_" + storage_type + "_" + get_uuid_str() _create_clickhouse_table( started_cluster_iceberg_with_spark, storage_type, TABLE_NAME, "(id Int32, s String)" @@ -408,12 +408,16 @@ def test_row_id_clickhouse_survives_delete(started_cluster_iceberg_with_spark, s instance.query( f"INSERT INTO {TABLE_NAME} SELECT number, 'a' FROM numbers(4)", settings=INSERT_SETTINGS ) - instance.query(f"ALTER TABLE {TABLE_NAME} DELETE WHERE id = 1", settings=INSERT_SETTINGS) - assert int(instance.query(f"SELECT count() FROM {TABLE_NAME}")) == 3 + # Iceberg v3 forbids new position-delete files and ClickHouse cannot write deletion vectors yet, + # so the mutation must fail before writing anything and leave the rows and their ids as they were. + error = instance.query_and_get_error( + f"ALTER TABLE {TABLE_NAME} DELETE WHERE id = 1", settings=INSERT_SETTINGS + ) + assert "SUPPORT_IS_DISABLED" in error - # Deleting a row does not renumber the rows that stay. - assert _row_ids(_clickhouse_lineage(instance, TABLE_NAME)) == {0: 0, 2: 2, 3: 3} + assert int(instance.query(f"SELECT count() FROM {TABLE_NAME}")) == 4 + assert _row_ids(_clickhouse_lineage(instance, TABLE_NAME)) == {k: k for k in range(4)} @pytest.mark.parametrize("storage_type", ["s3"])