From 2b3a3a5120444a287e5a4c23ed2fb56d0a322b93 Mon Sep 17 00:00:00 2001 From: Konstantin Morozov Date: Fri, 17 Jul 2026 17:21:23 +0200 Subject: [PATCH 1/7] update prewhere exclude columns list Signed-off-by: Konstantin Morozov --- .../DataLakes/DataLakeConfiguration.h | 11 +++++ .../DataLakes/IDataLakeMetadata.h | 2 + .../DataLakes/Iceberg/IcebergMetadata.cpp | 15 ++++++ .../DataLakes/Iceberg/IcebergMetadata.h | 2 + .../ObjectStorage/DataLakes/Iceberg/Utils.cpp | 46 +++++++++++++++++++ .../ObjectStorage/DataLakes/Iceberg/Utils.h | 1 + .../ObjectStorage/StorageObjectStorage.cpp | 25 +++++++++- .../ObjectStorage/StorageObjectStorage.h | 4 ++ .../StorageObjectStorageCluster.cpp | 3 ++ .../StorageObjectStorageConfiguration.h | 1 + 10 files changed, 109 insertions(+), 1 deletion(-) diff --git a/src/Storages/ObjectStorage/DataLakes/DataLakeConfiguration.h b/src/Storages/ObjectStorage/DataLakes/DataLakeConfiguration.h index 74d7ceea7f93..8d043a8af349 100644 --- a/src/Storages/ObjectStorage/DataLakes/DataLakeConfiguration.h +++ b/src/Storages/ObjectStorage/DataLakes/DataLakeConfiguration.h @@ -323,6 +323,15 @@ class DataLakeConfiguration : public BaseStorageConfiguration, public std::enabl return current_metadata->getColumnMapperForCurrentSchema(storage_metadata_snapshot, context); } + Names getIdentityPartitionColumnNames(ContextPtr context) const override + { + if (!current_metadata) + { + return {}; + } + return current_metadata->getIdentityPartitionColumnNames(context); + } + void drop(ContextPtr local_context) override { if (current_metadata) @@ -806,6 +815,8 @@ class StorageIcebergConfiguration : public StorageObjectStorageConfiguration, pu ColumnMapperPtr getColumnMapperForCurrentSchema(StorageMetadataPtr storage_metadata_snapshot, ContextPtr context) const override { return getImpl().getColumnMapperForCurrentSchema(storage_metadata_snapshot, context); } + Names getIdentityPartitionColumnNames(ContextPtr context) const override { return getImpl().getIdentityPartitionColumnNames(context); } + std::shared_ptr getCatalog(ContextPtr context, const StorageID & table_id) const override { return getImpl().getCatalog(context, table_id); } diff --git a/src/Storages/ObjectStorage/DataLakes/IDataLakeMetadata.h b/src/Storages/ObjectStorage/DataLakes/IDataLakeMetadata.h index 7fad60407ed5..7a76d89ea410 100644 --- a/src/Storages/ObjectStorage/DataLakes/IDataLakeMetadata.h +++ b/src/Storages/ObjectStorage/DataLakes/IDataLakeMetadata.h @@ -193,6 +193,8 @@ class IDataLakeMetadata : boost::noncopyable virtual ColumnMapperPtr getColumnMapperForObject(ObjectInfoPtr /**/) const { return nullptr; } virtual ColumnMapperPtr getColumnMapperForCurrentSchema(StorageMetadataPtr, ContextPtr) const { return nullptr; } + virtual Names getIdentityPartitionColumnNames(ContextPtr) const { return {}; } + virtual SinkToStoragePtr write( SharedHeader /*sample_block*/, const StorageID & /*table_id*/, diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.cpp index 5fc41baf7ee2..f742f98b6270 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.cpp @@ -1436,6 +1436,21 @@ ColumnMapperPtr IcebergMetadata::getColumnMapperForCurrentSchema(StorageMetadata return persistent_components.schema_processor->getColumnMapperById(iceberg_table_state->schema_id); } +Names IcebergMetadata::getIdentityPartitionColumnNames(ContextPtr local_context) const +{ + // check getPartitionKey + auto [data_snapshot, table_state_snapshot] = getRelevantState(local_context); + auto metadata_object = getMetadataJSONObject( + table_state_snapshot.metadata_file_path, + object_storage, + persistent_components.metadata_cache, + local_context, + log, + persistent_components.metadata_compression_method, + persistent_components.table_uuid); + return getIdentityPartitionColumnsFromMetadata(metadata_object); +} + std::optional IcebergMetadata::getPartitionKey(ContextPtr local_context, TableStateSnapshot actual_table_state_snapshot) const { auto metadata_object = getMetadataJSONObject( diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.h b/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.h index 5b2c64c07039..9f30a34fdb14 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.h +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.h @@ -106,6 +106,8 @@ class IcebergMetadata : public IDataLakeMetadata ColumnMapperPtr getColumnMapperForCurrentSchema(StorageMetadataPtr storage_metadata_snapshot, ContextPtr context) const override; + Names getIdentityPartitionColumnNames(ContextPtr context) const override; + SinkToStoragePtr write( SharedHeader sample_block, const StorageID & table_id, diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.cpp index 8564eebda3ef..f43d0d499d6b 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.cpp @@ -1471,6 +1471,52 @@ std::optional getPartitionKeyStringFromMetadata(Poco::JSON::Object::Ptr return result; } +Names getIdentityPartitionColumnsFromMetadata(Poco::JSON::Object::Ptr metadata_object) +{ + // @todo refactor dupli + if (!metadata_object->has(f_partition_specs) || !metadata_object->has(f_default_spec_id)) + return {}; + + auto partition_spec_id = metadata_object->getValue(f_default_spec_id); + Poco::JSON::Array::Ptr partition_specs = metadata_object->getArray(f_partition_specs); + std::unordered_map source_id_to_column_name; + auto [schema, current_schema_id] = parseTableSchemaV2Method(metadata_object); + auto mapper = createColumnMapper(schema)->getStorageColumnEncoding(); + for (const auto & [col_name, source_id] : mapper) + source_id_to_column_name[source_id] = col_name; + + Poco::JSON::Object::Ptr partition_spec; + for (size_t i = 0; i < partition_specs->size(); ++i) + { + auto spec = partition_specs->getObject(static_cast(i)); + if (spec->getValue(f_spec_id) == partition_spec_id) + { + partition_spec = spec; + break; + } + } + + if (!partition_spec || !partition_spec->has(f_fields)) + return {}; + + auto fields = partition_spec->getArray(f_fields); + if (fields->size() == 0) + return {}; + + Names result; + std::vector part_exprs; + for (UInt32 i = 0; i < fields->size(); ++i) + { + auto field = fields->getObject(i); + if (field->getValue(f_transform) != "identity") + continue; + + if (auto it = source_id_to_column_name.find(field->getValue(f_source_id)); it != source_id_to_column_name.end()) + result.push_back(it->second); + } + return result; +} + std::optional getSortingKeyDisplayStringFromMetadata(Poco::JSON::Object::Ptr metadata_object, const NamesAndTypesList & /* ch_schema */) { if (!metadata_object->has(f_sort_orders) || !metadata_object->has(f_default_sort_order_id)) diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.h b/src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.h index 2720e51547e0..835937fb5d35 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.h +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.h @@ -133,6 +133,7 @@ std::optional getSortingKeyDisplayStringFromMetadata( Poco::JSON::Object::Ptr metadata_object, const NamesAndTypesList & ch_schema); std::optional getPartitionKeyStringFromMetadata( Poco::JSON::Object::Ptr metadata_object, const NamesAndTypesList & ch_schema, ContextPtr local_context); +Names getIdentityPartitionColumnsFromMetadata(Poco::JSON::Object::Ptr metadata_object); void sortBlockByKeyDescription(Block & block, const KeyDescription & sort_description, ContextPtr context); } diff --git a/src/Storages/ObjectStorage/StorageObjectStorage.cpp b/src/Storages/ObjectStorage/StorageObjectStorage.cpp index de5e96577812..36905fd433f9 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorage.cpp +++ b/src/Storages/ObjectStorage/StorageObjectStorage.cpp @@ -3,6 +3,7 @@ #include #include +#include #include #include #include @@ -336,7 +337,21 @@ bool StorageObjectStorage::canMoveConditionsToPrewhere() const std::optional StorageObjectStorage::supportedPrewhereColumns() const { - return getInMemoryMetadataPtr()->getColumnsWithoutDefaultExpressions(/*exclude=*/ hive_partition_columns_to_read_from_file_path); + auto exclude = hive_partition_columns_to_read_from_file_path; + + const auto & cols = getInMemoryMetadataPtr()->getColumns(); + { + SharedLockGuard lock(mutex_file_constant_columns); + for (const auto & col : file_constant_columns) + { + if (getInMemoryMetadataPtr()->getColumns().has(col)) + // tryGetColumn + exclude.emplace_back(col, cols.get(col).type); + } + } + + LOG_DEBUG(log, "Prewhere exclude list: [{}]", exclude.toString()); + return getInMemoryMetadataPtr()->getColumnsWithoutDefaultExpressions(/*exclude=*/exclude); } IStorage::ColumnSizeByName StorageObjectStorage::getColumnSizes() const @@ -354,6 +369,12 @@ configuration->update(object_storage, query_context); return configuration->getExternalMetadata(); } +void StorageObjectStorage::updateFileConstantColumns(ContextPtr query_context) +{ + UniqueLock lock(mutex_file_constant_columns); + file_constant_columns = configuration->getIdentityPartitionColumnNames(query_context); +} + void StorageObjectStorage::updateExternalDynamicMetadataIfExists(ContextPtr query_context) { if (!configuration->isDataLakeConfiguration()) @@ -381,6 +402,8 @@ void StorageObjectStorage::updateExternalDynamicMetadataIfExists(ContextPtr quer new_metadata = *metadata_snapshot; } + updateFileConstantColumns(query_context); + setInMemoryMetadata(new_metadata); } diff --git a/src/Storages/ObjectStorage/StorageObjectStorage.h b/src/Storages/ObjectStorage/StorageObjectStorage.h index 6ace629b7085..183b8f6869db 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorage.h +++ b/src/Storages/ObjectStorage/StorageObjectStorage.h @@ -155,6 +155,7 @@ class StorageObjectStorage : public IStorage void addInferredEngineArgsToCreateQuery(ASTs & args, const ContextPtr & context) const override; + void updateFileConstantColumns(ContextPtr query_context); void updateExternalDynamicMetadataIfExists(ContextPtr query_context) override; IDataLakeMetadata * getExternalMetadata(ContextPtr query_context); @@ -214,6 +215,9 @@ class StorageObjectStorage : public IStorage NamesAndTypesList hive_partition_columns_to_read_from_file_path; NamesAndTypesList file_columns; + mutable SharedMutex mutex_file_constant_columns; + Names file_constant_columns TSA_GUARDED_BY(mutex_file_constant_columns); + LoggerPtr log; std::shared_ptr catalog; diff --git a/src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp b/src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp index bfb107c14d21..438fb2432aa8 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp +++ b/src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp @@ -569,7 +569,10 @@ void StorageObjectStorageCluster::updateExternalDynamicMetadataIfExists(ContextP setInMemoryMetadata(new_metadata); if (pure_storage) + { pure_storage->setInMemoryMetadata(IStorageCluster::getInMemoryMetadata()); + pure_storage->updateFileConstantColumns(query_context); + } } class TaskDistributor : public TaskIterator diff --git a/src/Storages/ObjectStorage/StorageObjectStorageConfiguration.h b/src/Storages/ObjectStorage/StorageObjectStorageConfiguration.h index 0ce750c6193a..112e7e949949 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorageConfiguration.h +++ b/src/Storages/ObjectStorage/StorageObjectStorageConfiguration.h @@ -312,6 +312,7 @@ class StorageObjectStorageConfiguration virtual ColumnMapperPtr getColumnMapperForCurrentSchema(StorageMetadataPtr /**/, ContextPtr /**/) const { return nullptr; } + virtual Names getIdentityPartitionColumnNames(ContextPtr) const { return {}; } virtual std::shared_ptr getCatalog(ContextPtr /*context*/, const StorageID & /*table_id*/) const { From b48c21860991d288ec21e8c8c8bb5e9326959d59 Mon Sep 17 00:00:00 2001 From: Konstantin Morozov Date: Tue, 21 Jul 2026 15:26:25 +0200 Subject: [PATCH 2/7] move ident columns to metadata Signed-off-by: Konstantin Morozov --- .../DataLakes/DataLakeConfiguration.h | 6 ++-- .../DataLakes/IDataLakeMetadata.h | 2 +- .../DataLakes/Iceberg/IcebergMetadata.cpp | 9 ++--- .../DataLakes/Iceberg/IcebergMetadata.h | 2 +- .../ObjectStorage/DataLakes/Iceberg/Utils.cpp | 3 +- .../ObjectStorage/StorageObjectStorage.cpp | 35 +++++++++++-------- .../ObjectStorage/StorageObjectStorage.h | 15 +++++--- .../StorageObjectStorageCluster.cpp | 9 +++-- .../StorageObjectStorageConfiguration.h | 2 +- .../StorageObjectStorageSource.cpp | 2 +- src/Storages/StorageInMemoryMetadata.cpp | 7 ++++ src/Storages/StorageInMemoryMetadata.h | 5 +++ 12 files changed, 64 insertions(+), 33 deletions(-) diff --git a/src/Storages/ObjectStorage/DataLakes/DataLakeConfiguration.h b/src/Storages/ObjectStorage/DataLakes/DataLakeConfiguration.h index 8d043a8af349..1ee40645ae24 100644 --- a/src/Storages/ObjectStorage/DataLakes/DataLakeConfiguration.h +++ b/src/Storages/ObjectStorage/DataLakes/DataLakeConfiguration.h @@ -323,13 +323,13 @@ class DataLakeConfiguration : public BaseStorageConfiguration, public std::enabl return current_metadata->getColumnMapperForCurrentSchema(storage_metadata_snapshot, context); } - Names getIdentityPartitionColumnNames(ContextPtr context) const override + Names getIdentityPartitionColumnNames(const DataLakeTableStateSnapshot & state, ContextPtr context) const override { if (!current_metadata) { return {}; } - return current_metadata->getIdentityPartitionColumnNames(context); + return current_metadata->getIdentityPartitionColumnNames(state, context); } void drop(ContextPtr local_context) override @@ -815,7 +815,7 @@ class StorageIcebergConfiguration : public StorageObjectStorageConfiguration, pu ColumnMapperPtr getColumnMapperForCurrentSchema(StorageMetadataPtr storage_metadata_snapshot, ContextPtr context) const override { return getImpl().getColumnMapperForCurrentSchema(storage_metadata_snapshot, context); } - Names getIdentityPartitionColumnNames(ContextPtr context) const override { return getImpl().getIdentityPartitionColumnNames(context); } + Names getIdentityPartitionColumnNames(const DataLakeTableStateSnapshot & state, ContextPtr context) const override { return getImpl().getIdentityPartitionColumnNames(state, context); } std::shared_ptr getCatalog(ContextPtr context, const StorageID & table_id) const override { return getImpl().getCatalog(context, table_id); } diff --git a/src/Storages/ObjectStorage/DataLakes/IDataLakeMetadata.h b/src/Storages/ObjectStorage/DataLakes/IDataLakeMetadata.h index 7a76d89ea410..e64db238a62b 100644 --- a/src/Storages/ObjectStorage/DataLakes/IDataLakeMetadata.h +++ b/src/Storages/ObjectStorage/DataLakes/IDataLakeMetadata.h @@ -193,7 +193,7 @@ class IDataLakeMetadata : boost::noncopyable virtual ColumnMapperPtr getColumnMapperForObject(ObjectInfoPtr /**/) const { return nullptr; } virtual ColumnMapperPtr getColumnMapperForCurrentSchema(StorageMetadataPtr, ContextPtr) const { return nullptr; } - virtual Names getIdentityPartitionColumnNames(ContextPtr) const { return {}; } + virtual Names getIdentityPartitionColumnNames(const DataLakeTableStateSnapshot &, ContextPtr) const { return {}; } virtual SinkToStoragePtr write( SharedHeader /*sample_block*/, diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.cpp index f742f98b6270..661c45606f6f 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.cpp @@ -1436,12 +1436,13 @@ ColumnMapperPtr IcebergMetadata::getColumnMapperForCurrentSchema(StorageMetadata return persistent_components.schema_processor->getColumnMapperById(iceberg_table_state->schema_id); } -Names IcebergMetadata::getIdentityPartitionColumnNames(ContextPtr local_context) const +Names IcebergMetadata::getIdentityPartitionColumnNames(const DataLakeTableStateSnapshot & state, ContextPtr local_context) const { - // check getPartitionKey - auto [data_snapshot, table_state_snapshot] = getRelevantState(local_context); + // @todo + chassert(std::holds_alternative(state)); + const auto & iceberg_state = std::get(state); auto metadata_object = getMetadataJSONObject( - table_state_snapshot.metadata_file_path, + iceberg_state.metadata_file_path, object_storage, persistent_components.metadata_cache, local_context, diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.h b/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.h index 9f30a34fdb14..13f3849d35a2 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.h +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.h @@ -106,7 +106,7 @@ class IcebergMetadata : public IDataLakeMetadata ColumnMapperPtr getColumnMapperForCurrentSchema(StorageMetadataPtr storage_metadata_snapshot, ContextPtr context) const override; - Names getIdentityPartitionColumnNames(ContextPtr context) const override; + Names getIdentityPartitionColumnNames(const DataLakeTableStateSnapshot & state, ContextPtr context) const override; SinkToStoragePtr write( SharedHeader sample_block, diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.cpp index f43d0d499d6b..997816680932 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.cpp @@ -1508,7 +1508,8 @@ Names getIdentityPartitionColumnsFromMetadata(Poco::JSON::Object::Ptr metadata_o for (UInt32 i = 0; i < fields->size(); ++i) { auto field = fields->getObject(i); - if (field->getValue(f_transform) != "identity") + + if (Poco::toLower(field->getValue(f_transform)) != "identity") continue; if (auto it = source_id_to_column_name.find(field->getValue(f_source_id)); it != source_id_to_column_name.end()) diff --git a/src/Storages/ObjectStorage/StorageObjectStorage.cpp b/src/Storages/ObjectStorage/StorageObjectStorage.cpp index 36905fd433f9..c94f846ae4a2 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorage.cpp +++ b/src/Storages/ObjectStorage/StorageObjectStorage.cpp @@ -3,7 +3,6 @@ #include #include -#include #include #include #include @@ -52,6 +51,7 @@ namespace Setting extern const SettingsInt64 delta_lake_snapshot_start_version; extern const SettingsInt64 delta_lake_snapshot_end_version; extern const SettingsUInt64 max_streams_for_files_processing_in_cluster_functions; + extern const SettingsBool allow_experimental_iceberg_read_optimization; } namespace ErrorCodes @@ -285,6 +285,12 @@ StorageObjectStorage::StorageObjectStorage( if (auto metadata_snapshot = configuration->buildStorageMetadataFromState(*state, context)) metadata = *metadata_snapshot; } + + /// Table functions bypass updateExternalDynamicMetadataIfExists (see above), + /// so identity-partition columns must be resolved here as well -- otherwise + /// PREWHERE could be built over a column that gets erased from the read set + /// as file-constant (see StorageObjectStorageSource's constant-column handling). + updateIdentityPartitionColumns(metadata, configuration, *state, context); } } @@ -339,19 +345,15 @@ std::optional StorageObjectStorage::supportedPrewhereColumns() const { auto exclude = hive_partition_columns_to_read_from_file_path; - const auto & cols = getInMemoryMetadataPtr()->getColumns(); + const auto metadata = getInMemoryMetadataPtr(); + for (const auto & identity_partition_column : metadata->identity_partition_columns) { - SharedLockGuard lock(mutex_file_constant_columns); - for (const auto & col : file_constant_columns) - { - if (getInMemoryMetadataPtr()->getColumns().has(col)) - // tryGetColumn - exclude.emplace_back(col, cols.get(col).type); - } + if (metadata->getColumns().has(identity_partition_column)) + exclude.emplace_back(identity_partition_column, metadata->getColumns().get(identity_partition_column).type); } LOG_DEBUG(log, "Prewhere exclude list: [{}]", exclude.toString()); - return getInMemoryMetadataPtr()->getColumnsWithoutDefaultExpressions(/*exclude=*/exclude); + return metadata->getColumnsWithoutDefaultExpressions(/*exclude=*/exclude); } IStorage::ColumnSizeByName StorageObjectStorage::getColumnSizes() const @@ -369,10 +371,14 @@ configuration->update(object_storage, query_context); return configuration->getExternalMetadata(); } -void StorageObjectStorage::updateFileConstantColumns(ContextPtr query_context) +void StorageObjectStorage::updateIdentityPartitionColumns( + StorageInMemoryMetadata & metadata, + const StorageObjectStorageConfigurationPtr & configuration, + const DataLakeTableStateSnapshot & state, + ContextPtr context) { - UniqueLock lock(mutex_file_constant_columns); - file_constant_columns = configuration->getIdentityPartitionColumnNames(query_context); + if (context->getSettingsRef()[Setting::allow_experimental_iceberg_read_optimization]) + metadata.setIdentityPartitionColumns(configuration->getIdentityPartitionColumnNames(state, context)); } void StorageObjectStorage::updateExternalDynamicMetadataIfExists(ContextPtr query_context) @@ -402,7 +408,8 @@ void StorageObjectStorage::updateExternalDynamicMetadataIfExists(ContextPtr quer new_metadata = *metadata_snapshot; } - updateFileConstantColumns(query_context); + /// Resolved from the same pinned `state` above -- no second state resolution. + updateIdentityPartitionColumns(new_metadata, configuration, *state, query_context); setInMemoryMetadata(new_metadata); } diff --git a/src/Storages/ObjectStorage/StorageObjectStorage.h b/src/Storages/ObjectStorage/StorageObjectStorage.h index 183b8f6869db..5208536f48e9 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorage.h +++ b/src/Storages/ObjectStorage/StorageObjectStorage.h @@ -155,7 +155,17 @@ class StorageObjectStorage : public IStorage void addInferredEngineArgsToCreateQuery(ASTs & args, const ContextPtr & context) const override; - void updateFileConstantColumns(ContextPtr query_context); + /// Resolves identity-partition column names for the given already-pinned table + /// state snapshot (no extra state resolution) and stores them into `metadata`, + /// gated by the `allow_experimental_iceberg_read_optimization` setting. Shared by + /// the table-function constructor, `updateExternalDynamicMetadataIfExists`, and + /// `StorageObjectStorageCluster`. + static void updateIdentityPartitionColumns( + StorageInMemoryMetadata & metadata, + const StorageObjectStorageConfigurationPtr & configuration, + const DataLakeTableStateSnapshot & state, + ContextPtr context); + void updateExternalDynamicMetadataIfExists(ContextPtr query_context) override; IDataLakeMetadata * getExternalMetadata(ContextPtr query_context); @@ -215,9 +225,6 @@ class StorageObjectStorage : public IStorage NamesAndTypesList hive_partition_columns_to_read_from_file_path; NamesAndTypesList file_columns; - mutable SharedMutex mutex_file_constant_columns; - Names file_constant_columns TSA_GUARDED_BY(mutex_file_constant_columns); - LoggerPtr log; std::shared_ptr catalog; diff --git a/src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp b/src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp index 438fb2432aa8..7d9888b3f1db 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp +++ b/src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp @@ -566,13 +566,16 @@ void StorageObjectStorageCluster::updateExternalDynamicMetadataIfExists(ContextP new_metadata = *metadata_snapshot; } + /// Resolved from the same pinned `state` above -- no second state resolution. + /// Stored directly in `new_metadata` so it reaches `pure_storage` (and, through + /// it, distributed workers constructed via the rewritten table function) via the + /// same setInMemoryMetadata propagation below, with no separate call needed. + StorageObjectStorage::updateIdentityPartitionColumns(new_metadata, configuration, *state, query_context); + setInMemoryMetadata(new_metadata); if (pure_storage) - { pure_storage->setInMemoryMetadata(IStorageCluster::getInMemoryMetadata()); - pure_storage->updateFileConstantColumns(query_context); - } } class TaskDistributor : public TaskIterator diff --git a/src/Storages/ObjectStorage/StorageObjectStorageConfiguration.h b/src/Storages/ObjectStorage/StorageObjectStorageConfiguration.h index 112e7e949949..ae9a3af4b5bf 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorageConfiguration.h +++ b/src/Storages/ObjectStorage/StorageObjectStorageConfiguration.h @@ -312,7 +312,7 @@ class StorageObjectStorageConfiguration virtual ColumnMapperPtr getColumnMapperForCurrentSchema(StorageMetadataPtr /**/, ContextPtr /**/) const { return nullptr; } - virtual Names getIdentityPartitionColumnNames(ContextPtr) const { return {}; } + virtual Names getIdentityPartitionColumnNames(const DataLakeTableStateSnapshot &, ContextPtr) const { return {}; } virtual std::shared_ptr getCatalog(ContextPtr /*context*/, const StorageID & /*table_id*/) const { diff --git a/src/Storages/ObjectStorage/StorageObjectStorageSource.cpp b/src/Storages/ObjectStorage/StorageObjectStorageSource.cpp index e881bdfeaf21..616d3fb5ae31 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorageSource.cpp +++ b/src/Storages/ObjectStorage/StorageObjectStorageSource.cpp @@ -955,8 +955,8 @@ StorageObjectStorageSource::ReaderHolder StorageObjectStorageSource::createReade if (requested_columns_copy.empty() && (!format_filter_info || (!format_filter_info->row_level_filter && !format_filter_info->prewhere_info))) need_only_count = true; + } } - } std::optional num_rows_from_cache = need_only_count && context_->getSettingsRef()[Setting::use_cache_for_count_from_files] ? try_get_num_rows_from_cache() : std::nullopt; diff --git a/src/Storages/StorageInMemoryMetadata.cpp b/src/Storages/StorageInMemoryMetadata.cpp index bb0e26c4cea0..007a91d43b0f 100644 --- a/src/Storages/StorageInMemoryMetadata.cpp +++ b/src/Storages/StorageInMemoryMetadata.cpp @@ -62,6 +62,7 @@ StorageInMemoryMetadata::StorageInMemoryMetadata(const StorageInMemoryMetadata & , comment(other.comment) , metadata_version(other.metadata_version) , datalake_table_state(other.datalake_table_state) + , identity_partition_columns(other.identity_partition_columns) { } @@ -99,6 +100,7 @@ StorageInMemoryMetadata & StorageInMemoryMetadata::operator=(const StorageInMemo comment = other.comment; metadata_version = other.metadata_version; datalake_table_state = other.datalake_table_state; + identity_partition_columns = other.identity_partition_columns; return *this; } @@ -245,6 +247,11 @@ void StorageInMemoryMetadata::setDataLakeTableState(const DataLakeTableStateSnap datalake_table_state = datalake_table_state_; } +void StorageInMemoryMetadata::setIdentityPartitionColumns(const Names & identity_partition_columns_) +{ + identity_partition_columns = identity_partition_columns_; +} + StorageInMemoryMetadata StorageInMemoryMetadata::withMetadataVersion(int32_t metadata_version_) const { StorageInMemoryMetadata copy(*this); diff --git a/src/Storages/StorageInMemoryMetadata.h b/src/Storages/StorageInMemoryMetadata.h index aa577f8d02ff..feec86049f89 100644 --- a/src/Storages/StorageInMemoryMetadata.h +++ b/src/Storages/StorageInMemoryMetadata.h @@ -78,6 +78,10 @@ struct StorageInMemoryMetadata /// Current state of a datalake table. std::optional datalake_table_state; + /// Names of identity-partition columns (constant within every data file), + /// resolved for the same pinned `datalake_table_state` snapshot above. + Names identity_partition_columns; + StorageInMemoryMetadata() = default; StorageInMemoryMetadata(const StorageInMemoryMetadata & other); @@ -130,6 +134,7 @@ struct StorageInMemoryMetadata void setSQLSecurity(const ASTSQLSecurity & sql_security); void setDataLakeTableState(const DataLakeTableStateSnapshot & datalake_table_state_); + void setIdentityPartitionColumns(const Names & identity_partition_columns_); UUID getDefinerID(ContextPtr context) const; /// Returns a copy of the context with the correct user from SQL security options. From 70b3658b9f354730d74a9a12b4cdcad659556fde Mon Sep 17 00:00:00 2001 From: Konstantin Morozov Date: Wed, 22 Jul 2026 12:46:11 +0200 Subject: [PATCH 3/7] optimize only automatic prewhere Signed-off-by: Konstantin Morozov --- src/Interpreters/InterpreterSelectQuery.cpp | 2 +- .../QueryPlan/Optimizations/optimizePrewhere.cpp | 2 +- src/Storages/IStorage.h | 2 ++ src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.cpp | 8 +++++--- src/Storages/ObjectStorage/StorageObjectStorage.cpp | 7 ++++++- src/Storages/ObjectStorage/StorageObjectStorage.h | 1 + .../ObjectStorage/StorageObjectStorageCluster.cpp | 9 +++++++++ src/Storages/ObjectStorage/StorageObjectStorageCluster.h | 1 + 8 files changed, 26 insertions(+), 6 deletions(-) diff --git a/src/Interpreters/InterpreterSelectQuery.cpp b/src/Interpreters/InterpreterSelectQuery.cpp index bc7119f0b132..c75aaa782c10 100644 --- a/src/Interpreters/InterpreterSelectQuery.cpp +++ b/src/Interpreters/InterpreterSelectQuery.cpp @@ -831,7 +831,7 @@ InterpreterSelectQuery::InterpreterSelectQuery( current_info.syntax_analyzer_result = syntax_analyzer_result; Names queried_columns = syntax_analyzer_result->requiredSourceColumns(); - const auto & supported_prewhere_columns = storage->supportedPrewhereColumns(); + const auto & supported_prewhere_columns = storage->supportedAutomaticPrewhereColumns(); RangesInDataParts parts_for_estimator; if (storage_snapshot->data) diff --git a/src/Processors/QueryPlan/Optimizations/optimizePrewhere.cpp b/src/Processors/QueryPlan/Optimizations/optimizePrewhere.cpp index 1b0267e8e084..c4e904a7f1a6 100644 --- a/src/Processors/QueryPlan/Optimizations/optimizePrewhere.cpp +++ b/src/Processors/QueryPlan/Optimizations/optimizePrewhere.cpp @@ -186,7 +186,7 @@ void optimizePrewhere(QueryPlan::Node & parent_node, const bool remove_unused_co storage_snapshot, read_from_merge_tree_step ? read_from_merge_tree_step->getConditionSelectivityEstimator(queried_columns) : nullptr, queried_columns, - storage.supportedPrewhereColumns(), + storage.supportedAutomaticPrewhereColumns(), getLogger("QueryPlanOptimizePrewhere")}; auto optimize_result = where_optimizer.optimize(filter_step->getExpression(), diff --git a/src/Storages/IStorage.h b/src/Storages/IStorage.h index 8d3cb20b8c58..4363b6b3a478 100644 --- a/src/Storages/IStorage.h +++ b/src/Storages/IStorage.h @@ -146,6 +146,8 @@ class IStorage : public std::enable_shared_from_this, public TypePromo /// This is needed for engines whose aggregates data from multiple tables, like Merge. virtual std::optional supportedPrewhereColumns() const { return std::nullopt; } + virtual std::optional supportedAutomaticPrewhereColumns() const { return supportedPrewhereColumns(); } + /// Returns true if the storage supports optimization of moving conditions to PREWHERE section. virtual bool canMoveConditionsToPrewhere() const { return supportsPrewhere(); } diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.cpp index 997816680932..8bd320de36ec 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.cpp @@ -1473,13 +1473,16 @@ std::optional getPartitionKeyStringFromMetadata(Poco::JSON::Object::Ptr Names getIdentityPartitionColumnsFromMetadata(Poco::JSON::Object::Ptr metadata_object) { - // @todo refactor dupli if (!metadata_object->has(f_partition_specs) || !metadata_object->has(f_default_spec_id)) return {}; auto partition_spec_id = metadata_object->getValue(f_default_spec_id); Poco::JSON::Array::Ptr partition_specs = metadata_object->getArray(f_partition_specs); + if (!partition_specs) + return {}; + std::unordered_map source_id_to_column_name; + auto [schema, current_schema_id] = parseTableSchemaV2Method(metadata_object); auto mapper = createColumnMapper(schema)->getStorageColumnEncoding(); for (const auto & [col_name, source_id] : mapper) @@ -1500,11 +1503,10 @@ Names getIdentityPartitionColumnsFromMetadata(Poco::JSON::Object::Ptr metadata_o return {}; auto fields = partition_spec->getArray(f_fields); - if (fields->size() == 0) + if (!fields || fields->size() == 0) return {}; Names result; - std::vector part_exprs; for (UInt32 i = 0; i < fields->size(); ++i) { auto field = fields->getObject(i); diff --git a/src/Storages/ObjectStorage/StorageObjectStorage.cpp b/src/Storages/ObjectStorage/StorageObjectStorage.cpp index c94f846ae4a2..2c8dd9c81ceb 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorage.cpp +++ b/src/Storages/ObjectStorage/StorageObjectStorage.cpp @@ -342,6 +342,11 @@ bool StorageObjectStorage::canMoveConditionsToPrewhere() const } std::optional StorageObjectStorage::supportedPrewhereColumns() const +{ + return getInMemoryMetadataPtr()->getColumnsWithoutDefaultExpressions(/*exclude=*/hive_partition_columns_to_read_from_file_path); +} + +std::optional StorageObjectStorage::supportedAutomaticPrewhereColumns() const { auto exclude = hive_partition_columns_to_read_from_file_path; @@ -352,7 +357,7 @@ std::optional StorageObjectStorage::supportedPrewhereColumns() const exclude.emplace_back(identity_partition_column, metadata->getColumns().get(identity_partition_column).type); } - LOG_DEBUG(log, "Prewhere exclude list: [{}]", exclude.toString()); + LOG_DEBUG(log, "Automatic-PREWHERE exclude list: [{}]", exclude.toString()); return metadata->getColumnsWithoutDefaultExpressions(/*exclude=*/exclude); } diff --git a/src/Storages/ObjectStorage/StorageObjectStorage.h b/src/Storages/ObjectStorage/StorageObjectStorage.h index 5208536f48e9..619d2997f7db 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorage.h +++ b/src/Storages/ObjectStorage/StorageObjectStorage.h @@ -124,6 +124,7 @@ class StorageObjectStorage : public IStorage bool supportsPrewhere() const override; bool canMoveConditionsToPrewhere() const override; std::optional supportedPrewhereColumns() const override; + std::optional supportedAutomaticPrewhereColumns() const override; ColumnSizeByName getColumnSizes() const override; bool prefersLargeBlocks() const override; diff --git a/src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp b/src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp index 7d9888b3f1db..b62e1dec1ee6 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp +++ b/src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp @@ -212,6 +212,8 @@ StorageObjectStorageCluster::StorageObjectStorageCluster( if (auto metadata_snapshot = configuration->buildStorageMetadataFromState(*state, context_)) metadata = *metadata_snapshot; } + + StorageObjectStorage::updateIdentityPartitionColumns(metadata, configuration, *state, context_); } } @@ -1068,6 +1070,13 @@ std::optional StorageObjectStorageCluster::supportedPrewhereColumns() c return IStorageCluster::supportedPrewhereColumns(); } +std::optional StorageObjectStorageCluster::supportedAutomaticPrewhereColumns() const +{ + if (pure_storage) + return pure_storage->supportedAutomaticPrewhereColumns(); + return IStorageCluster::supportedAutomaticPrewhereColumns(); +} + IStorageCluster::ColumnSizeByName StorageObjectStorageCluster::getColumnSizes() const { if (pure_storage) diff --git a/src/Storages/ObjectStorage/StorageObjectStorageCluster.h b/src/Storages/ObjectStorage/StorageObjectStorageCluster.h index 91d07b5c498a..29ae5da41fe5 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorageCluster.h +++ b/src/Storages/ObjectStorage/StorageObjectStorageCluster.h @@ -163,6 +163,7 @@ class StorageObjectStorageCluster : public IStorageCluster bool supportsPrewhere() const override; bool canMoveConditionsToPrewhere() const override; std::optional supportedPrewhereColumns() const override; + std::optional supportedAutomaticPrewhereColumns() const override; ColumnSizeByName getColumnSizes() const override; bool parallelizeOutputAfterReading(ContextPtr context) const override; From 21666f5d2b52ea5db7a4614e4c861eea62720023 Mon Sep 17 00:00:00 2001 From: Konstantin Morozov Date: Wed, 22 Jul 2026 13:09:30 +0200 Subject: [PATCH 4/7] no metadata race Signed-off-by: Konstantin Morozov --- src/Interpreters/InterpreterSelectQuery.cpp | 2 +- src/Processors/QueryPlan/Optimizations/optimizePrewhere.cpp | 2 +- src/Storages/IStorage.h | 2 +- src/Storages/ObjectStorage/StorageObjectStorage.cpp | 3 +-- src/Storages/ObjectStorage/StorageObjectStorage.h | 2 +- src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp | 6 +++--- src/Storages/ObjectStorage/StorageObjectStorageCluster.h | 2 +- 7 files changed, 9 insertions(+), 10 deletions(-) diff --git a/src/Interpreters/InterpreterSelectQuery.cpp b/src/Interpreters/InterpreterSelectQuery.cpp index c75aaa782c10..48bfb1a57b2e 100644 --- a/src/Interpreters/InterpreterSelectQuery.cpp +++ b/src/Interpreters/InterpreterSelectQuery.cpp @@ -831,7 +831,7 @@ InterpreterSelectQuery::InterpreterSelectQuery( current_info.syntax_analyzer_result = syntax_analyzer_result; Names queried_columns = syntax_analyzer_result->requiredSourceColumns(); - const auto & supported_prewhere_columns = storage->supportedAutomaticPrewhereColumns(); + const auto & supported_prewhere_columns = storage->supportedAutomaticPrewhereColumns(storage_snapshot->metadata); RangesInDataParts parts_for_estimator; if (storage_snapshot->data) diff --git a/src/Processors/QueryPlan/Optimizations/optimizePrewhere.cpp b/src/Processors/QueryPlan/Optimizations/optimizePrewhere.cpp index c4e904a7f1a6..a2269ae0ca28 100644 --- a/src/Processors/QueryPlan/Optimizations/optimizePrewhere.cpp +++ b/src/Processors/QueryPlan/Optimizations/optimizePrewhere.cpp @@ -186,7 +186,7 @@ void optimizePrewhere(QueryPlan::Node & parent_node, const bool remove_unused_co storage_snapshot, read_from_merge_tree_step ? read_from_merge_tree_step->getConditionSelectivityEstimator(queried_columns) : nullptr, queried_columns, - storage.supportedAutomaticPrewhereColumns(), + storage.supportedAutomaticPrewhereColumns(storage_snapshot->metadata), getLogger("QueryPlanOptimizePrewhere")}; auto optimize_result = where_optimizer.optimize(filter_step->getExpression(), diff --git a/src/Storages/IStorage.h b/src/Storages/IStorage.h index 4363b6b3a478..6162f6aa6f7d 100644 --- a/src/Storages/IStorage.h +++ b/src/Storages/IStorage.h @@ -146,7 +146,7 @@ class IStorage : public std::enable_shared_from_this, public TypePromo /// This is needed for engines whose aggregates data from multiple tables, like Merge. virtual std::optional supportedPrewhereColumns() const { return std::nullopt; } - virtual std::optional supportedAutomaticPrewhereColumns() const { return supportedPrewhereColumns(); } + virtual std::optional supportedAutomaticPrewhereColumns(const StorageMetadataPtr & /* metadata */) const { return supportedPrewhereColumns(); } /// Returns true if the storage supports optimization of moving conditions to PREWHERE section. virtual bool canMoveConditionsToPrewhere() const { return supportsPrewhere(); } diff --git a/src/Storages/ObjectStorage/StorageObjectStorage.cpp b/src/Storages/ObjectStorage/StorageObjectStorage.cpp index 2c8dd9c81ceb..0f6063b3a942 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorage.cpp +++ b/src/Storages/ObjectStorage/StorageObjectStorage.cpp @@ -346,11 +346,10 @@ std::optional StorageObjectStorage::supportedPrewhereColumns() const return getInMemoryMetadataPtr()->getColumnsWithoutDefaultExpressions(/*exclude=*/hive_partition_columns_to_read_from_file_path); } -std::optional StorageObjectStorage::supportedAutomaticPrewhereColumns() const +std::optional StorageObjectStorage::supportedAutomaticPrewhereColumns(const StorageMetadataPtr & metadata) const { auto exclude = hive_partition_columns_to_read_from_file_path; - const auto metadata = getInMemoryMetadataPtr(); for (const auto & identity_partition_column : metadata->identity_partition_columns) { if (metadata->getColumns().has(identity_partition_column)) diff --git a/src/Storages/ObjectStorage/StorageObjectStorage.h b/src/Storages/ObjectStorage/StorageObjectStorage.h index 619d2997f7db..8128188c7927 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorage.h +++ b/src/Storages/ObjectStorage/StorageObjectStorage.h @@ -124,7 +124,7 @@ class StorageObjectStorage : public IStorage bool supportsPrewhere() const override; bool canMoveConditionsToPrewhere() const override; std::optional supportedPrewhereColumns() const override; - std::optional supportedAutomaticPrewhereColumns() const override; + std::optional supportedAutomaticPrewhereColumns(const StorageMetadataPtr & metadata) const override; ColumnSizeByName getColumnSizes() const override; bool prefersLargeBlocks() const override; diff --git a/src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp b/src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp index b62e1dec1ee6..d5d73b5fab47 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp +++ b/src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp @@ -1070,11 +1070,11 @@ std::optional StorageObjectStorageCluster::supportedPrewhereColumns() c return IStorageCluster::supportedPrewhereColumns(); } -std::optional StorageObjectStorageCluster::supportedAutomaticPrewhereColumns() const +std::optional StorageObjectStorageCluster::supportedAutomaticPrewhereColumns(const StorageMetadataPtr & metadata) const { if (pure_storage) - return pure_storage->supportedAutomaticPrewhereColumns(); - return IStorageCluster::supportedAutomaticPrewhereColumns(); + return pure_storage->supportedAutomaticPrewhereColumns(metadata); + return IStorageCluster::supportedAutomaticPrewhereColumns(metadata); } IStorageCluster::ColumnSizeByName StorageObjectStorageCluster::getColumnSizes() const diff --git a/src/Storages/ObjectStorage/StorageObjectStorageCluster.h b/src/Storages/ObjectStorage/StorageObjectStorageCluster.h index 29ae5da41fe5..e0461c72eda3 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorageCluster.h +++ b/src/Storages/ObjectStorage/StorageObjectStorageCluster.h @@ -163,7 +163,7 @@ class StorageObjectStorageCluster : public IStorageCluster bool supportsPrewhere() const override; bool canMoveConditionsToPrewhere() const override; std::optional supportedPrewhereColumns() const override; - std::optional supportedAutomaticPrewhereColumns() const override; + std::optional supportedAutomaticPrewhereColumns(const StorageMetadataPtr & metadata) const override; ColumnSizeByName getColumnSizes() const override; bool parallelizeOutputAfterReading(ContextPtr context) const override; From bc0f7f83760e71e968a43fabfff98573b568ce75 Mon Sep 17 00:00:00 2001 From: Konstantin Morozov Date: Wed, 22 Jul 2026 15:16:30 +0200 Subject: [PATCH 5/7] update iceberg metadata instead storage Signed-off-by: Konstantin Morozov --- .../DataLakes/DataLakeConfiguration.h | 11 ----- .../DataLakes/IDataLakeMetadata.h | 2 - .../DataLakes/Iceberg/IcebergMetadata.cpp | 41 ++++++++++--------- .../DataLakes/Iceberg/IcebergMetadata.h | 3 +- .../ObjectStorage/StorageObjectStorage.cpp | 19 --------- .../ObjectStorage/StorageObjectStorage.h | 11 ----- .../StorageObjectStorageCluster.cpp | 8 ---- .../StorageObjectStorageConfiguration.h | 2 - 8 files changed, 23 insertions(+), 74 deletions(-) diff --git a/src/Storages/ObjectStorage/DataLakes/DataLakeConfiguration.h b/src/Storages/ObjectStorage/DataLakes/DataLakeConfiguration.h index 1ee40645ae24..74d7ceea7f93 100644 --- a/src/Storages/ObjectStorage/DataLakes/DataLakeConfiguration.h +++ b/src/Storages/ObjectStorage/DataLakes/DataLakeConfiguration.h @@ -323,15 +323,6 @@ class DataLakeConfiguration : public BaseStorageConfiguration, public std::enabl return current_metadata->getColumnMapperForCurrentSchema(storage_metadata_snapshot, context); } - Names getIdentityPartitionColumnNames(const DataLakeTableStateSnapshot & state, ContextPtr context) const override - { - if (!current_metadata) - { - return {}; - } - return current_metadata->getIdentityPartitionColumnNames(state, context); - } - void drop(ContextPtr local_context) override { if (current_metadata) @@ -815,8 +806,6 @@ class StorageIcebergConfiguration : public StorageObjectStorageConfiguration, pu ColumnMapperPtr getColumnMapperForCurrentSchema(StorageMetadataPtr storage_metadata_snapshot, ContextPtr context) const override { return getImpl().getColumnMapperForCurrentSchema(storage_metadata_snapshot, context); } - Names getIdentityPartitionColumnNames(const DataLakeTableStateSnapshot & state, ContextPtr context) const override { return getImpl().getIdentityPartitionColumnNames(state, context); } - std::shared_ptr getCatalog(ContextPtr context, const StorageID & table_id) const override { return getImpl().getCatalog(context, table_id); } diff --git a/src/Storages/ObjectStorage/DataLakes/IDataLakeMetadata.h b/src/Storages/ObjectStorage/DataLakes/IDataLakeMetadata.h index e64db238a62b..7fad60407ed5 100644 --- a/src/Storages/ObjectStorage/DataLakes/IDataLakeMetadata.h +++ b/src/Storages/ObjectStorage/DataLakes/IDataLakeMetadata.h @@ -193,8 +193,6 @@ class IDataLakeMetadata : boost::noncopyable virtual ColumnMapperPtr getColumnMapperForObject(ObjectInfoPtr /**/) const { return nullptr; } virtual ColumnMapperPtr getColumnMapperForCurrentSchema(StorageMetadataPtr, ContextPtr) const { return nullptr; } - virtual Names getIdentityPartitionColumnNames(const DataLakeTableStateSnapshot &, ContextPtr) const { return {}; } - virtual SinkToStoragePtr write( SharedHeader /*sample_block*/, const StorageID & /*table_id*/, diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.cpp index 661c45606f6f..b32b07a04aca 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.cpp @@ -134,6 +134,7 @@ extern const SettingsBool allow_insert_into_iceberg; extern const SettingsBool allow_experimental_iceberg_compaction; extern const SettingsBool allow_experimental_expire_snapshots; extern const SettingsBool iceberg_delete_data_on_drop; +extern const SettingsBool allow_experimental_iceberg_read_optimization; } static constexpr size_t MAX_TRANSACTION_RETRIES = 100; @@ -1236,7 +1237,21 @@ std::unique_ptr IcebergMetadata::buildStorageMetadataFr result->setColumns( ColumnsDescription{*persistent_components.schema_processor->getClickhouseTableSchemaById(iceberg_state.schema_id)}); result->setDataLakeTableState(state); - result->sorting_key = getSortingKey(local_context, iceberg_state); + + auto metadata_object = getMetadataJSONObject( + iceberg_state.metadata_file_path, + object_storage, + persistent_components.metadata_cache, + local_context, + log, + persistent_components.metadata_compression_method, + persistent_components.table_uuid); + + result->sorting_key = getSortingKeyFromMetadata(metadata_object, local_context); + + if (local_context->getSettingsRef()[Setting::allow_experimental_iceberg_read_optimization]) + result->setIdentityPartitionColumns(getIdentityPartitionColumnsFromMetadata(metadata_object)); + return result; } @@ -1436,22 +1451,6 @@ ColumnMapperPtr IcebergMetadata::getColumnMapperForCurrentSchema(StorageMetadata return persistent_components.schema_processor->getColumnMapperById(iceberg_table_state->schema_id); } -Names IcebergMetadata::getIdentityPartitionColumnNames(const DataLakeTableStateSnapshot & state, ContextPtr local_context) const -{ - // @todo - chassert(std::holds_alternative(state)); - const auto & iceberg_state = std::get(state); - auto metadata_object = getMetadataJSONObject( - iceberg_state.metadata_file_path, - object_storage, - persistent_components.metadata_cache, - local_context, - log, - persistent_components.metadata_compression_method, - persistent_components.table_uuid); - return getIdentityPartitionColumnsFromMetadata(metadata_object); -} - std::optional IcebergMetadata::getPartitionKey(ContextPtr local_context, TableStateSnapshot actual_table_state_snapshot) const { auto metadata_object = getMetadataJSONObject( @@ -1480,10 +1479,14 @@ KeyDescription IcebergMetadata::getSortingKey(ContextPtr local_context, TableSta persistent_components.metadata_compression_method, persistent_components.table_uuid); + return getSortingKeyFromMetadata(metadata_object, local_context); +} + +KeyDescription IcebergMetadata::getSortingKeyFromMetadata(const Poco::JSON::Object::Ptr & metadata_object, ContextPtr local_context) const +{ auto [schema, current_schema_id] = parseTableSchemaV2Method(metadata_object); auto result = getSortingKeyDescriptionFromMetadata(metadata_object, *persistent_components.schema_processor->getClickhouseTableSchemaById(current_schema_id), local_context); - auto sort_order_id = metadata_object->getValue(f_default_sort_order_id); - result.sort_order_id = sort_order_id; + result.sort_order_id = metadata_object->getValue(f_default_sort_order_id); return result; } diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.h b/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.h index 13f3849d35a2..3324e9ec5015 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.h +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.h @@ -106,8 +106,6 @@ class IcebergMetadata : public IDataLakeMetadata ColumnMapperPtr getColumnMapperForCurrentSchema(StorageMetadataPtr storage_metadata_snapshot, ContextPtr context) const override; - Names getIdentityPartitionColumnNames(const DataLakeTableStateSnapshot & state, ContextPtr context) const override; - SinkToStoragePtr write( SharedHeader sample_block, const StorageID & table_id, @@ -211,6 +209,7 @@ class IcebergMetadata : public IDataLakeMetadata std::optional getPartitionKey(ContextPtr local_context, Iceberg::TableStateSnapshot actual_table_state_snapshot) const; KeyDescription getSortingKey(ContextPtr local_context, Iceberg::TableStateSnapshot actual_table_state_snapshot) const; + KeyDescription getSortingKeyFromMetadata(const Poco::JSON::Object::Ptr & metadata_object, ContextPtr local_context) const; /// Non-empty return value means the attempt succeeded (covers both the normal /// publish path and the `isExportPartitionTransactionAlreadyCommitted` short-circuit). diff --git a/src/Storages/ObjectStorage/StorageObjectStorage.cpp b/src/Storages/ObjectStorage/StorageObjectStorage.cpp index 0f6063b3a942..40656cdffbbc 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorage.cpp +++ b/src/Storages/ObjectStorage/StorageObjectStorage.cpp @@ -285,12 +285,6 @@ StorageObjectStorage::StorageObjectStorage( if (auto metadata_snapshot = configuration->buildStorageMetadataFromState(*state, context)) metadata = *metadata_snapshot; } - - /// Table functions bypass updateExternalDynamicMetadataIfExists (see above), - /// so identity-partition columns must be resolved here as well -- otherwise - /// PREWHERE could be built over a column that gets erased from the read set - /// as file-constant (see StorageObjectStorageSource's constant-column handling). - updateIdentityPartitionColumns(metadata, configuration, *state, context); } } @@ -375,16 +369,6 @@ configuration->update(object_storage, query_context); return configuration->getExternalMetadata(); } -void StorageObjectStorage::updateIdentityPartitionColumns( - StorageInMemoryMetadata & metadata, - const StorageObjectStorageConfigurationPtr & configuration, - const DataLakeTableStateSnapshot & state, - ContextPtr context) -{ - if (context->getSettingsRef()[Setting::allow_experimental_iceberg_read_optimization]) - metadata.setIdentityPartitionColumns(configuration->getIdentityPartitionColumnNames(state, context)); -} - void StorageObjectStorage::updateExternalDynamicMetadataIfExists(ContextPtr query_context) { if (!configuration->isDataLakeConfiguration()) @@ -412,9 +396,6 @@ void StorageObjectStorage::updateExternalDynamicMetadataIfExists(ContextPtr quer new_metadata = *metadata_snapshot; } - /// Resolved from the same pinned `state` above -- no second state resolution. - updateIdentityPartitionColumns(new_metadata, configuration, *state, query_context); - setInMemoryMetadata(new_metadata); } diff --git a/src/Storages/ObjectStorage/StorageObjectStorage.h b/src/Storages/ObjectStorage/StorageObjectStorage.h index 8128188c7927..e0d3f6f99ba1 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorage.h +++ b/src/Storages/ObjectStorage/StorageObjectStorage.h @@ -156,17 +156,6 @@ class StorageObjectStorage : public IStorage void addInferredEngineArgsToCreateQuery(ASTs & args, const ContextPtr & context) const override; - /// Resolves identity-partition column names for the given already-pinned table - /// state snapshot (no extra state resolution) and stores them into `metadata`, - /// gated by the `allow_experimental_iceberg_read_optimization` setting. Shared by - /// the table-function constructor, `updateExternalDynamicMetadataIfExists`, and - /// `StorageObjectStorageCluster`. - static void updateIdentityPartitionColumns( - StorageInMemoryMetadata & metadata, - const StorageObjectStorageConfigurationPtr & configuration, - const DataLakeTableStateSnapshot & state, - ContextPtr context); - void updateExternalDynamicMetadataIfExists(ContextPtr query_context) override; IDataLakeMetadata * getExternalMetadata(ContextPtr query_context); diff --git a/src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp b/src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp index d5d73b5fab47..1247c7a320cb 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp +++ b/src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp @@ -212,8 +212,6 @@ StorageObjectStorageCluster::StorageObjectStorageCluster( if (auto metadata_snapshot = configuration->buildStorageMetadataFromState(*state, context_)) metadata = *metadata_snapshot; } - - StorageObjectStorage::updateIdentityPartitionColumns(metadata, configuration, *state, context_); } } @@ -568,12 +566,6 @@ void StorageObjectStorageCluster::updateExternalDynamicMetadataIfExists(ContextP new_metadata = *metadata_snapshot; } - /// Resolved from the same pinned `state` above -- no second state resolution. - /// Stored directly in `new_metadata` so it reaches `pure_storage` (and, through - /// it, distributed workers constructed via the rewritten table function) via the - /// same setInMemoryMetadata propagation below, with no separate call needed. - StorageObjectStorage::updateIdentityPartitionColumns(new_metadata, configuration, *state, query_context); - setInMemoryMetadata(new_metadata); if (pure_storage) diff --git a/src/Storages/ObjectStorage/StorageObjectStorageConfiguration.h b/src/Storages/ObjectStorage/StorageObjectStorageConfiguration.h index ae9a3af4b5bf..708cdeb91f65 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorageConfiguration.h +++ b/src/Storages/ObjectStorage/StorageObjectStorageConfiguration.h @@ -312,8 +312,6 @@ class StorageObjectStorageConfiguration virtual ColumnMapperPtr getColumnMapperForCurrentSchema(StorageMetadataPtr /**/, ContextPtr /**/) const { return nullptr; } - virtual Names getIdentityPartitionColumnNames(const DataLakeTableStateSnapshot &, ContextPtr) const { return {}; } - virtual std::shared_ptr getCatalog(ContextPtr /*context*/, const StorageID & /*table_id*/) const { return nullptr; From 979957bf95c7cd15f2b0654ed17ee5eeaa04597f Mon Sep 17 00:00:00 2001 From: Konstantin Morozov Date: Wed, 22 Jul 2026 15:38:50 +0200 Subject: [PATCH 6/7] refactor utils Signed-off-by: Konstantin Morozov --- .../ObjectStorage/DataLakes/Iceberg/Utils.cpp | 127 +++++++++--------- 1 file changed, 61 insertions(+), 66 deletions(-) diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.cpp index 8bd320de36ec..87852dfb04fa 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.cpp @@ -108,6 +108,53 @@ static constexpr auto MAX_TRANSACTION_RETRIES = 100; namespace DB::Iceberg { using namespace DB; + +namespace +{ +struct ResolvedPartitionSpec +{ + Poco::JSON::Array::Ptr fields; + std::unordered_map source_id_to_column_name; +}; + +std::optional resolveDefaultPartitionSpec(const Poco::JSON::Object::Ptr & metadata_object) +{ + if (!metadata_object->has(f_partition_specs) || !metadata_object->has(f_default_spec_id)) + return std::nullopt; + + auto partition_spec_id = metadata_object->getValue(f_default_spec_id); + Poco::JSON::Array::Ptr partition_specs = metadata_object->getArray(f_partition_specs); + if (!partition_specs) + return std::nullopt; + + std::unordered_map source_id_to_column_name; + auto [schema, current_schema_id] = parseTableSchemaV2Method(metadata_object); + auto mapper = createColumnMapper(schema)->getStorageColumnEncoding(); + for (const auto & [col_name, source_id] : mapper) + source_id_to_column_name[source_id] = col_name; + + Poco::JSON::Object::Ptr partition_spec; + for (size_t i = 0; i < partition_specs->size(); ++i) + { + auto spec = partition_specs->getObject(static_cast(i)); + if (spec && spec->getValue(f_spec_id) == partition_spec_id) + { + partition_spec = spec; + break; + } + } + + if (!partition_spec || !partition_spec->has(f_fields)) + return std::nullopt; + + auto fields = partition_spec->getArray(f_fields); + if (!fields || fields->size() == 0) + return std::nullopt; + + return ResolvedPartitionSpec{fields, std::move(source_id_to_column_name)}; +} +} + static CompressionMethod getCompressionMethodFromMetadataFile(const String & path) { constexpr std::string_view metadata_suffix = ".metadata.json"; @@ -1423,44 +1470,20 @@ static String formatPartitionFieldDisplay(const String & iceberg_transform_name, std::optional getPartitionKeyStringFromMetadata(Poco::JSON::Object::Ptr metadata_object, const NamesAndTypesList & /* ch_schema */, ContextPtr /* local_context */) { - if (!metadata_object->has(f_partition_specs) || !metadata_object->has(f_default_spec_id)) - return std::nullopt; - auto partition_spec_id = metadata_object->getValue(f_default_spec_id); - Poco::JSON::Array::Ptr partition_specs = metadata_object->getArray(f_partition_specs); - std::unordered_map source_id_to_column_name; - auto [schema, current_schema_id] = parseTableSchemaV2Method(metadata_object); - auto mapper = createColumnMapper(schema)->getStorageColumnEncoding(); - for (const auto & [col_name, source_id] : mapper) - source_id_to_column_name[source_id] = col_name; - - Poco::JSON::Object::Ptr partition_spec; - for (size_t i = 0; i < partition_specs->size(); ++i) - { - auto spec = partition_specs->getObject(static_cast(i)); - if (spec->getValue(f_spec_id) == partition_spec_id) - { - partition_spec = spec; - break; - } - } - if (!partition_spec || !partition_spec->has(f_fields)) - return std::nullopt; - auto fields = partition_spec->getArray(f_fields); - if (fields->size() == 0) + auto resolved = resolveDefaultPartitionSpec(metadata_object); + if (!resolved) return std::nullopt; std::vector part_exprs; - for (UInt32 i = 0; i < fields->size(); ++i) + for (UInt32 i = 0; i < resolved->fields->size(); ++i) { - auto field = fields->getObject(i); - auto source_id = field->getValue(f_source_id); - auto it = source_id_to_column_name.find(source_id); - if (it == source_id_to_column_name.end()) + auto field = resolved->fields->getObject(i); + auto it = resolved->source_id_to_column_name.find(field->getValue(f_source_id)); + if (it == resolved->source_id_to_column_name.end()) return std::nullopt; - String column_name = it->second; - auto iceberg_transform_name = field->getValue(f_transform); - part_exprs.push_back(formatPartitionFieldDisplay(iceberg_transform_name, column_name)); + part_exprs.push_back(formatPartitionFieldDisplay(field->getValue(f_transform), it->second)); } + String result; for (size_t i = 0; i < part_exprs.size(); ++i) { @@ -1473,48 +1496,20 @@ std::optional getPartitionKeyStringFromMetadata(Poco::JSON::Object::Ptr Names getIdentityPartitionColumnsFromMetadata(Poco::JSON::Object::Ptr metadata_object) { - if (!metadata_object->has(f_partition_specs) || !metadata_object->has(f_default_spec_id)) - return {}; - - auto partition_spec_id = metadata_object->getValue(f_default_spec_id); - Poco::JSON::Array::Ptr partition_specs = metadata_object->getArray(f_partition_specs); - if (!partition_specs) - return {}; - - std::unordered_map source_id_to_column_name; - - auto [schema, current_schema_id] = parseTableSchemaV2Method(metadata_object); - auto mapper = createColumnMapper(schema)->getStorageColumnEncoding(); - for (const auto & [col_name, source_id] : mapper) - source_id_to_column_name[source_id] = col_name; - - Poco::JSON::Object::Ptr partition_spec; - for (size_t i = 0; i < partition_specs->size(); ++i) - { - auto spec = partition_specs->getObject(static_cast(i)); - if (spec->getValue(f_spec_id) == partition_spec_id) - { - partition_spec = spec; - break; - } - } - - if (!partition_spec || !partition_spec->has(f_fields)) - return {}; - - auto fields = partition_spec->getArray(f_fields); - if (!fields || fields->size() == 0) + auto resolved = resolveDefaultPartitionSpec(metadata_object); + if (!resolved) return {}; Names result; - for (UInt32 i = 0; i < fields->size(); ++i) + for (UInt32 i = 0; i < resolved->fields->size(); ++i) { - auto field = fields->getObject(i); + auto field = resolved->fields->getObject(i); if (Poco::toLower(field->getValue(f_transform)) != "identity") continue; - if (auto it = source_id_to_column_name.find(field->getValue(f_source_id)); it != source_id_to_column_name.end()) + if (auto it = resolved->source_id_to_column_name.find(field->getValue(f_source_id)); + it != resolved->source_id_to_column_name.end()) result.push_back(it->second); } return result; From 89f8e8017b9f9aac82f005d146d42bac977e4bd5 Mon Sep 17 00:00:00 2001 From: Konstantin Morozov Date: Wed, 22 Jul 2026 17:13:52 +0200 Subject: [PATCH 7/7] add tests Signed-off-by: Konstantin Morozov --- ...st_read_optimization_partition_prewhere.py | 262 ++++++++++++++++++ 1 file changed, 262 insertions(+) create mode 100644 tests/integration/test_database_iceberg/test_read_optimization_partition_prewhere.py diff --git a/tests/integration/test_database_iceberg/test_read_optimization_partition_prewhere.py b/tests/integration/test_database_iceberg/test_read_optimization_partition_prewhere.py new file mode 100644 index 000000000000..d641b480fe60 --- /dev/null +++ b/tests/integration/test_database_iceberg/test_read_optimization_partition_prewhere.py @@ -0,0 +1,262 @@ +import logging +import time +import uuid +from datetime import date + +import pyarrow as pa +import pytest +from pyiceberg.catalog import load_catalog +from pyiceberg.partitioning import PartitionField, PartitionSpec +from pyiceberg.schema import Schema +from pyiceberg.transforms import BucketTransform, IdentityTransform +from pyiceberg.types import DateType, LongType, NestedField + +from helpers.cluster import ClickHouseCluster +from helpers.config_cluster import minio_access_key, minio_secret_key + +ICEBERG_PORT = 8184 + +BASE_URL = "http://rest:8181/v1" +BASE_URL_LOCAL_RAW = f"http://localhost:{ICEBERG_PORT}" + +CATALOG_NAME = "demo" + +SCHEMA = Schema( + NestedField(field_id=1, name="id", field_type=LongType(), required=False), + NestedField(field_id=2, name="date", field_type=DateType(), required=False), + NestedField(field_id=3, name="value", field_type=LongType(), required=False), +) + +IDENTITY_ON_DATE = PartitionSpec( + PartitionField(source_id=2, field_id=1000, transform=IdentityTransform(), name="date") +) + + +@pytest.fixture(scope="module") +def started_cluster(): + try: + cluster = ClickHouseCluster(__file__) + cluster.iceberg_rest_external_port = ICEBERG_PORT + cluster.add_instance( + "node1", + main_configs=["configs/cluster.xml"], + user_configs=[], + stay_alive=True, + with_iceberg_catalog=True, + ) + + logging.info("Starting cluster...") + cluster.start() + + time.sleep(10) + + yield cluster + + finally: + cluster.shutdown() + + +def load_catalog_impl(started_cluster): + return load_catalog( + CATALOG_NAME, + **{ + "uri": BASE_URL_LOCAL_RAW, + "type": "rest", + "s3.endpoint": f"http://{started_cluster.get_instance_ip('minio')}:9000", + "s3.access-key-id": minio_access_key, + "s3.secret-access-key": minio_secret_key, + }, + ) + + +def create_table(catalog, namespace, table, schema, partition_spec): + return catalog.create_table( + identifier=f"{namespace}.{table}", + schema=schema, + location="s3://warehouse-rest/data", + partition_spec=partition_spec, + ) + + +def create_clickhouse_iceberg_database(node, name): + settings = { + "catalog_type": "rest", + "warehouse": "demo", + "storage_endpoint": "http://minio:9000/warehouse-rest", + } + node.query( + f""" +DROP DATABASE IF EXISTS {name}; +SET allow_database_iceberg=true; +SET write_full_path_in_iceberg_metadata=1; +CREATE DATABASE {name} ENGINE = DataLakeCatalog('{BASE_URL}', 'minio', '{minio_secret_key}') +SETTINGS {",".join((k + "=" + repr(v) for k, v in settings.items()))} +""" + ) + + +def append_row(table, id_, date_, value_): + append_rows(table, [(id_, date_, value_)]) + + +def append_rows(table, rows): + """Writes all `rows` in a single `.append()` call, i.e. into the same + physical data file when they share a partition value.""" + table.append( + pa.Table.from_pylist( + [{"id": id_, "date": date_, "value": value_} for id_, date_, value_ in rows], + schema=table.schema().as_arrow(), + ) + ) + + +def _run_and_get_profile_events(instance, select_expression, settings, event_names): + query_id = f"read-opt-{uuid.uuid4()}" + result = instance.query(select_expression, query_id=query_id, settings=settings) + instance.query("SYSTEM FLUSH LOGS") + columns = ", ".join(f"ProfileEvents['{name}']" for name in event_names) + events_row = instance.query( + f""" + SELECT {columns} + FROM system.query_log + WHERE query_id = '{query_id}' AND type = 'QueryFinish' + """ + ).strip() + events = tuple(int(x) for x in events_row.split("\t")) + return result.strip(), events + + +def test_iceberg_read_optimization_count_with_partition_filter(started_cluster): + node = started_cluster.instances["node1"] + catalog = load_catalog_impl(started_cluster) + + namespace = f"read_opt_count_ns_{uuid.uuid4().hex}" + table_name = "t" + catalog.create_namespace(namespace) + table = create_table(catalog, namespace, table_name, SCHEMA, IDENTITY_ON_DATE) + + rows = [ + (1, date(2024, 1, 10), 100), + (2, date(2024, 1, 20), 200), + (3, date(2024, 2, 5), 300), + (4, date(2024, 2, 15), 400), + (5, date(2024, 3, 1), 500), + ] + for id_, date_, value_ in rows: + append_row(table, id_, date_, value_) + + create_clickhouse_iceberg_database(node, CATALOG_NAME) + full_name = f"{CATALOG_NAME}.`{namespace}.{table_name}`" + + select_expression = f"SELECT count() FROM {full_name} WHERE date > '2024-01-15'" + events = ("ObjectStorageReadObjects", "ParquetReadRowGroups", "IcebergPartitionPrunedFiles") + + baseline_result, (baseline_reads, baseline_row_groups, baseline_pruned) = _run_and_get_profile_events( + node, select_expression, {"allow_experimental_iceberg_read_optimization": 0}, events + ) + optimized_result, (optimized_reads, optimized_row_groups, optimized_pruned) = _run_and_get_profile_events( + node, select_expression, {"allow_experimental_iceberg_read_optimization": 1}, events + ) + + assert baseline_result == optimized_result == "4" + + assert baseline_pruned == optimized_pruned == 1 + + assert baseline_reads > 0 + assert baseline_row_groups > 0 + assert optimized_reads == 0 + assert optimized_row_groups == 0 + + +def test_iceberg_read_optimization_partition_filter_excludes_all_files(started_cluster): + node = started_cluster.instances["node1"] + catalog = load_catalog_impl(started_cluster) + + namespace = f"read_opt_excl_ns_{uuid.uuid4().hex}" + table_name = "t" + catalog.create_namespace(namespace) + table = create_table(catalog, namespace, table_name, SCHEMA, IDENTITY_ON_DATE) + + rows = [ + (1, date(2024, 1, 10), 100), + (2, date(2024, 1, 20), 200), + (3, date(2024, 2, 5), 300), + ] + for id_, date_, value_ in rows: + append_row(table, id_, date_, value_) + + create_clickhouse_iceberg_database(node, CATALOG_NAME) + full_name = f"{CATALOG_NAME}.`{namespace}.{table_name}`" + + select_expression = f"SELECT count() FROM {full_name} WHERE date > '2099-01-01'" + events = ("ObjectStorageReadObjects", "IcebergPartitionPrunedFiles") + + result, (reads, pruned) = _run_and_get_profile_events( + node, select_expression, {"allow_experimental_iceberg_read_optimization": 1}, events + ) + + assert result == "0" + assert pruned == 3 + assert reads == 0 + + +def test_iceberg_read_optimization_real_column_query_still_reads_data(started_cluster): + node = started_cluster.instances["node1"] + catalog = load_catalog_impl(started_cluster) + + namespace = f"read_opt_realcol_ns_{uuid.uuid4().hex}" + table_name = "t" + catalog.create_namespace(namespace) + table = create_table(catalog, namespace, table_name, SCHEMA, IDENTITY_ON_DATE) + + append_row(table, 1, date(2024, 1, 10), 100) + append_rows(table, [(2, date(2024, 1, 20), 200), (3, date(2024, 1, 20), 250)]) + append_row(table, 4, date(2024, 2, 15), 400) + + create_clickhouse_iceberg_database(node, CATALOG_NAME) + full_name = f"{CATALOG_NAME}.`{namespace}.{table_name}`" + + select_expression = f"SELECT count(), sum(value) FROM {full_name} WHERE date > '2024-01-15'" + events = ("ObjectStorageReadObjects",) + + baseline_result, (baseline_reads,) = _run_and_get_profile_events( + node, select_expression, {"allow_experimental_iceberg_read_optimization": 0}, events + ) + optimized_result, (optimized_reads,) = _run_and_get_profile_events( + node, select_expression, {"allow_experimental_iceberg_read_optimization": 1}, events + ) + + assert baseline_result == optimized_result == "3\t850" + assert baseline_reads > 0 + assert optimized_reads > 0 + + +def test_iceberg_read_optimization_non_identity_transform_keeps_prewhere_eligible(started_cluster): + node = started_cluster.instances["node1"] + catalog = load_catalog_impl(started_cluster) + + namespace = f"read_opt_nonident_ns_{uuid.uuid4().hex}" + table_name = "t" + catalog.create_namespace(namespace) + + bucket_on_id = PartitionSpec( + PartitionField(source_id=1, field_id=1000, transform=BucketTransform(4), name="id_bucket") + ) + table = create_table(catalog, namespace, table_name, SCHEMA, bucket_on_id) + + for id_ in range(1, 9): + append_row(table, id_, date(2024, 1, 1), id_ * 100) + + create_clickhouse_iceberg_database(node, CATALOG_NAME) + full_name = f"{CATALOG_NAME}.`{namespace}.{table_name}`" + + select_expression = f"SELECT id, value FROM {full_name} WHERE id = 5 ORDER BY ALL" + + result_baseline = node.query( + select_expression, settings={"allow_experimental_iceberg_read_optimization": 0} + ).strip() + result_optimized = node.query( + select_expression, settings={"allow_experimental_iceberg_read_optimization": 1} + ).strip() + + assert result_baseline == result_optimized == "5\t500"