diff --git a/src/Interpreters/InterpreterSelectQuery.cpp b/src/Interpreters/InterpreterSelectQuery.cpp index bc7119f0b132..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->supportedPrewhereColumns(); + 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 1b0267e8e084..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.supportedPrewhereColumns(), + 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 8d3cb20b8c58..6162f6aa6f7d 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 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/DataLakes/Iceberg/IcebergMetadata.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.cpp index 5fc41baf7ee2..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; } @@ -1464,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 5b2c64c07039..3324e9ec5015 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.h +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.h @@ -209,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/DataLakes/Iceberg/Utils.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.cpp index 8564eebda3ef..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) { @@ -1471,6 +1494,27 @@ std::optional getPartitionKeyStringFromMetadata(Poco::JSON::Object::Ptr return result; } +Names getIdentityPartitionColumnsFromMetadata(Poco::JSON::Object::Ptr metadata_object) +{ + auto resolved = resolveDefaultPartitionSpec(metadata_object); + if (!resolved) + return {}; + + Names result; + for (UInt32 i = 0; i < resolved->fields->size(); ++i) + { + auto field = resolved->fields->getObject(i); + + if (Poco::toLower(field->getValue(f_transform)) != "identity") + continue; + + 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; +} + 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..40656cdffbbc 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorage.cpp +++ b/src/Storages/ObjectStorage/StorageObjectStorage.cpp @@ -51,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 @@ -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); + return getInMemoryMetadataPtr()->getColumnsWithoutDefaultExpressions(/*exclude=*/hive_partition_columns_to_read_from_file_path); +} + +std::optional StorageObjectStorage::supportedAutomaticPrewhereColumns(const StorageMetadataPtr & metadata) const +{ + auto exclude = hive_partition_columns_to_read_from_file_path; + + for (const auto & identity_partition_column : metadata->identity_partition_columns) + { + if (metadata->getColumns().has(identity_partition_column)) + exclude.emplace_back(identity_partition_column, metadata->getColumns().get(identity_partition_column).type); + } + + LOG_DEBUG(log, "Automatic-PREWHERE exclude list: [{}]", exclude.toString()); + return metadata->getColumnsWithoutDefaultExpressions(/*exclude=*/exclude); } IStorage::ColumnSizeByName StorageObjectStorage::getColumnSizes() const diff --git a/src/Storages/ObjectStorage/StorageObjectStorage.h b/src/Storages/ObjectStorage/StorageObjectStorage.h index 6ace629b7085..e0d3f6f99ba1 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 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 bfb107c14d21..1247c7a320cb 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp +++ b/src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp @@ -1062,6 +1062,13 @@ std::optional StorageObjectStorageCluster::supportedPrewhereColumns() c return IStorageCluster::supportedPrewhereColumns(); } +std::optional StorageObjectStorageCluster::supportedAutomaticPrewhereColumns(const StorageMetadataPtr & metadata) const +{ + if (pure_storage) + return pure_storage->supportedAutomaticPrewhereColumns(metadata); + return IStorageCluster::supportedAutomaticPrewhereColumns(metadata); +} + IStorageCluster::ColumnSizeByName StorageObjectStorageCluster::getColumnSizes() const { if (pure_storage) diff --git a/src/Storages/ObjectStorage/StorageObjectStorageCluster.h b/src/Storages/ObjectStorage/StorageObjectStorageCluster.h index 91d07b5c498a..e0461c72eda3 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 StorageMetadataPtr & metadata) const override; ColumnSizeByName getColumnSizes() const override; bool parallelizeOutputAfterReading(ContextPtr context) const override; diff --git a/src/Storages/ObjectStorage/StorageObjectStorageConfiguration.h b/src/Storages/ObjectStorage/StorageObjectStorageConfiguration.h index 0ce750c6193a..708cdeb91f65 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorageConfiguration.h +++ b/src/Storages/ObjectStorage/StorageObjectStorageConfiguration.h @@ -312,7 +312,6 @@ class StorageObjectStorageConfiguration virtual ColumnMapperPtr getColumnMapperForCurrentSchema(StorageMetadataPtr /**/, ContextPtr /**/) const { return nullptr; } - virtual std::shared_ptr getCatalog(ContextPtr /*context*/, const StorageID & /*table_id*/) const { return nullptr; 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. 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"