From b3cc1f3e619c8c3be9519727dafa770023e56a9a Mon Sep 17 00:00:00 2001 From: Andrey Zvonov <32552679+zvonand@users.noreply.github.com> Date: Mon, 2 Feb 2026 07:14:54 +0100 Subject: [PATCH 1/6] Merge pull request #1345 from Altinity/backports/25.8.15/87303 25.8.15 Backport of #87303 - Fix condition not being moved to PREWHERE in case there is a row policy (version 2) --- src/Formats/FormatFactory.cpp | 5 +- src/Formats/FormatFilterInfo.cpp | 36 ++-- src/Formats/FormatFilterInfo.h | 10 +- src/Interpreters/ExpressionAnalyzer.cpp | 12 +- src/Interpreters/ExpressionAnalyzer.h | 6 +- src/Interpreters/InterpreterSelectQuery.cpp | 126 +++++--------- src/Interpreters/InterpreterSelectQuery.h | 2 +- .../getHeaderForProcessingStage.cpp | 2 +- src/Planner/PlannerJoinTree.cpp | 51 +----- .../Formats/Impl/Parquet/ReadManager.cpp | 2 +- .../Formats/Impl/Parquet/Reader.cpp | 42 +++-- .../optimizeLazyMaterialization.cpp | 8 +- .../Optimizations/optimizePrewhere.cpp | 9 +- .../optimizePrimaryKeyConditionAndLimit.cpp | 7 +- .../Optimizations/projectionsCommon.cpp | 17 +- .../QueryPlan/ReadFromMergeTree.cpp | 163 ++++++++++-------- .../QueryPlan/ReadFromObjectStorageStep.cpp | 12 +- .../QueryPlan/SourceStepWithFilter.cpp | 121 +++++++------ .../QueryPlan/SourceStepWithFilter.h | 8 +- src/Storages/IStorage.cpp | 5 - .../MergeTree/MergeTreeBlockReadUtils.cpp | 5 +- .../MergeTree/MergeTreeBlockReadUtils.h | 1 + .../MergeTree/MergeTreePrefetchedReadPool.cpp | 2 + .../MergeTree/MergeTreePrefetchedReadPool.h | 1 + src/Storages/MergeTree/MergeTreeReadPool.cpp | 2 + src/Storages/MergeTree/MergeTreeReadPool.h | 1 + .../MergeTree/MergeTreeReadPoolBase.cpp | 3 + .../MergeTree/MergeTreeReadPoolBase.h | 2 + .../MergeTree/MergeTreeReadPoolInOrder.cpp | 2 + .../MergeTree/MergeTreeReadPoolInOrder.h | 1 + .../MergeTreeReadPoolParallelReplicas.cpp | 2 + .../MergeTreeReadPoolParallelReplicas.h | 1 + ...rgeTreeReadPoolParallelReplicasInOrder.cpp | 2 + ...MergeTreeReadPoolParallelReplicasInOrder.h | 1 + .../MergeTree/MergeTreeSelectProcessor.cpp | 73 ++++---- .../MergeTree/MergeTreeSelectProcessor.h | 6 +- .../DataLakes/Iceberg/Compaction.cpp | 2 +- .../DataLakes/Iceberg/Mutations.cpp | 2 +- .../Iceberg/PositionDeleteTransform.cpp | 2 +- .../ObjectStorage/StorageObjectStorage.cpp | 6 +- .../StorageObjectStorageSource.cpp | 2 +- src/Storages/SelectQueryInfo.cpp | 64 +++---- src/Storages/SelectQueryInfo.h | 16 +- src/Storages/StorageBuffer.cpp | 62 ++++--- src/Storages/StorageDummy.cpp | 2 +- src/Storages/StorageFile.cpp | 10 +- src/Storages/StorageURL.cpp | 18 +- src/Storages/prepareReadingFromFormat.cpp | 8 +- src/Storages/prepareReadingFromFormat.h | 6 +- ...n_merge_tree_prewhere_row_policy.reference | 4 +- ...591_optimize_prewhere_row_policy.reference | 144 ++++++++++++++++ .../03591_optimize_prewhere_row_policy.sql | 34 ++++ .../03641_analyzer_issue_85834.reference | 1 + .../03641_analyzer_issue_85834.sql | 14 ++ 54 files changed, 676 insertions(+), 470 deletions(-) create mode 100644 tests/queries/0_stateless/03591_optimize_prewhere_row_policy.reference create mode 100644 tests/queries/0_stateless/03591_optimize_prewhere_row_policy.sql create mode 100644 tests/queries/0_stateless/03641_analyzer_issue_85834.reference create mode 100644 tests/queries/0_stateless/03641_analyzer_issue_85834.sql diff --git a/src/Formats/FormatFactory.cpp b/src/Formats/FormatFactory.cpp index 5d605eac01ea..e2e200342ffb 100644 --- a/src/Formats/FormatFactory.cpp +++ b/src/Formats/FormatFactory.cpp @@ -415,8 +415,9 @@ InputFormatPtr FormatFactory::getInput( const FormatSettings format_settings = _format_settings ? *_format_settings : getFormatSettings(context); const Settings & settings = context->getSettingsRef(); - if (format_filter_info && format_filter_info->prewhere_info && (!creators.random_access_input_creator || !creators.prewhere_support_checker || !creators.prewhere_support_checker(format_settings))) - throw Exception(ErrorCodes::LOGICAL_ERROR, "PREWHERE passed to format that doesn't support it"); + if (format_filter_info && (format_filter_info->prewhere_info || format_filter_info->row_level_filter) && (!creators.random_access_input_creator || !creators.prewhere_support_checker || !creators.prewhere_support_checker(format_settings))) + throw Exception(ErrorCodes::LOGICAL_ERROR, "{} passed to format that doesn't support it", + format_filter_info->prewhere_info ? "PREWHERE" : "ROW LEVEL FILTER"); if (!parser_shared_resources) parser_shared_resources = std::make_shared( diff --git a/src/Formats/FormatFilterInfo.cpp b/src/Formats/FormatFilterInfo.cpp index ac130771fbb9..32f197d79b16 100644 --- a/src/Formats/FormatFilterInfo.cpp +++ b/src/Formats/FormatFilterInfo.cpp @@ -48,19 +48,21 @@ std::pair, std::unordered_map return {clickhouse_to_parquet_names, parquet_names_to_clickhouse}; } - FormatFilterInfo::FormatFilterInfo(std::shared_ptr filter_actions_dag_, const ContextPtr & context_, ColumnMapperPtr column_mapper_) - : filter_actions_dag(filter_actions_dag_) - , context(context_) - , column_mapper(column_mapper_) - { - } +FormatFilterInfo::FormatFilterInfo( + std::shared_ptr filter_actions_dag_, + const ContextPtr & context_, + ColumnMapperPtr column_mapper_, + FilterDAGInfoPtr row_level_filter_, + PrewhereInfoPtr prewhere_info_) + : filter_actions_dag(filter_actions_dag_) + , context(context_) + , row_level_filter(std::move(row_level_filter_)) + , prewhere_info(std::move(prewhere_info_)) + , column_mapper(column_mapper_) +{ +} - FormatFilterInfo::FormatFilterInfo() - : filter_actions_dag(nullptr) - , context(static_cast(nullptr)) - , column_mapper(nullptr) - { - } +FormatFilterInfo::FormatFilterInfo() = default; bool FormatFilterInfo::hasFilter() const @@ -78,7 +80,7 @@ void FormatFilterInfo::initKeyCondition(const Block & keys) if (!ctx) throw Exception(ErrorCodes::LOGICAL_ERROR, "Context has expired"); - if (prewhere_info) + if (prewhere_info || row_level_filter) { auto add_columns = [&](const ActionsDAG & dag) { @@ -88,9 +90,11 @@ void FormatFilterInfo::initKeyCondition(const Block & keys) additional_columns.insert({col.type->createColumn(), col.type, col.name}); } }; - if (prewhere_info->row_level_filter.has_value()) - add_columns(prewhere_info->row_level_filter.value()); - add_columns(prewhere_info->prewhere_actions); + + if (row_level_filter) + add_columns(row_level_filter->actions); + if (prewhere_info) + add_columns(prewhere_info->prewhere_actions); } ColumnsWithTypeAndName columns = keys.getColumnsWithTypeAndName(); diff --git a/src/Formats/FormatFilterInfo.h b/src/Formats/FormatFilterInfo.h index 024c4e5321aa..712e21fbf695 100644 --- a/src/Formats/FormatFilterInfo.h +++ b/src/Formats/FormatFilterInfo.h @@ -12,6 +12,8 @@ struct Settings; class KeyCondition; struct PrewhereInfo; using PrewhereInfoPtr = std::shared_ptr; +struct FilterDAGInfo; +using FilterDAGInfoPtr = std::shared_ptr; /// Some formats needs to custom mapping between columns in file and clickhouse columns. class ColumnMapper @@ -47,12 +49,18 @@ using FormatFilterInfoPtr = std::shared_ptr; /// because most implementations don't use most of this struct. struct FormatFilterInfo { - FormatFilterInfo(std::shared_ptr filter_actions_dag_, const ContextPtr & context_, ColumnMapperPtr column_mapper_); + FormatFilterInfo( + std::shared_ptr filter_actions_dag_, + const ContextPtr & context_, + ColumnMapperPtr column_mapper_, + FilterDAGInfoPtr row_level_filter_, + PrewhereInfoPtr prewhere_info_); FormatFilterInfo(); std::shared_ptr filter_actions_dag; ContextWeakPtr context; // required only if `filter_actions_dag` is set + FilterDAGInfoPtr row_level_filter; PrewhereInfoPtr prewhere_info; // assigned only if the format supports prewhere /// Optionally created from filter_actions_dag, if the format needs it. diff --git a/src/Interpreters/ExpressionAnalyzer.cpp b/src/Interpreters/ExpressionAnalyzer.cpp index a74a0484d0eb..9fd9911b9cc2 100644 --- a/src/Interpreters/ExpressionAnalyzer.cpp +++ b/src/Interpreters/ExpressionAnalyzer.cpp @@ -1936,7 +1936,7 @@ ExpressionAnalysisResult::ExpressionAnalysisResult( bool first_stage_, bool second_stage_, bool only_types, - const FilterDAGInfoPtr & filter_info_, + const FilterDAGInfoPtr & row_policy_info_, const FilterDAGInfoPtr & additional_filter, const Block & source_header) : first_stage(first_stage_) @@ -2034,10 +2034,10 @@ ExpressionAnalysisResult::ExpressionAnalysisResult( columns_for_additional_filter.begin(), columns_for_additional_filter.end()); } - if (storage && filter_info_) + if (storage && row_policy_info_) { - filter_info = filter_info_; - filter_info->do_remove_column = true; + row_policy_info = row_policy_info_; + row_policy_info->do_remove_column = true; } if (prewhere_dag_and_flags = query_analyzer.appendPrewhere(chain, !first_stage); prewhere_dag_and_flags) @@ -2376,9 +2376,9 @@ std::string ExpressionAnalysisResult::dump() const ss << "prewhere_info " << prewhere_info->dump() << "\n"; } - if (filter_info) + if (row_policy_info) { - ss << "filter_info " << filter_info->dump() << "\n"; + ss << "filter_info " << row_policy_info->dump() << "\n"; } if (before_aggregation) diff --git a/src/Interpreters/ExpressionAnalyzer.h b/src/Interpreters/ExpressionAnalyzer.h index 9665bb1c32e1..e8eea1fbf635 100644 --- a/src/Interpreters/ExpressionAnalyzer.h +++ b/src/Interpreters/ExpressionAnalyzer.h @@ -273,7 +273,7 @@ struct ExpressionAnalysisResult NameSet columns_to_remove_after_prewhere; PrewhereInfoPtr prewhere_info; - FilterDAGInfoPtr filter_info; + FilterDAGInfoPtr row_policy_info; ConstantFilterDescription prewhere_constant_filter_description; ConstantFilterDescription where_constant_filter_description; /// Actions by every element of ORDER BY @@ -288,12 +288,12 @@ struct ExpressionAnalysisResult bool first_stage, bool second_stage, bool only_types, - const FilterDAGInfoPtr & filter_info, + const FilterDAGInfoPtr & row_policy_info, const FilterDAGInfoPtr & additional_filter, /// for setting additional_filters const Block & source_header); /// Filter for row-level security. - bool hasFilter() const { return filter_info.get(); } + bool hasRowPolicyFilter() const { return row_policy_info.get(); } bool hasJoin() const { return join.get(); } bool hasPrewhere() const { return prewhere_info.get(); } diff --git a/src/Interpreters/InterpreterSelectQuery.cpp b/src/Interpreters/InterpreterSelectQuery.cpp index d728868c2dea..929e69f57cbb 100644 --- a/src/Interpreters/InterpreterSelectQuery.cpp +++ b/src/Interpreters/InterpreterSelectQuery.cpp @@ -897,7 +897,7 @@ InterpreterSelectQuery::InterpreterSelectQuery( /// Fix source_header for filter actions. if (row_policy_filter && !row_policy_filter->empty()) { - filter_info = generateFilterActions( + row_policy_info = generateFilterActions( table_id, row_policy_filter->expression, context, storage, storage_snapshot, metadata_snapshot, required_columns, prepared_sets); @@ -1066,8 +1066,6 @@ bool InterpreterSelectQuery::adjustParallelReplicasAfterAnalysis() max_rows = max_rows ? std::min(max_rows, settings[Setting::max_rows_to_read].value) : settings[Setting::max_rows_to_read]; query_info_copy.trivial_limit = max_rows; - /// Apply filters to prewhere and add them to the query_info so we can filter out parts efficiently during row estimation - applyFiltersToPrewhereInAnalysis(analysis_copy); if (analysis_copy.prewhere_info) { query_info_copy.prewhere_info = analysis_copy.prewhere_info; @@ -1083,13 +1081,13 @@ bool InterpreterSelectQuery::adjustParallelReplicasAfterAnalysis() = query_info_copy.prewhere_info->prewhere_actions.findInOutputs(query_info_copy.prewhere_info->prewhere_column_name); added_filter_nodes.nodes.push_back(&node); } + } - if (query_info_copy.prewhere_info->row_level_filter) - { - const auto & node - = query_info_copy.prewhere_info->row_level_filter->findInOutputs(query_info_copy.prewhere_info->row_level_column_name); - added_filter_nodes.nodes.push_back(&node); - } + if (query_info_copy.row_level_filter) + { + const auto & node + = query_info_copy.row_level_filter->actions.findInOutputs(query_info_copy.row_level_filter->column_name); + added_filter_nodes.nodes.push_back(&node); } if (auto filter_actions_dag = ActionsDAG::buildFilterActionsDAG(added_filter_nodes.nodes)) @@ -1192,7 +1190,7 @@ Block InterpreterSelectQuery::getSampleBlockImpl() && options.to_stage > QueryProcessingStage::WithMergeableState; analysis_result = ExpressionAnalysisResult( - *query_analyzer, metadata_snapshot, first_stage, second_stage, options.only_analyze, filter_info, additional_filter_info, *source_header); + *query_analyzer, metadata_snapshot, first_stage, second_stage, options.only_analyze, row_policy_info, additional_filter_info, *source_header); if (options.to_stage == QueryProcessingStage::Enum::FetchColumns) { @@ -1635,13 +1633,13 @@ void InterpreterSelectQuery::executeImpl(QueryPlan & query_plan, std::optional

(source_header); query_plan.addStep(std::move(read_nothing)); - if (expressions.filter_info) + if (expressions.row_policy_info) { auto row_level_security_step = std::make_unique( query_plan.getCurrentHeader(), - expressions.filter_info->actions.clone(), - expressions.filter_info->column_name, - expressions.filter_info->do_remove_column); + expressions.row_policy_info->actions.clone(), + expressions.row_policy_info->column_name, + expressions.row_policy_info->do_remove_column); row_level_security_step->setStepDescription("Row-level security filter"); query_plan.addStep(std::move(row_level_security_step)); @@ -1649,18 +1647,6 @@ void InterpreterSelectQuery::executeImpl(QueryPlan & query_plan, std::optional

row_level_filter) - { - auto row_level_filter_step = std::make_unique( - query_plan.getCurrentHeader(), - expressions.prewhere_info->row_level_filter->clone(), - expressions.prewhere_info->row_level_column_name, - true); - - row_level_filter_step->setStepDescription("Row-level security filter (PREWHERE)"); - query_plan.addStep(std::move(row_level_filter_step)); - } - auto prewhere_step = std::make_unique( query_plan.getCurrentHeader(), expressions.prewhere_info->prewhere_actions.clone(), @@ -1762,13 +1748,13 @@ void InterpreterSelectQuery::executeImpl(QueryPlan & query_plan, std::optional

supportsPrewhere())) { auto row_level_security_step = std::make_unique( query_plan.getCurrentHeader(), - expressions.filter_info->actions.clone(), - expressions.filter_info->column_name, - expressions.filter_info->do_remove_column); + expressions.row_policy_info->actions.clone(), + expressions.row_policy_info->column_name, + expressions.row_policy_info->do_remove_column); row_level_security_step->setStepDescription("Row-level security filter"); query_plan.addStep(std::move(row_level_security_step)); @@ -2225,21 +2211,21 @@ void InterpreterSelectQuery::addEmptySourceToQueryPlan(QueryPlan & query_plan, c { Pipe pipe(std::make_shared(std::make_shared(source_header))); - if (query_info.prewhere_info) + if (query_info.row_level_filter) { - auto & prewhere_info = *query_info.prewhere_info; - - if (prewhere_info.row_level_filter) + auto row_level_actions = std::make_shared(query_info.row_level_filter->actions.clone()); + pipe.addSimpleTransform([&](const SharedHeader & header) { - auto row_level_actions = std::make_shared(prewhere_info.row_level_filter->clone()); - pipe.addSimpleTransform([&](const SharedHeader & header) - { - return std::make_shared(header, - row_level_actions, - prewhere_info.row_level_column_name, true); - }); - } + return std::make_shared(header, + row_level_actions, + query_info.row_level_filter->column_name, + query_info.row_level_filter->do_remove_column); + }); + } + if (query_info.prewhere_info) + { + auto & prewhere_info = *query_info.prewhere_info; auto filter_actions = std::make_shared(prewhere_info.prewhere_actions.clone()); pipe.addSimpleTransform([&](const SharedHeader & header) { @@ -2278,38 +2264,9 @@ bool InterpreterSelectQuery::shouldMoveToPrewhere() const return settings[Setting::optimize_move_to_prewhere] && (!query.final() || settings[Setting::optimize_move_to_prewhere_if_final]); } -/// Note that this is const and accepts the analysis ref to be able to use it to do analysis for parallel replicas -/// without affecting the final analysis multiple times -void InterpreterSelectQuery::applyFiltersToPrewhereInAnalysis(ExpressionAnalysisResult & analysis) const -{ - if (!analysis.filter_info) - return; - - if (!analysis.prewhere_info) - { - const bool does_storage_support_prewhere = !input_pipe && storage && storage->supportsPrewhere(); - if (does_storage_support_prewhere && shouldMoveToPrewhere()) - { - /// Execute row level filter in prewhere as a part of "move to prewhere" optimization. - analysis.prewhere_info = std::make_shared(std::move(analysis.filter_info->actions), analysis.filter_info->column_name); - analysis.prewhere_info->remove_prewhere_column = std::move(analysis.filter_info->do_remove_column); - analysis.prewhere_info->need_filter = true; - analysis.filter_info = nullptr; - } - } - else - { - /// Add row level security actions to prewhere. - analysis.prewhere_info->row_level_filter = std::move(analysis.filter_info->actions); - analysis.prewhere_info->row_level_column_name = std::move(analysis.filter_info->column_name); - analysis.filter_info = nullptr; - } -} - - void InterpreterSelectQuery::addPrewhereAliasActions() { - applyFiltersToPrewhereInAnalysis(analysis_result); + auto & row_level_filter = analysis_result.row_policy_info; auto & prewhere_info = analysis_result.prewhere_info; auto & columns_to_remove_after_prewhere = analysis_result.columns_to_remove_after_prewhere; @@ -2336,12 +2293,12 @@ void InterpreterSelectQuery::addPrewhereAliasActions() /// Get some columns directly from PREWHERE expression actions auto prewhere_required_columns = prewhere_info->prewhere_actions.getRequiredColumns().getNames(); columns.insert(prewhere_required_columns.begin(), prewhere_required_columns.end()); + } - if (prewhere_info->row_level_filter) - { - auto row_level_required_columns = prewhere_info->row_level_filter->getRequiredColumns().getNames(); - columns.insert(row_level_required_columns.begin(), row_level_required_columns.end()); - } + if (row_level_filter) + { + auto row_level_required_columns = row_level_filter->actions.getRequiredColumns().getNames(); + columns.insert(row_level_required_columns.begin(), row_level_required_columns.end()); } return columns; @@ -2499,13 +2456,15 @@ std::optional InterpreterSelectQuery::getTrivialCount(UInt64 allow_exper // It's possible to optimize count() given only partition predicates ActionsDAG::NodeRawConstPtrs filter_nodes; + if (analysis_result.hasRowPolicyFilter()) + { + auto & row_level_filter = analysis_result.row_policy_info; + filter_nodes.push_back(&row_level_filter->actions.findInOutputs(row_level_filter->column_name)); + } if (analysis_result.hasPrewhere()) { auto & prewhere_info = analysis_result.prewhere_info; filter_nodes.push_back(&prewhere_info->prewhere_actions.findInOutputs(prewhere_info->prewhere_column_name)); - - if (prewhere_info->row_level_filter) - filter_nodes.push_back(&prewhere_info->row_level_filter->findInOutputs(prewhere_info->row_level_column_name)); } if (analysis_result.hasWhere()) { @@ -2696,10 +2655,11 @@ void InterpreterSelectQuery::executeFetchColumns(QueryProcessingStage::Enum proc if (max_streams == 0) max_streams = 1; - auto & prewhere_info = analysis_result.prewhere_info; + if (analysis_result.row_policy_info && (!input_pipe && storage && storage->supportsPrewhere())) + query_info.row_level_filter = analysis_result.row_policy_info; - if (prewhere_info) - query_info.prewhere_info = prewhere_info; + if (analysis_result.prewhere_info) + query_info.prewhere_info = analysis_result.prewhere_info; bool optimize_read_in_order = analysis_result.optimize_read_in_order; bool optimize_aggregation_in_order = analysis_result.optimize_read_in_order && !query_analyzer->useGroupingSetKey(); diff --git a/src/Interpreters/InterpreterSelectQuery.h b/src/Interpreters/InterpreterSelectQuery.h index af85ac4b9f30..18084cc43ee0 100644 --- a/src/Interpreters/InterpreterSelectQuery.h +++ b/src/Interpreters/InterpreterSelectQuery.h @@ -220,7 +220,7 @@ class InterpreterSelectQuery : public IInterpreterUnionOrSelectQuery ExpressionAnalysisResult analysis_result; /// For row-level security. RowPolicyFilterPtr row_policy_filter; - FilterDAGInfoPtr filter_info; + FilterDAGInfoPtr row_policy_info; /// For additional_filter setting. FilterDAGInfoPtr additional_filter_info; diff --git a/src/Interpreters/getHeaderForProcessingStage.cpp b/src/Interpreters/getHeaderForProcessingStage.cpp index c686e3f042dd..bb1eec0ba044 100644 --- a/src/Interpreters/getHeaderForProcessingStage.cpp +++ b/src/Interpreters/getHeaderForProcessingStage.cpp @@ -104,7 +104,7 @@ SharedHeader getHeaderForProcessingStage( case QueryProcessingStage::FetchColumns: { Block header = storage_snapshot->getSampleBlockForColumns(column_names); - header = SourceStepWithFilter::applyPrewhereActions(header, query_info.prewhere_info); + header = SourceStepWithFilter::applyPrewhereActions(header, query_info.row_level_filter, query_info.prewhere_info); return std::make_shared(std::move(header)); } case QueryProcessingStage::WithMergeableState: diff --git a/src/Planner/PlannerJoinTree.cpp b/src/Planner/PlannerJoinTree.cpp index ed54b80de514..1488d1ea5118 100644 --- a/src/Planner/PlannerJoinTree.cpp +++ b/src/Planner/PlannerJoinTree.cpp @@ -913,6 +913,7 @@ JoinTreeQueryPlan buildQueryPlanForTableExpression(QueryTreeNodePtr table_expres { if (!select_query_options.only_analyze) { + auto & row_level_filter = table_expression_query_info.row_level_filter; auto & prewhere_info = table_expression_query_info.prewhere_info; const auto & prewhere_actions = table_expression_data.getPrewhereFilterActions(); const auto & columns_names = table_expression_data.getColumnNames(); @@ -979,52 +980,16 @@ JoinTreeQueryPlan buildQueryPlanForTableExpression(QueryTreeNodePtr table_expres updatePrewhereOutputsIfNeeded(table_expression_query_info, table_expression_data.getColumnNames(), storage_snapshot); - const auto add_filter = [&](FilterDAGInfo & filter_info, std::string description) - { - bool is_final = table_expression_query_info.table_expression_modifiers - && table_expression_query_info.table_expression_modifiers->hasFinal(); - bool optimize_move_to_prewhere - = settings[Setting::optimize_move_to_prewhere] && (!is_final || settings[Setting::optimize_move_to_prewhere_if_final]); - - auto supported_prewhere_columns = storage->supportedPrewhereColumns(); - bool has_table_virtual_column = - filter_info.column_name == "_table" && storage->isVirtualColumn(filter_info.column_name, storage_snapshot->metadata); - if (!select_query_options.build_logical_plan && storage->canMoveConditionsToPrewhere() && optimize_move_to_prewhere - && (!supported_prewhere_columns || supported_prewhere_columns->contains(filter_info.column_name)) - && !has_table_virtual_column) - { - if (!prewhere_info) - { - prewhere_info = std::make_shared(); - prewhere_info->prewhere_actions = std::move(filter_info.actions); - prewhere_info->prewhere_column_name = filter_info.column_name; - prewhere_info->remove_prewhere_column = filter_info.do_remove_column; - prewhere_info->need_filter = true; - } - else if (!prewhere_info->row_level_filter) - { - prewhere_info->row_level_filter = std::move(filter_info.actions); - prewhere_info->row_level_column_name = filter_info.column_name; - prewhere_info->need_filter = true; - } - else - { - where_filters.emplace_back(std::move(filter_info), std::move(description)); - } - - } - else - { - where_filters.emplace_back(std::move(filter_info), std::move(description)); - } - }; - auto row_policy_filter_info = buildRowPolicyFilterIfNeeded(storage, table_expression_query_info, planner_context, used_row_policies); if (row_policy_filter_info) { table_expression_data.setRowLevelFilterActions(row_policy_filter_info->actions.clone()); - add_filter(*row_policy_filter_info, "Row-level security filter"); + /// TODO: Never put row-level security filter in WHERE clause for storages that do not support PREWHERE to avoid merging of filters. + if (storage->supportsPrewhere()) + row_level_filter = std::make_shared(std::move(*row_policy_filter_info)); + else + where_filters.emplace_back(std::move(*row_policy_filter_info), "Row-level security filter"); } if (query_context->canUseParallelReplicasCustomKey()) @@ -1032,7 +997,7 @@ JoinTreeQueryPlan buildQueryPlanForTableExpression(QueryTreeNodePtr table_expres if (settings[Setting::parallel_replicas_count] > 1) { if (auto parallel_replicas_custom_key_filter_info= buildCustomKeyFilterIfNeeded(storage, table_expression_query_info, planner_context)) - add_filter(*parallel_replicas_custom_key_filter_info, "Parallel replicas custom key filter"); + where_filters.emplace_back(std::move(*parallel_replicas_custom_key_filter_info), "Parallel replicas custom key filter"); } else if (auto * distributed = typeid_cast(storage.get()); distributed && query_context->canUseParallelReplicasCustomKeyForCluster(*distributed->getCluster())) @@ -1049,7 +1014,7 @@ JoinTreeQueryPlan buildQueryPlanForTableExpression(QueryTreeNodePtr table_expres if (auto additional_filters_info = buildAdditionalFiltersIfNeeded(table_expression_query_info, prewhere_info, planner_context)) { appendSetsFromActionsDAG(additional_filters_info->actions, useful_sets); - add_filter(*additional_filters_info, "additional filter"); + where_filters.emplace_back(std::move(*additional_filters_info), "additional filter"); } if (!select_query_options.build_logical_plan) diff --git a/src/Processors/Formats/Impl/Parquet/ReadManager.cpp b/src/Processors/Formats/Impl/Parquet/ReadManager.cpp index f6cd46092594..65448a8ec2b7 100644 --- a/src/Processors/Formats/Impl/Parquet/ReadManager.cpp +++ b/src/Processors/Formats/Impl/Parquet/ReadManager.cpp @@ -73,7 +73,7 @@ void ReadManager::init(FormatParserSharedResourcesPtr parser_shared_resources_, /// eat all memory, and MainData would have to execute in one thread to minimize memory usage. double sum = 0; stages[size_t(ReadStage::MainData)].memory_target_fraction *= 10; - if (reader.format_filter_info->prewhere_info) + if (reader.format_filter_info->prewhere_info || reader.format_filter_info->row_level_filter) stages[size_t(ReadStage::PrewhereData)].memory_target_fraction *= 5; else { diff --git a/src/Processors/Formats/Impl/Parquet/Reader.cpp b/src/Processors/Formats/Impl/Parquet/Reader.cpp index 01cd1fedc8c5..33b822f9e1f6 100644 --- a/src/Processors/Formats/Impl/Parquet/Reader.cpp +++ b/src/Processors/Formats/Impl/Parquet/Reader.cpp @@ -318,10 +318,13 @@ void Reader::prefilterAndInitRowGroups(const std::optionaladditional_columns) extended_sample_block.insert(col); extended_sample_block_data_types = extended_sample_block.getDataTypes(); - PrewhereInfoPtr prewhere_info = format_filter_info->prewhere_info; + const auto & row_level_filter = format_filter_info->row_level_filter; + const auto & prewhere_info = format_filter_info->prewhere_info; /// Process schema. SchemaConverter schemer(file_metadata, options, &extended_sample_block); + if (row_level_filter && !row_level_filter->do_remove_column) + schemer.external_columns.push_back(row_level_filter->column_name); if (prewhere_info && !prewhere_info->remove_prewhere_column) schemer.external_columns.push_back(prewhere_info->prewhere_column_name); schemer.column_mapper = format_filter_info->column_mapper.get(); @@ -513,7 +516,7 @@ void Reader::prepareBloomFilterCondition() void Reader::initializePrefetches() { - bool use_offset_index = options.format.parquet.use_offset_index || format_filter_info->prewhere_info + bool use_offset_index = options.format.parquet.use_offset_index || format_filter_info->prewhere_info || format_filter_info->row_level_filter || std::any_of(primitive_columns.begin(), primitive_columns.end(), [](const auto & c) { return c.column_index_condition; }); bool need_to_find_bloom_filter_lengths_the_hard_way = false; @@ -659,8 +662,9 @@ void Reader::initializePrefetches() void Reader::preparePrewhere() { - PrewhereInfoPtr prewhere_info = format_filter_info->prewhere_info; - if (prewhere_info) + const auto & row_level_filter = format_filter_info->row_level_filter; + const auto & prewhere_info = format_filter_info->prewhere_info; + if (row_level_filter || prewhere_info) { /// TODO [parquet]: We currently run prewhere after reading all prewhere columns of the row /// subgroup, in one thread per row group. Instead, we could extract single-column conditions @@ -671,24 +675,30 @@ void Reader::preparePrewhere() /// Convert ActionsDAG to ExpressionActions. ExpressionActionsSettings actions_settings; - if (prewhere_info->row_level_filter.has_value()) + if (row_level_filter) { - ExpressionActions actions(prewhere_info->row_level_filter->clone(), actions_settings); + ExpressionActions actions(row_level_filter->actions.clone(), actions_settings); prewhere_steps.push_back(PrewhereStep { .actions = std::move(actions), - .result_column_name = prewhere_info->row_level_column_name, + .result_column_name = row_level_filter->column_name, }); + + if (!row_level_filter->do_remove_column) + prewhere_steps.back().idx_in_output_block = sample_block->getPositionByName(row_level_filter->column_name); + } + if (prewhere_info) + { + ExpressionActions actions(prewhere_info->prewhere_actions.clone(), actions_settings); + prewhere_steps.push_back(PrewhereStep + { + .actions = std::move(actions), + .result_column_name = prewhere_info->prewhere_column_name, + .need_filter = prewhere_info->need_filter, + }); + if (!prewhere_info->remove_prewhere_column) + prewhere_steps.back().idx_in_output_block = sample_block->getPositionByName(prewhere_info->prewhere_column_name); } - ExpressionActions actions(prewhere_info->prewhere_actions.clone(), actions_settings); - prewhere_steps.push_back(PrewhereStep - { - .actions = std::move(actions), - .result_column_name = prewhere_info->prewhere_column_name, - .need_filter = prewhere_info->need_filter, - }); - if (!prewhere_info->remove_prewhere_column) - prewhere_steps.back().idx_in_output_block = sample_block->getPositionByName(prewhere_info->prewhere_column_name); } /// Look up expression inputs in extended_sample_block. for (PrewhereStep & step : prewhere_steps) diff --git a/src/Processors/QueryPlan/Optimizations/optimizeLazyMaterialization.cpp b/src/Processors/QueryPlan/Optimizations/optimizeLazyMaterialization.cpp index 950ac49c1055..a2f5275fd0f1 100644 --- a/src/Processors/QueryPlan/Optimizations/optimizeLazyMaterialization.cpp +++ b/src/Processors/QueryPlan/Optimizations/optimizeLazyMaterialization.cpp @@ -131,13 +131,11 @@ static void collectLazilyReadColumnNames( for (const auto & column_name : lazily_read_column_name_set) alias_index.emplace(column_name, column_name); - if (const auto & prewhere_info = read_from_merge_tree->getPrewhereInfo()) - { - if (prewhere_info->row_level_filter) - removeUsedColumnNames(*prewhere_info->row_level_filter, lazily_read_column_name_set, alias_index, prewhere_info->row_level_column_name); + if (const auto & row_level_filter = read_from_merge_tree->getRowLevelFilter()) + removeUsedColumnNames(row_level_filter->actions, lazily_read_column_name_set, alias_index, row_level_filter->column_name); + if (const auto & prewhere_info = read_from_merge_tree->getPrewhereInfo()) removeUsedColumnNames(prewhere_info->prewhere_actions, lazily_read_column_name_set, alias_index, prewhere_info->prewhere_column_name); - } for (auto step_it = steps.rbegin(); step_it != steps.rend(); ++step_it) { diff --git a/src/Processors/QueryPlan/Optimizations/optimizePrewhere.cpp b/src/Processors/QueryPlan/Optimizations/optimizePrewhere.cpp index 70eb278f2332..38b5becae90a 100644 --- a/src/Processors/QueryPlan/Optimizations/optimizePrewhere.cpp +++ b/src/Processors/QueryPlan/Optimizations/optimizePrewhere.cpp @@ -140,8 +140,7 @@ void optimizePrewhere(Stack & stack, QueryPlan::Nodes &) if (!storage.canMoveConditionsToPrewhere()) return; - const auto & storage_prewhere_info = source_step_with_filter->getPrewhereInfo(); - if (storage_prewhere_info) + if (source_step_with_filter->getPrewhereInfo()) return; /// TODO: We can also check for UnionStep, such as StorageBuffer and local distributed plans. @@ -196,11 +195,7 @@ void optimizePrewhere(Stack & stack, QueryPlan::Nodes &) if (optimize_result.prewhere_nodes.empty()) return; - PrewhereInfoPtr prewhere_info; - if (storage_prewhere_info) - prewhere_info = storage_prewhere_info->clone(); - else - prewhere_info = std::make_shared(); + PrewhereInfoPtr prewhere_info = std::make_shared(); auto remaining_expr = splitAndFillPrewhereInfo( prewhere_info, diff --git a/src/Processors/QueryPlan/Optimizations/optimizePrimaryKeyConditionAndLimit.cpp b/src/Processors/QueryPlan/Optimizations/optimizePrimaryKeyConditionAndLimit.cpp index 33408e02df87..850451c39a0a 100644 --- a/src/Processors/QueryPlan/Optimizations/optimizePrimaryKeyConditionAndLimit.cpp +++ b/src/Processors/QueryPlan/Optimizations/optimizePrimaryKeyConditionAndLimit.cpp @@ -17,12 +17,11 @@ void optimizePrimaryKeyConditionAndLimit(const Stack & stack) return; const auto & storage_prewhere_info = source_step_with_filter->getPrewhereInfo(); + const auto & storage_row_level_filter = source_step_with_filter->getRowLevelFilter(); + if (storage_row_level_filter) + source_step_with_filter->addFilter(storage_row_level_filter->actions.clone(), storage_row_level_filter->column_name); if (storage_prewhere_info) - { source_step_with_filter->addFilter(storage_prewhere_info->prewhere_actions.clone(), storage_prewhere_info->prewhere_column_name); - if (storage_prewhere_info->row_level_filter) - source_step_with_filter->addFilter(storage_prewhere_info->row_level_filter->clone(), storage_prewhere_info->row_level_column_name); - } for (auto iter = stack.rbegin() + 1; iter != stack.rend(); ++iter) { diff --git a/src/Processors/QueryPlan/Optimizations/projectionsCommon.cpp b/src/Processors/QueryPlan/Optimizations/projectionsCommon.cpp index 666c5609c5b0..31f4aa6ce9bb 100644 --- a/src/Processors/QueryPlan/Optimizations/projectionsCommon.cpp +++ b/src/Processors/QueryPlan/Optimizations/projectionsCommon.cpp @@ -151,17 +151,16 @@ bool QueryDAG::buildImpl(QueryPlan::Node & node, ActionsDAG::NodeRawConstPtrs & IQueryPlanStep * step = node.step.get(); if (auto * reading = typeid_cast(step)) { + if (const auto & row_level_filter = reading->getRowLevelFilter()) + { + appendExpression(row_level_filter->actions); + if (const auto * filter_expression = findInOutputs(*dag, row_level_filter->column_name, row_level_filter->do_remove_column)) + filter_nodes.push_back(filter_expression); + else + return false; + } if (const auto & prewhere_info = reading->getPrewhereInfo()) { - if (prewhere_info->row_level_filter) - { - appendExpression(*prewhere_info->row_level_filter); - if (const auto * filter_expression = findInOutputs(*dag, prewhere_info->row_level_column_name, false)) - filter_nodes.push_back(filter_expression); - else - return false; - } - appendExpression(prewhere_info->prewhere_actions); if (const auto * filter_expression = findInOutputs(*dag, prewhere_info->prewhere_column_name, prewhere_info->remove_prewhere_column)) diff --git a/src/Processors/QueryPlan/ReadFromMergeTree.cpp b/src/Processors/QueryPlan/ReadFromMergeTree.cpp index 309fcf1beeed..a34a54997d1b 100644 --- a/src/Processors/QueryPlan/ReadFromMergeTree.cpp +++ b/src/Processors/QueryPlan/ReadFromMergeTree.cpp @@ -106,13 +106,14 @@ bool restoreDAGInputs(ActionsDAG & dag, const NameSet & inputs) return added; } -bool restorePrewhereInputs(PrewhereInfo & info, const NameSet & inputs) +bool restorePrewhereInputs(FilterDAGInfo * row_level_filter, PrewhereInfo * info, const NameSet & inputs) { bool added = false; - if (info.row_level_filter) - added = added || restoreDAGInputs(*info.row_level_filter, inputs); + if (row_level_filter) + added = added || restoreDAGInputs(row_level_filter->actions, inputs); - added = added || restoreDAGInputs(info.prewhere_actions, inputs); + if (info) + added = added || restoreDAGInputs(info->prewhere_actions, inputs); return added; } @@ -212,7 +213,8 @@ static SortDescription getSortDescriptionForOutputHeader( const std::vector & reverse_flags, const int sort_direction, InputOrderInfoPtr input_order_info, - PrewhereInfoPtr prewhere_info, + const FilterDAGInfoPtr & row_level_filter, + const PrewhereInfoPtr & prewhere_info, bool enable_vertical_final) { /// Updating sort description can be done after PREWHERE actions are applied to the header. @@ -231,16 +233,16 @@ static SortDescription getSortDescriptionForOutputHeader( column.name = original_node->result_name; } } + } - if (prewhere_info->row_level_filter) + if (row_level_filter) + { + FindOriginalNodeForOutputName original_column_finder(row_level_filter->actions); + for (auto & column : original_header) { - FindOriginalNodeForOutputName original_column_finder(*prewhere_info->row_level_filter); - for (auto & column : original_header) - { - const auto * original_node = original_column_finder.find(column.name); - if (original_node) - column.name = original_node->result_name; - } + const auto * original_node = original_column_finder.find(column.name); + if (original_node) + column.name = original_node->result_name; } } @@ -327,6 +329,7 @@ ReadFromMergeTree::ReadFromMergeTree( : SourceStepWithFilter(std::make_shared(MergeTreeSelectProcessor::transformHeader( storage_snapshot_->getSampleBlockForColumns(all_column_names_), {}, + query_info_.row_level_filter, query_info_.prewhere_info)), all_column_names_, query_info_, storage_snapshot_, context_) , reader_settings(MergeTreeReaderSettings::create(context_, query_info_)) , prepared_parts(std::move(parts_)) @@ -426,7 +429,8 @@ Pipe ReadFromMergeTree::readFromPoolParallelReplicas(RangesInDataParts parts_wit mutations_snapshot, shared_virtual_fields, storage_snapshot, - prewhere_info, + query_info.row_level_filter, + query_info.prewhere_info, actions_settings, reader_settings, required_columns, @@ -441,7 +445,7 @@ Pipe ReadFromMergeTree::readFromPoolParallelReplicas(RangesInDataParts parts_wit auto algorithm = std::make_unique(i); auto processor = std::make_unique( - pool, std::move(algorithm), prewhere_info, lazily_read_info, actions_settings, reader_settings); + pool, std::move(algorithm), query_info.row_level_filter, query_info.prewhere_info, lazily_read_info, actions_settings, reader_settings); auto source = std::make_shared(std::move(processor), data.getLogName()); pipes.emplace_back(std::move(source)); @@ -503,7 +507,8 @@ Pipe ReadFromMergeTree::readFromPool( mutations_snapshot, shared_virtual_fields, storage_snapshot, - prewhere_info, + query_info.row_level_filter, + query_info.prewhere_info, actions_settings, reader_settings, required_columns, @@ -518,7 +523,8 @@ Pipe ReadFromMergeTree::readFromPool( mutations_snapshot, shared_virtual_fields, storage_snapshot, - prewhere_info, + query_info.row_level_filter, + query_info.prewhere_info, actions_settings, reader_settings, required_columns, @@ -535,7 +541,7 @@ Pipe ReadFromMergeTree::readFromPool( auto algorithm = std::make_unique(i); auto processor - = std::make_unique(pool, std::move(algorithm), prewhere_info, lazily_read_info, actions_settings, reader_settings); + = std::make_unique(pool, std::move(algorithm), query_info.row_level_filter, query_info.prewhere_info, lazily_read_info, actions_settings, reader_settings); auto source = std::make_shared(std::move(processor), data.getLogName()); @@ -583,7 +589,8 @@ Pipe ReadFromMergeTree::readInOrder( shared_virtual_fields, has_limit_below_one_block, storage_snapshot, - prewhere_info, + query_info.row_level_filter, + query_info.prewhere_info, actions_settings, reader_settings, required_columns, @@ -600,7 +607,8 @@ Pipe ReadFromMergeTree::readInOrder( mutations_snapshot, shared_virtual_fields, storage_snapshot, - prewhere_info, + query_info.row_level_filter, + query_info.prewhere_info, actions_settings, reader_settings, required_columns, @@ -638,7 +646,7 @@ Pipe ReadFromMergeTree::readInOrder( algorithm = std::make_unique(i); auto processor = std::make_unique( - pool, std::move(algorithm), prewhere_info, lazily_read_info, actions_settings, reader_settings); + pool, std::move(algorithm), query_info.row_level_filter, query_info.prewhere_info, lazily_read_info, actions_settings, reader_settings); processor->addPartLevelToChunk(isQueryWithFinal()); @@ -862,6 +870,7 @@ Pipe ReadFromMergeTree::readByLayers(const RangesInDataParts & parts_with_ranges auto header = std::make_shared(MergeTreeSelectProcessor::transformHeader( storage_snapshot->getSampleBlockForColumns(in_order_column_names_to_read), lazily_read_info, + query_info.row_level_filter, query_info.prewhere_info)); pipe = Pipe(std::make_shared(header)); } @@ -965,7 +974,8 @@ Pipe ReadFromMergeTree::spreadMarkRangesAmongStreams(RangesInDataParts && parts_ fault(thread_local_rng) && !isQueryWithFinal() && data.merging_params.is_deleted_column.empty() && - !prewhere_info && + !query_info.row_level_filter && + !query_info.prewhere_info && !lazily_read_info && !reader_settings.use_query_condition_cache && /// the query condition cache produces incorrect results with intersecting ranges !isVectorColumnReplaced()) /// Vector search optimization needs ranges & offsets to be stable @@ -1073,13 +1083,13 @@ Pipe ReadFromMergeTree::spreadMarkRangesAmongStreamsWithOrder( /// To fix this, we prohibit removing any input in prewhere actions. Instead, projection actions will be added after sorting. /// See 02354_read_in_order_prewhere.sql as an example. bool have_input_columns_removed_after_prewhere = false; - if (prewhere_info) + if (query_info.prewhere_info || query_info.row_level_filter) { NameSet sorting_columns; for (const auto & column : storage_snapshot->metadata->getSortingKey().expression->getRequiredColumnsWithTypes()) sorting_columns.insert(column.name); - have_input_columns_removed_after_prewhere = restorePrewhereInputs(*prewhere_info, sorting_columns); + have_input_columns_removed_after_prewhere = restorePrewhereInputs(query_info.row_level_filter.get(), query_info.prewhere_info.get(), sorting_columns); } /// Let's split ranges to avoid reading much data. @@ -1479,12 +1489,12 @@ Pipe ReadFromMergeTree::spreadMarkRangesAmongStreamsFinal( auto sorting_expr = storage_snapshot->metadata->getSortingKey().expression; - if (prewhere_info) + if (query_info.prewhere_info || query_info.row_level_filter) { NameSet sorting_columns; for (const auto & column : storage_snapshot->metadata->getSortingKey().expression->getRequiredColumnsWithTypes()) sorting_columns.insert(column.name); - restorePrewhereInputs(*prewhere_info, sorting_columns); + restorePrewhereInputs(query_info.row_level_filter.get(), query_info.prewhere_info.get(), sorting_columns); } for (size_t range_index = 0; range_index < parts_to_merge_ranges.size() - 1; ++range_index) @@ -2072,7 +2082,8 @@ void ReadFromMergeTree::updateSortDescription() storage_snapshot->metadata->getSortingKeyReverseFlags(), getSortDirection(), query_info.input_order_info, - prewhere_info, + query_info.row_level_filter, + query_info.prewhere_info, enable_vertical_final); } @@ -2122,11 +2133,11 @@ bool ReadFromMergeTree::readsInOrder() const void ReadFromMergeTree::updatePrewhereInfo(const PrewhereInfoPtr & prewhere_info_value) { query_info.prewhere_info = prewhere_info_value; - prewhere_info = prewhere_info_value; output_header = std::make_shared(MergeTreeSelectProcessor::transformHeader( storage_snapshot->getSampleBlockForColumns(all_column_names), lazily_read_info, + query_info.row_level_filter, prewhere_info_value)); updateSortDescription(); @@ -2157,7 +2168,8 @@ void ReadFromMergeTree::updateLazilyReadInfo(const LazilyReadInfoPtr & lazily_re output_header = std::make_shared(MergeTreeSelectProcessor::transformHeader( storage_snapshot->getSampleBlockForColumns(all_column_names), lazily_read_info, - prewhere_info)); + query_info.row_level_filter, + query_info.prewhere_info)); /// if analysis has already been done (like in optimization for projections), /// then update columns to read in analysis result @@ -2174,7 +2186,8 @@ void ReadFromMergeTree::replaceVectorColumnWithDistanceColumn(const String & vec output_header = std::make_shared(MergeTreeSelectProcessor::transformHeader( storage_snapshot->getSampleBlockForColumns(all_column_names), lazily_read_info, - prewhere_info)); + query_info.row_level_filter, + query_info.prewhere_info)); /// if analysis has already been done (like in optimization for projections), /// then update columns to read in analysis result @@ -2294,8 +2307,8 @@ Pipe ReadFromMergeTree::spreadMarkRanges( sampling_columns.insert(column); } - if (prewhere_info) - restorePrewhereInputs(*prewhere_info, sampling_columns); + if (query_info.prewhere_info || query_info.row_level_filter) + restorePrewhereInputs(query_info.row_level_filter.get(), query_info.prewhere_info.get(), sampling_columns); } if (final) @@ -2594,33 +2607,38 @@ void ReadFromMergeTree::describeActions(FormatSettings & format_settings) const format_settings.out << prefix << "Granules: " << result.index_stats.back().num_granules_after << '\n'; } - if (prewhere_info) + if (query_info.prewhere_info || query_info.row_level_filter) { format_settings.out << prefix << "Prewhere info" << '\n'; - format_settings.out << prefix << "Need filter: " << prewhere_info->need_filter << '\n'; + if (query_info.prewhere_info) + format_settings.out << prefix << "Need filter: " << query_info.prewhere_info->need_filter << '\n'; prefix.push_back(format_settings.indent_char); prefix.push_back(format_settings.indent_char); + } - { - format_settings.out << prefix << "Prewhere filter" << '\n'; - format_settings.out << prefix << "Prewhere filter column: " << prewhere_info->prewhere_column_name; - if (prewhere_info->remove_prewhere_column) - format_settings.out << " (removed)"; - format_settings.out << '\n'; + if (query_info.prewhere_info) + { + format_settings.out << prefix << "Prewhere filter" << '\n'; + format_settings.out << prefix << "Prewhere filter column: " << query_info.prewhere_info->prewhere_column_name; + if (query_info.prewhere_info->remove_prewhere_column) + format_settings.out << " (removed)"; + format_settings.out << '\n'; - auto expression = std::make_shared(prewhere_info->prewhere_actions.clone()); - expression->describeActions(format_settings.out, prefix); - } + auto expression = std::make_shared(query_info.prewhere_info->prewhere_actions.clone()); + expression->describeActions(format_settings.out, prefix); + } - if (prewhere_info->row_level_filter) - { - format_settings.out << prefix << "Row level filter" << '\n'; - format_settings.out << prefix << "Row level filter column: " << prewhere_info->row_level_column_name << '\n'; + if (query_info.row_level_filter) + { + format_settings.out << prefix << "Row level filter" << '\n'; + format_settings.out << prefix << "Row level filter column: " << query_info.row_level_filter->column_name; + if (query_info.row_level_filter->do_remove_column) + format_settings.out << " (removed)"; + format_settings.out << '\n'; - auto expression = std::make_shared(prewhere_info->row_level_filter->clone()); - expression->describeActions(format_settings.out, prefix); - } + auto expression = std::make_shared(query_info.row_level_filter->actions.clone()); + expression->describeActions(format_settings.out, prefix); } if (virtual_row_conversion) @@ -2639,34 +2657,37 @@ void ReadFromMergeTree::describeActions(JSONBuilder::JSONMap & map) const map.add("Parts", result.index_stats.back().num_parts_after); map.add("Granules", result.index_stats.back().num_granules_after); } - - if (prewhere_info) + std::unique_ptr prewhere_info_map; + if (query_info.prewhere_info || query_info.row_level_filter) { - std::unique_ptr prewhere_info_map = std::make_unique(); - prewhere_info_map->add("Need filter", prewhere_info->need_filter); + prewhere_info_map = std::make_unique(); + if (query_info.prewhere_info) + prewhere_info_map->add("Need filter", query_info.prewhere_info->need_filter); + } - { - std::unique_ptr prewhere_filter_map = std::make_unique(); - prewhere_filter_map->add("Prewhere filter column", prewhere_info->prewhere_column_name); - prewhere_filter_map->add("Prewhere filter remove filter column", prewhere_info->remove_prewhere_column); - auto expression = std::make_shared(prewhere_info->prewhere_actions.clone()); - prewhere_filter_map->add("Prewhere filter expression", expression->toTree()); + if (query_info.prewhere_info) + { + std::unique_ptr prewhere_filter_map = std::make_unique(); + prewhere_filter_map->add("Prewhere filter column", query_info.prewhere_info->prewhere_column_name); + prewhere_filter_map->add("Prewhere filter remove filter column", query_info.prewhere_info->remove_prewhere_column); + auto expression = std::make_shared(query_info.prewhere_info->prewhere_actions.clone()); + prewhere_filter_map->add("Prewhere filter expression", expression->toTree()); - prewhere_info_map->add("Prewhere filter", std::move(prewhere_filter_map)); - } + prewhere_info_map->add("Prewhere filter", std::move(prewhere_filter_map)); + } - if (prewhere_info->row_level_filter) - { - std::unique_ptr row_level_filter_map = std::make_unique(); - row_level_filter_map->add("Row level filter column", prewhere_info->row_level_column_name); - auto expression = std::make_shared(prewhere_info->row_level_filter->clone()); - row_level_filter_map->add("Row level filter expression", expression->toTree()); + if (query_info.row_level_filter) + { + std::unique_ptr row_level_filter_map = std::make_unique(); + row_level_filter_map->add("Row level filter column", query_info.row_level_filter->column_name); + auto expression = std::make_shared(query_info.row_level_filter->actions.clone()); + row_level_filter_map->add("Row level filter expression", expression->toTree()); - prewhere_info_map->add("Row level filter", std::move(row_level_filter_map)); - } + prewhere_info_map->add("Row level filter", std::move(row_level_filter_map)); + } + if (prewhere_info_map) map.add("Prewhere info", std::move(prewhere_info_map)); - } if (virtual_row_conversion) map.add("Virtual row conversions", virtual_row_conversion->toTree()); diff --git a/src/Processors/QueryPlan/ReadFromObjectStorageStep.cpp b/src/Processors/QueryPlan/ReadFromObjectStorageStep.cpp index 34914c39fc2b..b5deeb839f76 100644 --- a/src/Processors/QueryPlan/ReadFromObjectStorageStep.cpp +++ b/src/Processors/QueryPlan/ReadFromObjectStorageStep.cpp @@ -74,9 +74,8 @@ void ReadFromObjectStorageStep::applyFilters(ActionDAGNodes added_filter_nodes) void ReadFromObjectStorageStep::updatePrewhereInfo(const PrewhereInfoPtr & prewhere_info_value) { - info = updateFormatPrewhereInfo(info, prewhere_info_value); + info = updateFormatPrewhereInfo(info, query_info.row_level_filter, prewhere_info_value); query_info.prewhere_info = prewhere_info_value; - prewhere_info = prewhere_info_value; output_header = std::make_shared(info.source_header); } @@ -100,9 +99,12 @@ void ReadFromObjectStorageStep::initializePipeline(QueryPipelineBuilder & pipeli // const size_t max_parsing_threads = (distributed_processing || num_streams >= max_threads) ? 1 : (max_threads / std::max(num_streams, 1ul)); auto parser_shared_resources = std::make_shared(context->getSettingsRef(), num_streams); - auto format_filter_info - = std::make_shared(filter_actions_dag, context, configuration->getColumnMapperForCurrentSchema()); - format_filter_info->prewhere_info = prewhere_info; + auto format_filter_info = std::make_shared( + filter_actions_dag, + context, + configuration->getColumnMapperForCurrentSchema(), + query_info.row_level_filter, + query_info.prewhere_info); for (size_t i = 0; i < num_streams; ++i) { diff --git a/src/Processors/QueryPlan/SourceStepWithFilter.cpp b/src/Processors/QueryPlan/SourceStepWithFilter.cpp index 702ad28e71c2..00c67a0f4390 100644 --- a/src/Processors/QueryPlan/SourceStepWithFilter.cpp +++ b/src/Processors/QueryPlan/SourceStepWithFilter.cpp @@ -16,25 +16,26 @@ namespace ErrorCodes extern const int ILLEGAL_TYPE_OF_COLUMN_FOR_FILTER; } -Block SourceStepWithFilter::applyPrewhereActions(Block block, const PrewhereInfoPtr & prewhere_info) +Block SourceStepWithFilter::applyPrewhereActions(Block block, const FilterDAGInfoPtr & row_level_filter, const PrewhereInfoPtr & prewhere_info) { - if (prewhere_info) + if (row_level_filter) { - if (prewhere_info->row_level_filter) + block = row_level_filter->actions.updateHeader(block); + auto & row_level_column = block.getByName(row_level_filter->column_name); + if (!row_level_column.type->canBeUsedInBooleanContext()) { - block = prewhere_info->row_level_filter->updateHeader(block); - auto & row_level_column = block.getByName(prewhere_info->row_level_column_name); - if (!row_level_column.type->canBeUsedInBooleanContext()) - { - throw Exception( - ErrorCodes::ILLEGAL_TYPE_OF_COLUMN_FOR_FILTER, - "Invalid type for filter in PREWHERE: {}", - row_level_column.type->getName()); - } - - block.erase(prewhere_info->row_level_column_name); + throw Exception( + ErrorCodes::ILLEGAL_TYPE_OF_COLUMN_FOR_FILTER, + "Invalid type for filter in PREWHERE: {}", + row_level_column.type->getName()); } + if (row_level_filter->do_remove_column) + block.erase(row_level_filter->column_name); + } + + if (prewhere_info) + { { block = prewhere_info->prewhere_actions.updateHeader(block); @@ -93,73 +94,81 @@ void SourceStepWithFilter::applyFilters(ActionDAGNodes added_filter_nodes) void SourceStepWithFilter::updatePrewhereInfo(const PrewhereInfoPtr & prewhere_info_value) { query_info.prewhere_info = prewhere_info_value; - prewhere_info = prewhere_info_value; - output_header = std::make_shared(applyPrewhereActions(*output_header, prewhere_info)); + output_header = std::make_shared(applyPrewhereActions(*output_header, query_info.row_level_filter, query_info.prewhere_info)); } void SourceStepWithFilter::describeActions(FormatSettings & format_settings) const { std::string prefix(format_settings.offset, format_settings.indent_char); - if (prewhere_info) + if (query_info.prewhere_info || query_info.row_level_filter) { format_settings.out << prefix << "Prewhere info" << '\n'; - format_settings.out << prefix << "Need filter: " << prewhere_info->need_filter << '\n'; + if (query_info.prewhere_info) + format_settings.out << prefix << "Need filter: " << query_info.prewhere_info->need_filter << '\n'; prefix.push_back(format_settings.indent_char); prefix.push_back(format_settings.indent_char); + } - { - format_settings.out << prefix << "Prewhere filter" << '\n'; - format_settings.out << prefix << "Prewhere filter column: " << prewhere_info->prewhere_column_name; - if (prewhere_info->remove_prewhere_column) - format_settings.out << " (removed)"; - format_settings.out << '\n'; - - auto expression = std::make_shared(prewhere_info->prewhere_actions.clone()); - expression->describeActions(format_settings.out, prefix); - } - - if (prewhere_info->row_level_filter) - { - format_settings.out << prefix << "Row level filter" << '\n'; - format_settings.out << prefix << "Row level filter column: " << prewhere_info->row_level_column_name << '\n'; + if (query_info.prewhere_info) + { + format_settings.out << prefix << "Prewhere filter" << '\n'; + format_settings.out << prefix << "Prewhere filter column: " << query_info.prewhere_info->prewhere_column_name; + if (query_info.prewhere_info->remove_prewhere_column) + format_settings.out << " (removed)"; + format_settings.out << '\n'; + + auto expression = std::make_shared(query_info.prewhere_info->prewhere_actions.clone()); + expression->describeActions(format_settings.out, prefix); + } - auto expression = std::make_shared(prewhere_info->row_level_filter->clone()); - expression->describeActions(format_settings.out, prefix); - } + if (query_info.row_level_filter) + { + format_settings.out << prefix << "Row level filter" << '\n'; + format_settings.out << prefix << "Row level filter column: " << query_info.row_level_filter->column_name; + if (query_info.row_level_filter->do_remove_column) + format_settings.out << " (removed)"; + format_settings.out << '\n'; + + auto expression = std::make_shared(query_info.row_level_filter->actions.clone()); + expression->describeActions(format_settings.out, prefix); } } void SourceStepWithFilter::describeActions(JSONBuilder::JSONMap & map) const { - if (prewhere_info) + std::unique_ptr prewhere_info_map; + if (query_info.prewhere_info || query_info.row_level_filter) { - std::unique_ptr prewhere_info_map = std::make_unique(); - prewhere_info_map->add("Need filter", prewhere_info->need_filter); + prewhere_info_map = std::make_unique(); + if (query_info.prewhere_info) + prewhere_info_map->add("Need filter", query_info.prewhere_info->need_filter); + } - { - std::unique_ptr prewhere_filter_map = std::make_unique(); - prewhere_filter_map->add("Prewhere filter column", prewhere_info->prewhere_column_name); - prewhere_filter_map->add("Prewhere filter remove filter column", prewhere_info->remove_prewhere_column); - auto expression = std::make_shared(prewhere_info->prewhere_actions.clone()); - prewhere_filter_map->add("Prewhere filter expression", expression->toTree()); + if (query_info.prewhere_info) + { + std::unique_ptr prewhere_filter_map = std::make_unique(); + prewhere_filter_map->add("Prewhere filter column", query_info.prewhere_info->prewhere_column_name); + prewhere_filter_map->add("Prewhere filter remove filter column", query_info.prewhere_info->remove_prewhere_column); + auto expression = std::make_shared(query_info.prewhere_info->prewhere_actions.clone()); + prewhere_filter_map->add("Prewhere filter expression", expression->toTree()); - prewhere_info_map->add("Prewhere filter", std::move(prewhere_filter_map)); - } + prewhere_info_map->add("Prewhere filter", std::move(prewhere_filter_map)); + } - if (prewhere_info->row_level_filter) - { - std::unique_ptr row_level_filter_map = std::make_unique(); - row_level_filter_map->add("Row level filter column", prewhere_info->row_level_column_name); - auto expression = std::make_shared(prewhere_info->row_level_filter->clone()); - row_level_filter_map->add("Row level filter expression", expression->toTree()); + if (query_info.row_level_filter) + { + std::unique_ptr row_level_filter_map = std::make_unique(); + row_level_filter_map->add("Row level filter column", query_info.row_level_filter->column_name); + auto expression = std::make_shared(query_info.row_level_filter->actions.clone()); + row_level_filter_map->add("Row level filter expression", expression->toTree()); - prewhere_info_map->add("Row level filter", std::move(row_level_filter_map)); - } + prewhere_info_map->add("Row level filter", std::move(row_level_filter_map)); + } + if (prewhere_info_map) map.add("Prewhere info", std::move(prewhere_info_map)); - } } } diff --git a/src/Processors/QueryPlan/SourceStepWithFilter.h b/src/Processors/QueryPlan/SourceStepWithFilter.h index f92af5245494..c5bd3dab9077 100644 --- a/src/Processors/QueryPlan/SourceStepWithFilter.h +++ b/src/Processors/QueryPlan/SourceStepWithFilter.h @@ -62,6 +62,7 @@ class SourceStepWithFilterBase : public ISourceStep } virtual void applyFilters(ActionDAGNodes added_filter_nodes); + virtual FilterDAGInfoPtr getRowLevelFilter() const { return nullptr; } virtual PrewhereInfoPtr getPrewhereInfo() const { return nullptr; } const std::shared_ptr & getFilterActionsDAG() const { return filter_actions_dag; } @@ -103,7 +104,6 @@ class SourceStepWithFilter : public SourceStepWithFilterBase : SourceStepWithFilterBase(std::move(output_header_)) , required_source_columns(column_names_) , query_info(query_info_) - , prewhere_info(query_info.prewhere_info) , storage_snapshot(storage_snapshot_) , context(context_) { @@ -112,7 +112,8 @@ class SourceStepWithFilter : public SourceStepWithFilterBase SourceStepWithFilter(const SourceStepWithFilter &) = default; const SelectQueryInfo & getQueryInfo() const { return query_info; } - PrewhereInfoPtr getPrewhereInfo() const override { return prewhere_info; } + FilterDAGInfoPtr getRowLevelFilter() const override { return query_info.row_level_filter; } + PrewhereInfoPtr getPrewhereInfo() const override { return query_info.prewhere_info; } ContextPtr getContext() const { return context; } const StorageSnapshotPtr & getStorageSnapshot() const { return storage_snapshot; } @@ -128,12 +129,11 @@ class SourceStepWithFilter : public SourceStepWithFilterBase void describeActions(JSONBuilder::JSONMap & map) const override; - static Block applyPrewhereActions(Block block, const PrewhereInfoPtr & prewhere_info); + static Block applyPrewhereActions(Block block, const FilterDAGInfoPtr & row_level_filter, const PrewhereInfoPtr & prewhere_info); protected: Names required_source_columns; SelectQueryInfo query_info; - PrewhereInfoPtr prewhere_info; StorageSnapshotPtr storage_snapshot; ContextPtr context; }; diff --git a/src/Storages/IStorage.cpp b/src/Storages/IStorage.cpp index 3766afafe6b9..88768173c162 100644 --- a/src/Storages/IStorage.cpp +++ b/src/Storages/IStorage.cpp @@ -415,11 +415,6 @@ std::string PrewhereInfo::dump() const WriteBufferFromOwnString ss; ss << "PrewhereDagInfo\n"; - if (row_level_filter) - { - ss << "row_level_filter " << row_level_filter->dumpDAG() << "\n"; - } - { ss << "prewhere_actions " << prewhere_actions.dumpDAG() << "\n"; } diff --git a/src/Storages/MergeTree/MergeTreeBlockReadUtils.cpp b/src/Storages/MergeTree/MergeTreeBlockReadUtils.cpp index 5287e609cd91..c021809f7177 100644 --- a/src/Storages/MergeTree/MergeTreeBlockReadUtils.cpp +++ b/src/Storages/MergeTree/MergeTreeBlockReadUtils.cpp @@ -365,6 +365,7 @@ MergeTreeReadTaskColumns getReadTaskColumns( const IMergeTreeDataPartInfoForReader & data_part_info_for_reader, const StorageSnapshotPtr & storage_snapshot, const Names & required_columns, + const FilterDAGInfoPtr & row_level_filter, const PrewhereInfoPtr & prewhere_info, const PrewhereExprSteps & mutation_steps, const ExpressionActionsSettings & actions_settings, @@ -438,9 +439,10 @@ MergeTreeReadTaskColumns getReadTaskColumns( for (const auto & step : mutation_steps) add_step(*step); - if (prewhere_info) + if (prewhere_info || row_level_filter) { auto prewhere_actions = MergeTreeSelectProcessor::getPrewhereActions( + row_level_filter, prewhere_info, actions_settings, reader_settings.enable_multiple_prewhere_read_steps, reader_settings.force_short_circuit_execution); @@ -471,6 +473,7 @@ MergeTreeReadTaskColumns getReadTaskColumnsForMerge( data_part_info_for_reader, storage_snapshot, required_columns, + /*row_level_filter=*/ nullptr, /*prewhere_info=*/ nullptr, mutation_steps, /*actions_settings=*/ {}, diff --git a/src/Storages/MergeTree/MergeTreeBlockReadUtils.h b/src/Storages/MergeTree/MergeTreeBlockReadUtils.h index 6bf95b605a5a..9745069d454d 100644 --- a/src/Storages/MergeTree/MergeTreeBlockReadUtils.h +++ b/src/Storages/MergeTree/MergeTreeBlockReadUtils.h @@ -33,6 +33,7 @@ MergeTreeReadTaskColumns getReadTaskColumns( const IMergeTreeDataPartInfoForReader & data_part_info_for_reader, const StorageSnapshotPtr & storage_snapshot, const Names & required_columns, + const FilterDAGInfoPtr & row_level_filter, const PrewhereInfoPtr & prewhere_info, const PrewhereExprSteps & mutation_steps, const ExpressionActionsSettings & actions_settings, diff --git a/src/Storages/MergeTree/MergeTreePrefetchedReadPool.cpp b/src/Storages/MergeTree/MergeTreePrefetchedReadPool.cpp index af3e1aabaa36..2dd137676c7a 100644 --- a/src/Storages/MergeTree/MergeTreePrefetchedReadPool.cpp +++ b/src/Storages/MergeTree/MergeTreePrefetchedReadPool.cpp @@ -106,6 +106,7 @@ MergeTreePrefetchedReadPool::MergeTreePrefetchedReadPool( MutationsSnapshotPtr mutations_snapshot_, VirtualFields shared_virtual_fields_, const StorageSnapshotPtr & storage_snapshot_, + const FilterDAGInfoPtr & row_level_filter_, const PrewhereInfoPtr & prewhere_info_, const ExpressionActionsSettings & actions_settings_, const MergeTreeReaderSettings & reader_settings_, @@ -118,6 +119,7 @@ MergeTreePrefetchedReadPool::MergeTreePrefetchedReadPool( std::move(mutations_snapshot_), std::move(shared_virtual_fields_), storage_snapshot_, + row_level_filter_, prewhere_info_, actions_settings_, reader_settings_, diff --git a/src/Storages/MergeTree/MergeTreePrefetchedReadPool.h b/src/Storages/MergeTree/MergeTreePrefetchedReadPool.h index 2c6687002dc0..d07681a00f00 100644 --- a/src/Storages/MergeTree/MergeTreePrefetchedReadPool.h +++ b/src/Storages/MergeTree/MergeTreePrefetchedReadPool.h @@ -22,6 +22,7 @@ class MergeTreePrefetchedReadPool : public MergeTreeReadPoolBase MutationsSnapshotPtr mutations_snapshot_, VirtualFields shared_virtual_fields_, const StorageSnapshotPtr & storage_snapshot_, + const FilterDAGInfoPtr & row_level_filter_, const PrewhereInfoPtr & prewhere_info_, const ExpressionActionsSettings & actions_settings_, const MergeTreeReaderSettings & reader_settings_, diff --git a/src/Storages/MergeTree/MergeTreeReadPool.cpp b/src/Storages/MergeTree/MergeTreeReadPool.cpp index 9fb56a1e9381..db173e064f97 100644 --- a/src/Storages/MergeTree/MergeTreeReadPool.cpp +++ b/src/Storages/MergeTree/MergeTreeReadPool.cpp @@ -40,6 +40,7 @@ MergeTreeReadPool::MergeTreeReadPool( MutationsSnapshotPtr mutations_snapshot_, VirtualFields shared_virtual_fields_, const StorageSnapshotPtr & storage_snapshot_, + const FilterDAGInfoPtr & row_level_filter_, const PrewhereInfoPtr & prewhere_info_, const ExpressionActionsSettings & actions_settings_, const MergeTreeReaderSettings & reader_settings_, @@ -52,6 +53,7 @@ MergeTreeReadPool::MergeTreeReadPool( std::move(mutations_snapshot_), std::move(shared_virtual_fields_), storage_snapshot_, + row_level_filter_, prewhere_info_, actions_settings_, reader_settings_, diff --git a/src/Storages/MergeTree/MergeTreeReadPool.h b/src/Storages/MergeTree/MergeTreeReadPool.h index 1b55284c592b..8891158d2c14 100644 --- a/src/Storages/MergeTree/MergeTreeReadPool.h +++ b/src/Storages/MergeTree/MergeTreeReadPool.h @@ -29,6 +29,7 @@ class MergeTreeReadPool : public MergeTreeReadPoolBase MutationsSnapshotPtr mutations_snapshot_, VirtualFields shared_virtual_fields_, const StorageSnapshotPtr & storage_snapshot_, + const FilterDAGInfoPtr & row_level_filter_, const PrewhereInfoPtr & prewhere_info_, const ExpressionActionsSettings & actions_settings_, const MergeTreeReaderSettings & reader_settings_, diff --git a/src/Storages/MergeTree/MergeTreeReadPoolBase.cpp b/src/Storages/MergeTree/MergeTreeReadPoolBase.cpp index e2f615574ff6..3bbbe618e9dc 100644 --- a/src/Storages/MergeTree/MergeTreeReadPoolBase.cpp +++ b/src/Storages/MergeTree/MergeTreeReadPoolBase.cpp @@ -29,6 +29,7 @@ MergeTreeReadPoolBase::MergeTreeReadPoolBase( MutationsSnapshotPtr mutations_snapshot_, VirtualFields shared_virtual_fields_, const StorageSnapshotPtr & storage_snapshot_, + const FilterDAGInfoPtr & row_level_filter_, const PrewhereInfoPtr & prewhere_info_, const ExpressionActionsSettings & actions_settings_, const MergeTreeReaderSettings & reader_settings_, @@ -41,6 +42,7 @@ MergeTreeReadPoolBase::MergeTreeReadPoolBase( , mutations_snapshot(std::move(mutations_snapshot_)) , shared_virtual_fields(std::move(shared_virtual_fields_)) , storage_snapshot(storage_snapshot_) + , row_level_filter(row_level_filter_) , prewhere_info(prewhere_info_) , actions_settings(actions_settings_) , reader_settings(reader_settings_) @@ -188,6 +190,7 @@ void MergeTreeReadPoolBase::fillPerPartInfos(const Settings & settings) part_info, storage_snapshot, column_names, + row_level_filter, prewhere_info, read_task_info.mutation_steps, actions_settings, diff --git a/src/Storages/MergeTree/MergeTreeReadPoolBase.h b/src/Storages/MergeTree/MergeTreeReadPoolBase.h index eb4314642619..eae6df8b5e78 100644 --- a/src/Storages/MergeTree/MergeTreeReadPoolBase.h +++ b/src/Storages/MergeTree/MergeTreeReadPoolBase.h @@ -32,6 +32,7 @@ class MergeTreeReadPoolBase : public IMergeTreeReadPool, protected WithContext MutationsSnapshotPtr mutations_snapshot_, VirtualFields shared_virtual_fields_, const StorageSnapshotPtr & storage_snapshot_, + const FilterDAGInfoPtr & row_level_filter_, const PrewhereInfoPtr & prewhere_info_, const ExpressionActionsSettings & actions_settings_, const MergeTreeReaderSettings & reader_settings_, @@ -48,6 +49,7 @@ class MergeTreeReadPoolBase : public IMergeTreeReadPool, protected WithContext const MutationsSnapshotPtr mutations_snapshot; const VirtualFields shared_virtual_fields; const StorageSnapshotPtr storage_snapshot; + const FilterDAGInfoPtr row_level_filter; const PrewhereInfoPtr prewhere_info; const ExpressionActionsSettings actions_settings; const MergeTreeReaderSettings reader_settings; diff --git a/src/Storages/MergeTree/MergeTreeReadPoolInOrder.cpp b/src/Storages/MergeTree/MergeTreeReadPoolInOrder.cpp index c4244ecd9820..1aa6b5eedc19 100644 --- a/src/Storages/MergeTree/MergeTreeReadPoolInOrder.cpp +++ b/src/Storages/MergeTree/MergeTreeReadPoolInOrder.cpp @@ -15,6 +15,7 @@ MergeTreeReadPoolInOrder::MergeTreeReadPoolInOrder( MutationsSnapshotPtr mutations_snapshot_, VirtualFields shared_virtual_fields_, const StorageSnapshotPtr & storage_snapshot_, + const FilterDAGInfoPtr & row_level_filter_, const PrewhereInfoPtr & prewhere_info_, const ExpressionActionsSettings & actions_settings_, const MergeTreeReaderSettings & reader_settings_, @@ -27,6 +28,7 @@ MergeTreeReadPoolInOrder::MergeTreeReadPoolInOrder( std::move(mutations_snapshot_), std::move(shared_virtual_fields_), storage_snapshot_, + row_level_filter_, prewhere_info_, actions_settings_, reader_settings_, diff --git a/src/Storages/MergeTree/MergeTreeReadPoolInOrder.h b/src/Storages/MergeTree/MergeTreeReadPoolInOrder.h index 41f3ab1061c1..35c968acb28c 100644 --- a/src/Storages/MergeTree/MergeTreeReadPoolInOrder.h +++ b/src/Storages/MergeTree/MergeTreeReadPoolInOrder.h @@ -14,6 +14,7 @@ class MergeTreeReadPoolInOrder : public MergeTreeReadPoolBase MutationsSnapshotPtr mutations_snapshot_, VirtualFields shared_virtual_fields_, const StorageSnapshotPtr & storage_snapshot_, + const FilterDAGInfoPtr & row_level_filter_, const PrewhereInfoPtr & prewhere_info_, const ExpressionActionsSettings & actions_settings_, const MergeTreeReaderSettings & reader_settings_, diff --git a/src/Storages/MergeTree/MergeTreeReadPoolParallelReplicas.cpp b/src/Storages/MergeTree/MergeTreeReadPoolParallelReplicas.cpp index c3e5d609b8ff..20af6def427e 100644 --- a/src/Storages/MergeTree/MergeTreeReadPoolParallelReplicas.cpp +++ b/src/Storages/MergeTree/MergeTreeReadPoolParallelReplicas.cpp @@ -112,6 +112,7 @@ MergeTreeReadPoolParallelReplicas::MergeTreeReadPoolParallelReplicas( MutationsSnapshotPtr mutations_snapshot_, VirtualFields shared_virtual_fields_, const StorageSnapshotPtr & storage_snapshot_, + const FilterDAGInfoPtr & row_level_filter_, const PrewhereInfoPtr & prewhere_info_, const ExpressionActionsSettings & actions_settings_, const MergeTreeReaderSettings & reader_settings_, @@ -124,6 +125,7 @@ MergeTreeReadPoolParallelReplicas::MergeTreeReadPoolParallelReplicas( std::move(mutations_snapshot_), std::move(shared_virtual_fields_), storage_snapshot_, + row_level_filter_, prewhere_info_, actions_settings_, reader_settings_, diff --git a/src/Storages/MergeTree/MergeTreeReadPoolParallelReplicas.h b/src/Storages/MergeTree/MergeTreeReadPoolParallelReplicas.h index 63816340eb1d..b5fbd4efc1ba 100644 --- a/src/Storages/MergeTree/MergeTreeReadPoolParallelReplicas.h +++ b/src/Storages/MergeTree/MergeTreeReadPoolParallelReplicas.h @@ -14,6 +14,7 @@ class MergeTreeReadPoolParallelReplicas : public MergeTreeReadPoolBase MutationsSnapshotPtr mutations_snapshot_, VirtualFields shared_virtual_fields_, const StorageSnapshotPtr & storage_snapshot_, + const FilterDAGInfoPtr & row_level_filter_, const PrewhereInfoPtr & prewhere_info_, const ExpressionActionsSettings & actions_settings_, const MergeTreeReaderSettings & reader_settings_, diff --git a/src/Storages/MergeTree/MergeTreeReadPoolParallelReplicasInOrder.cpp b/src/Storages/MergeTree/MergeTreeReadPoolParallelReplicasInOrder.cpp index cdc84ff6c0eb..6ff17f078e29 100644 --- a/src/Storages/MergeTree/MergeTreeReadPoolParallelReplicasInOrder.cpp +++ b/src/Storages/MergeTree/MergeTreeReadPoolParallelReplicasInOrder.cpp @@ -28,6 +28,7 @@ MergeTreeReadPoolParallelReplicasInOrder::MergeTreeReadPoolParallelReplicasInOrd VirtualFields shared_virtual_fields_, bool has_limit_below_one_block_, const StorageSnapshotPtr & storage_snapshot_, + const FilterDAGInfoPtr & row_level_filter_, const PrewhereInfoPtr & prewhere_info_, const ExpressionActionsSettings & actions_settings_, const MergeTreeReaderSettings & reader_settings_, @@ -40,6 +41,7 @@ MergeTreeReadPoolParallelReplicasInOrder::MergeTreeReadPoolParallelReplicasInOrd std::move(mutations_snapshot_), std::move(shared_virtual_fields_), storage_snapshot_, + row_level_filter_, prewhere_info_, actions_settings_, reader_settings_, diff --git a/src/Storages/MergeTree/MergeTreeReadPoolParallelReplicasInOrder.h b/src/Storages/MergeTree/MergeTreeReadPoolParallelReplicasInOrder.h index e7fad180657b..c8fb1f2c89ee 100644 --- a/src/Storages/MergeTree/MergeTreeReadPoolParallelReplicasInOrder.h +++ b/src/Storages/MergeTree/MergeTreeReadPoolParallelReplicasInOrder.h @@ -16,6 +16,7 @@ class MergeTreeReadPoolParallelReplicasInOrder : public MergeTreeReadPoolBase VirtualFields shared_virtual_fields_, bool has_limit_below_one_block_, const StorageSnapshotPtr & storage_snapshot_, + const FilterDAGInfoPtr & row_level_filter_, const PrewhereInfoPtr & prewhere_info_, const ExpressionActionsSettings & actions_settings_, const MergeTreeReaderSettings & reader_settings_, diff --git a/src/Storages/MergeTree/MergeTreeSelectProcessor.cpp b/src/Storages/MergeTree/MergeTreeSelectProcessor.cpp index 642a0f4ad4e2..75529a5fc8d8 100644 --- a/src/Storages/MergeTree/MergeTreeSelectProcessor.cpp +++ b/src/Storages/MergeTree/MergeTreeSelectProcessor.cpp @@ -92,22 +92,25 @@ std::optional ParallelReadingExtension::sendReadRequest( MergeTreeSelectProcessor::MergeTreeSelectProcessor( MergeTreeReadPoolPtr pool_, MergeTreeSelectAlgorithmPtr algorithm_, + const FilterDAGInfoPtr & row_level_filter_, const PrewhereInfoPtr & prewhere_info_, const LazilyReadInfoPtr & lazily_read_info_, const ExpressionActionsSettings & actions_settings_, const MergeTreeReaderSettings & reader_settings_) : pool(std::move(pool_)) , algorithm(std::move(algorithm_)) + , row_level_filter(row_level_filter_) , prewhere_info(prewhere_info_) , actions_settings(actions_settings_) , prewhere_actions(getPrewhereActions( + row_level_filter, prewhere_info, actions_settings, reader_settings_.enable_multiple_prewhere_read_steps, reader_settings_.force_short_circuit_execution)) , lazily_read_info(lazily_read_info_) , reader_settings(reader_settings_) - , result_header(transformHeader(pool->getHeader(), lazily_read_info, prewhere_info)) + , result_header(transformHeader(pool->getHeader(), lazily_read_info, row_level_filter, prewhere_info)) { bool has_prewhere_actions_steps = !prewhere_actions.steps.empty(); if (has_prewhere_actions_steps) @@ -127,43 +130,46 @@ String MergeTreeSelectProcessor::getName() const bool tryBuildPrewhereSteps(PrewhereInfoPtr prewhere_info, const ExpressionActionsSettings & actions_settings, PrewhereExprInfo & prewhere, bool force_short_circuit_execution); -PrewhereExprInfo MergeTreeSelectProcessor::getPrewhereActions(PrewhereInfoPtr prewhere_info, const ExpressionActionsSettings & actions_settings, bool enable_multiple_prewhere_read_steps, bool force_short_circuit_execution) +PrewhereExprInfo MergeTreeSelectProcessor::getPrewhereActions( + const FilterDAGInfoPtr & row_level_filter, + const PrewhereInfoPtr & prewhere_info, + const ExpressionActionsSettings & actions_settings, + bool enable_multiple_prewhere_read_steps, + bool force_short_circuit_execution) { PrewhereExprInfo prewhere_actions; - if (prewhere_info) + + if (row_level_filter) { - if (prewhere_info->row_level_filter) + PrewhereExprStep row_level_filter_step { - PrewhereExprStep row_level_filter_step - { - .type = PrewhereExprStep::Filter, - .actions = std::make_shared(prewhere_info->row_level_filter->clone(), actions_settings), - .filter_column_name = prewhere_info->row_level_column_name, - .remove_filter_column = true, - .need_filter = true, - .perform_alter_conversions = true, - .mutation_version = std::nullopt, - }; - - prewhere_actions.steps.emplace_back(std::make_shared(std::move(row_level_filter_step))); - } + .type = PrewhereExprStep::Filter, + .actions = std::make_shared(row_level_filter->actions.clone(), actions_settings), + .filter_column_name = row_level_filter->column_name, + .remove_filter_column = row_level_filter->do_remove_column, + .need_filter = true, + .perform_alter_conversions = true, + .mutation_version = std::nullopt, + }; + + prewhere_actions.steps.emplace_back(std::make_shared(std::move(row_level_filter_step))); + } - if (!enable_multiple_prewhere_read_steps || - !tryBuildPrewhereSteps(prewhere_info, actions_settings, prewhere_actions, force_short_circuit_execution)) + if (prewhere_info && + (!enable_multiple_prewhere_read_steps || !tryBuildPrewhereSteps(prewhere_info, actions_settings, prewhere_actions, force_short_circuit_execution))) + { + PrewhereExprStep prewhere_step { - PrewhereExprStep prewhere_step - { - .type = PrewhereExprStep::Filter, - .actions = std::make_shared(prewhere_info->prewhere_actions.clone(), actions_settings), - .filter_column_name = prewhere_info->prewhere_column_name, - .remove_filter_column = prewhere_info->remove_prewhere_column, - .need_filter = prewhere_info->need_filter, - .perform_alter_conversions = true, - .mutation_version = std::nullopt, - }; - - prewhere_actions.steps.emplace_back(std::make_shared(std::move(prewhere_step))); - } + .type = PrewhereExprStep::Filter, + .actions = std::make_shared(prewhere_info->prewhere_actions.clone(), actions_settings), + .filter_column_name = prewhere_info->prewhere_column_name, + .remove_filter_column = prewhere_info->remove_prewhere_column, + .need_filter = prewhere_info->need_filter, + .perform_alter_conversions = true, + .mutation_version = std::nullopt, + }; + + prewhere_actions.steps.emplace_back(std::make_shared(std::move(prewhere_step))); } return prewhere_actions; @@ -329,9 +335,10 @@ void MergeTreeSelectProcessor::injectLazilyReadColumns( Block MergeTreeSelectProcessor::transformHeader( Block block, const LazilyReadInfoPtr & lazily_read_info, + const FilterDAGInfoPtr & row_level_filter, const PrewhereInfoPtr & prewhere_info) { - auto transformed = SourceStepWithFilter::applyPrewhereActions(std::move(block), prewhere_info); + auto transformed = SourceStepWithFilter::applyPrewhereActions(std::move(block), row_level_filter, prewhere_info); injectLazilyReadColumns(0, transformed, -1, lazily_read_info); return transformed; } diff --git a/src/Storages/MergeTree/MergeTreeSelectProcessor.h b/src/Storages/MergeTree/MergeTreeSelectProcessor.h index 7c898c91e36e..de04d99f9c74 100644 --- a/src/Storages/MergeTree/MergeTreeSelectProcessor.h +++ b/src/Storages/MergeTree/MergeTreeSelectProcessor.h @@ -63,6 +63,7 @@ class MergeTreeSelectProcessor : private boost::noncopyable MergeTreeSelectProcessor( MergeTreeReadPoolPtr pool_, MergeTreeSelectAlgorithmPtr algorithm_, + const FilterDAGInfoPtr & row_level_filter_, const PrewhereInfoPtr & prewhere_info_, const LazilyReadInfoPtr & lazily_read_info_, const ExpressionActionsSettings & actions_settings_, @@ -73,6 +74,7 @@ class MergeTreeSelectProcessor : private boost::noncopyable static Block transformHeader( Block block, const LazilyReadInfoPtr & lazily_read_info, + const FilterDAGInfoPtr & row_level_filter, const PrewhereInfoPtr & prewhere_info); Block getHeader() const { return result_header; } @@ -84,7 +86,8 @@ class MergeTreeSelectProcessor : private boost::noncopyable const MergeTreeReaderSettings & getSettings() const { return reader_settings; } static PrewhereExprInfo getPrewhereActions( - PrewhereInfoPtr prewhere_info, + const FilterDAGInfoPtr & row_level_filter, + const PrewhereInfoPtr & prewhere_info, const ExpressionActionsSettings & actions_settings, bool enable_multiple_prewhere_read_steps, bool force_short_circuit_execution); @@ -106,6 +109,7 @@ class MergeTreeSelectProcessor : private boost::noncopyable const MergeTreeReadPoolPtr pool; const MergeTreeSelectAlgorithmPtr algorithm; + const FilterDAGInfoPtr row_level_filter; const PrewhereInfoPtr prewhere_info; const ExpressionActionsSettings actions_settings; const PrewhereExprInfo prewhere_actions; diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/Compaction.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/Compaction.cpp index c2e783dc73e9..0fe29c16575e 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/Compaction.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/Compaction.cpp @@ -265,7 +265,7 @@ void writeDataFiles( 8192, format_settings, parser_shared_resources, - std::make_shared(nullptr, context, nullptr), + std::make_shared(nullptr, context, nullptr, nullptr, nullptr), true /* is_remote_fs */, chooseCompressionMethod(data_file->data_object_info->getPath(), configuration->getCompressionMethod()), false); diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/Mutations.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/Mutations.cpp index 5b5df91860ee..99fa70c81727 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/Mutations.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/Mutations.cpp @@ -143,7 +143,7 @@ std::optional writeDataFiles( field_ids[IcebergPositionDeleteTransform::positions_column_name] = IcebergPositionDeleteTransform::positions_column_field_id; field_ids[IcebergPositionDeleteTransform::data_file_path_column_name] = IcebergPositionDeleteTransform::data_file_path_column_field_id; column_mapper->setStorageColumnEncoding(std::move(field_ids)); - FormatFilterInfoPtr format_filter_info = std::make_shared(nullptr, context, column_mapper); + FormatFilterInfoPtr format_filter_info = std::make_shared(nullptr, context, column_mapper, nullptr, nullptr); auto output_format = FormatFactory::instance().getOutputFormat( configuration->getFormat(), *write_buffer, delete_file_sample_block, context, format_settings, format_filter_info); diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/PositionDeleteTransform.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/PositionDeleteTransform.cpp index d97b786d1ff0..adc3b4f73c52 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/PositionDeleteTransform.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/PositionDeleteTransform.cpp @@ -115,7 +115,7 @@ void IcebergPositionDeleteTransform::initializeDeleteSources() context->getSettingsRef()[DB::Setting::max_block_size], format_settings, std::make_shared(context->getSettingsRef(), 1), - std::make_shared(actions_dag_ptr, context, nullptr), + std::make_shared(actions_dag_ptr, context, nullptr, nullptr, nullptr), true /* is_remote_fs */, compression_method); diff --git a/src/Storages/ObjectStorage/StorageObjectStorage.cpp b/src/Storages/ObjectStorage/StorageObjectStorage.cpp index 889186adb022..f7861c83b8b5 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorage.cpp +++ b/src/Storages/ObjectStorage/StorageObjectStorage.cpp @@ -391,10 +391,10 @@ void StorageObjectStorage::read( supports_tuple_elements, local_context, PrepareReadingFromFormatHiveParams { file_columns, hive_partition_columns_to_read_from_file_path.getNameToTypeMap() }); - if (query_info.prewhere_info) - read_from_format_info = updateFormatPrewhereInfo(read_from_format_info, query_info.prewhere_info); + if (query_info.prewhere_info || query_info.row_level_filter) + read_from_format_info = updateFormatPrewhereInfo(read_from_format_info, query_info.row_level_filter, query_info.prewhere_info); - const bool need_only_count = (query_info.optimize_trivial_count || (read_from_format_info.requested_columns.empty() && !read_from_format_info.prewhere_info)) + const bool need_only_count = (query_info.optimize_trivial_count || (read_from_format_info.requested_columns.empty() && !read_from_format_info.prewhere_info && !read_from_format_info.row_level_filter)) && local_context->getSettingsRef()[Setting::optimize_count_from_files]; auto modified_format_settings{format_settings}; diff --git a/src/Storages/ObjectStorage/StorageObjectStorageSource.cpp b/src/Storages/ObjectStorage/StorageObjectStorageSource.cpp index ebab6fe533b4..77ed6608da93 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorageSource.cpp +++ b/src/Storages/ObjectStorage/StorageObjectStorageSource.cpp @@ -769,7 +769,7 @@ StorageObjectStorageSource::ReaderHolder StorageObjectStorageSource::createReade auto mapper = configuration->getColumnMapperForObject(object_info); if (!mapper) return format_filter_info; - return std::make_shared(format_filter_info->filter_actions_dag, format_filter_info->context.lock(), mapper); + return std::make_shared(format_filter_info->filter_actions_dag, format_filter_info->context.lock(), mapper, nullptr, nullptr); }(); auto input_format = FormatFactory::instance().getInput( diff --git a/src/Storages/SelectQueryInfo.cpp b/src/Storages/SelectQueryInfo.cpp index 6965822427d6..84254910ac62 100644 --- a/src/Storages/SelectQueryInfo.cpp +++ b/src/Storages/SelectQueryInfo.cpp @@ -38,51 +38,53 @@ std::unordered_map SelectQueryInfo::buildNod return node_name_to_input_node_column; } -PrewhereInfoPtr PrewhereInfo::clone() const +PrewhereInfo PrewhereInfo::clone() const { - PrewhereInfoPtr prewhere_info = std::make_shared(); + PrewhereInfo prewhere_info; - if (row_level_filter) - prewhere_info->row_level_filter = row_level_filter->clone(); - - prewhere_info->prewhere_actions = prewhere_actions.clone(); - - prewhere_info->row_level_column_name = row_level_column_name; - prewhere_info->prewhere_column_name = prewhere_column_name; - prewhere_info->remove_prewhere_column = remove_prewhere_column; - prewhere_info->need_filter = need_filter; - prewhere_info->generated_by_optimizer = generated_by_optimizer; + prewhere_info.prewhere_actions = prewhere_actions.clone(); + prewhere_info.prewhere_column_name = prewhere_column_name; + prewhere_info.remove_prewhere_column = remove_prewhere_column; + prewhere_info.need_filter = need_filter; return prewhere_info; } void PrewhereInfo::serialize(IQueryPlanStep::Serialization & ctx) const { - writeBinary(row_level_filter.has_value(), ctx.out); - if (row_level_filter.has_value()) - row_level_filter->serialize(ctx.out, ctx.registry); prewhere_actions.serialize(ctx.out, ctx.registry); - writeStringBinary(row_level_column_name, ctx.out); writeStringBinary(prewhere_column_name, ctx.out); writeBinary(remove_prewhere_column, ctx.out); - writeBinary(need_filter, ctx.out); - writeBinary(generated_by_optimizer, ctx.out); } -PrewhereInfoPtr PrewhereInfo::deserialize(IQueryPlanStep::Deserialization & ctx) +PrewhereInfo PrewhereInfo::deserialize(IQueryPlanStep::Deserialization & ctx) { - PrewhereInfoPtr result = std::make_shared(); - bool has_row_level_filter; - readBinary(has_row_level_filter, ctx.in); - if (has_row_level_filter) - result->row_level_filter = ActionsDAG::deserialize(ctx.in, ctx.registry, ctx.context); - result->prewhere_actions = ActionsDAG::deserialize(ctx.in, ctx.registry, ctx.context); - readStringBinary(result->row_level_column_name, ctx.in); - readStringBinary(result->prewhere_column_name, ctx.in); - readBinary(result->remove_prewhere_column, ctx.in); - readBinary(result->need_filter, ctx.in); - readBinary(result->generated_by_optimizer, ctx.in); - return result; + PrewhereInfo prewhere_info; + + prewhere_info.prewhere_actions = ActionsDAG::deserialize(ctx.in, ctx.registry, ctx.context); + readStringBinary(prewhere_info.prewhere_column_name, ctx.in); + readBinary(prewhere_info.remove_prewhere_column, ctx.in); + prewhere_info.need_filter = true; + + return prewhere_info; +} + +void FilterDAGInfo::serialize(IQueryPlanStep::Serialization & ctx) const +{ + actions.serialize(ctx.out, ctx.registry); + writeStringBinary(column_name, ctx.out); + writeBinary(do_remove_column, ctx.out); +} + +FilterDAGInfo FilterDAGInfo::deserialize(IQueryPlanStep::Deserialization & ctx) +{ + FilterDAGInfo filter_dag_info; + + filter_dag_info.actions = ActionsDAG::deserialize(ctx.in, ctx.registry, ctx.context); + readStringBinary(filter_dag_info.column_name, ctx.in); + readBinary(filter_dag_info.do_remove_column, ctx.in); + + return filter_dag_info; } } diff --git a/src/Storages/SelectQueryInfo.h b/src/Storages/SelectQueryInfo.h index 77448e4c1ee1..2252424e259e 100644 --- a/src/Storages/SelectQueryInfo.h +++ b/src/Storages/SelectQueryInfo.h @@ -42,16 +42,11 @@ using PreparedSetsPtr = std::shared_ptr; struct PrewhereInfo { - /// Actions for row level security filter. Applied separately before prewhere_actions. - /// This actions are separate because prewhere condition should not be executed over filtered rows. - std::optional row_level_filter; /// Actions which are executed on block in order to get filter column for prewhere step. ActionsDAG prewhere_actions; - String row_level_column_name; String prewhere_column_name; bool remove_prewhere_column = false; bool need_filter = false; - bool generated_by_optimizer = false; PrewhereInfo() = default; explicit PrewhereInfo(ActionsDAG prewhere_actions_, String prewhere_column_name_) @@ -59,10 +54,10 @@ struct PrewhereInfo std::string dump() const; - PrewhereInfoPtr clone() const; + PrewhereInfo clone() const; void serialize(IQueryPlanStep::Serialization & ctx) const; - static PrewhereInfoPtr deserialize(IQueryPlanStep::Deserialization & ctx); + static PrewhereInfo deserialize(IQueryPlanStep::Deserialization & ctx); }; /// Same as FilterInfo, but with ActionsDAG. @@ -73,6 +68,9 @@ struct FilterDAGInfo bool do_remove_column = false; std::string dump() const; + + void serialize(IQueryPlanStep::Serialization & ctx) const; + static FilterDAGInfo deserialize(IQueryPlanStep::Deserialization & ctx); }; struct InputOrderInfo @@ -190,6 +188,10 @@ struct SelectQueryInfo bool has_window = false; bool has_order_by = false; bool need_aggregate = false; + + /// Actions for row level security filter. Applied separately before prewhere. + /// This actions are separate because prewhere condition should not be executed over filtered rows. + FilterDAGInfoPtr row_level_filter; PrewhereInfoPtr prewhere_info; /// If query has aggregate functions diff --git a/src/Storages/StorageBuffer.cpp b/src/Storages/StorageBuffer.cpp index 7e7e43db21c7..9a7ed4ef046d 100644 --- a/src/Storages/StorageBuffer.cpp +++ b/src/Storages/StorageBuffer.cpp @@ -355,27 +355,35 @@ void StorageBuffer::read( else { auto src_table_query_info = query_info; - if (src_table_query_info.prewhere_info) + ActionsDAG converting_dag; + if (src_table_query_info.prewhere_info || src_table_query_info.row_level_filter) { - src_table_query_info.prewhere_info = src_table_query_info.prewhere_info->clone(); + converting_dag = ActionsDAG::makeConvertingActions( + header_after_adding_defaults.getColumnsWithTypeAndName(), + header.getColumnsWithTypeAndName(), + ActionsDAG::MatchColumnsMode::Name); + } - auto actions_dag = ActionsDAG::makeConvertingActions( - header_after_adding_defaults.getColumnsWithTypeAndName(), - header.getColumnsWithTypeAndName(), - ActionsDAG::MatchColumnsMode::Name); + if (src_table_query_info.row_level_filter) + { + auto row_level_filter = std::make_shared(); + row_level_filter->column_name = src_table_query_info.row_level_filter->column_name; + row_level_filter->do_remove_column = src_table_query_info.row_level_filter->do_remove_column; - if (src_table_query_info.prewhere_info->row_level_filter) - { - src_table_query_info.prewhere_info->row_level_filter = ActionsDAG::merge( - actions_dag.clone(), - std::move(*src_table_query_info.prewhere_info->row_level_filter)); + row_level_filter->actions = ActionsDAG::merge( + converting_dag.clone(), + src_table_query_info.row_level_filter->actions.clone()); - src_table_query_info.prewhere_info->row_level_filter->removeUnusedActions(); - } + row_level_filter->actions.removeUnusedActions(); + src_table_query_info.row_level_filter = std::move(row_level_filter); + } + if (src_table_query_info.prewhere_info) + { + src_table_query_info.prewhere_info = std::make_shared(src_table_query_info.prewhere_info->clone()); { src_table_query_info.prewhere_info->prewhere_actions = ActionsDAG::merge( - actions_dag.clone(), + converting_dag.clone(), std::move(src_table_query_info.prewhere_info->prewhere_actions)); src_table_query_info.prewhere_info->prewhere_actions.removeUnusedActions(); @@ -480,23 +488,23 @@ void StorageBuffer::read( } else { - if (query_info.prewhere_info) + if (query_info.row_level_filter) { ExpressionActionsSettings actions_settings(local_context); - - if (query_info.prewhere_info->row_level_filter) + auto actions = std::make_shared(query_info.row_level_filter->actions.clone(), actions_settings); + pipe_from_buffers.addSimpleTransform([&](const SharedHeader & header) { - auto actions = std::make_shared(query_info.prewhere_info->row_level_filter->clone(), actions_settings); - pipe_from_buffers.addSimpleTransform([&](const SharedHeader & header) - { - return std::make_shared( - header, - actions, - query_info.prewhere_info->row_level_column_name, - false); - }); - } + return std::make_shared( + header, + actions, + query_info.row_level_filter->column_name, + query_info.row_level_filter->do_remove_column); + }); + } + if (query_info.prewhere_info) + { + ExpressionActionsSettings actions_settings(local_context); auto actions = std::make_shared(query_info.prewhere_info->prewhere_actions.clone(), actions_settings); pipe_from_buffers.addSimpleTransform([&](const SharedHeader & header) { diff --git a/src/Storages/StorageDummy.cpp b/src/Storages/StorageDummy.cpp index 95878e69d3b1..dc98e4e4a0ff 100644 --- a/src/Storages/StorageDummy.cpp +++ b/src/Storages/StorageDummy.cpp @@ -57,7 +57,7 @@ ReadFromDummy::ReadFromDummy( const ContextPtr & context_, const StorageDummy & storage_) : SourceStepWithFilter(std::make_shared(SourceStepWithFilter::applyPrewhereActions( - storage_snapshot_->getSampleBlockForColumns(column_names_), query_info_.prewhere_info)), + storage_snapshot_->getSampleBlockForColumns(column_names_), query_info_.row_level_filter, query_info_.prewhere_info)), column_names_, query_info_, storage_snapshot_, diff --git a/src/Storages/StorageFile.cpp b/src/Storages/StorageFile.cpp index ea9fef7b4dd4..76269bd06f0c 100644 --- a/src/Storages/StorageFile.cpp +++ b/src/Storages/StorageFile.cpp @@ -1641,9 +1641,8 @@ void ReadFromFile::applyFilters(ActionDAGNodes added_filter_nodes) void ReadFromFile::updatePrewhereInfo(const PrewhereInfoPtr & prewhere_info_value) { - info = updateFormatPrewhereInfo(info, prewhere_info_value); + info = updateFormatPrewhereInfo(info, query_info.row_level_filter, prewhere_info_value); query_info.prewhere_info = prewhere_info_value; - prewhere_info = prewhere_info_value; output_header = std::make_shared(info.source_header); } @@ -1692,9 +1691,9 @@ void StorageFile::read( PrepareReadingFromFormatHiveParams {file_columns, hive_partition_columns_to_read_from_file_path.getNameToTypeMap()}); if (query_info.prewhere_info) - read_from_format_info = updateFormatPrewhereInfo(read_from_format_info, query_info.prewhere_info); + read_from_format_info = updateFormatPrewhereInfo(read_from_format_info, query_info.row_level_filter, query_info.prewhere_info); - bool need_only_count = (query_info.optimize_trivial_count || (read_from_format_info.requested_columns.empty() && !read_from_format_info.prewhere_info)) + bool need_only_count = (query_info.optimize_trivial_count || (read_from_format_info.requested_columns.empty() && !read_from_format_info.prewhere_info && !read_from_format_info.row_level_filter)) && context->getSettingsRef()[Setting::optimize_count_from_files]; auto reading = std::make_unique( @@ -1753,8 +1752,7 @@ void ReadFromFile::initializePipeline(QueryPipelineBuilder & pipeline, const Bui progress_callback(FileProgress(0, storage->total_bytes_to_read)); auto parser_shared_resources = std::make_shared(ctx->getSettingsRef(), num_streams); - auto format_filter_info = std::make_shared(filter_actions_dag, ctx, nullptr); - format_filter_info->prewhere_info = prewhere_info; + auto format_filter_info = std::make_shared(filter_actions_dag, ctx, nullptr, query_info.row_level_filter, query_info.prewhere_info); for (size_t i = 0; i < num_streams; ++i) { diff --git a/src/Storages/StorageURL.cpp b/src/Storages/StorageURL.cpp index b377fe5448d3..49035e56739d 100644 --- a/src/Storages/StorageURL.cpp +++ b/src/Storages/StorageURL.cpp @@ -1170,9 +1170,8 @@ void ReadFromURL::applyFilters(ActionDAGNodes added_filter_nodes) void ReadFromURL::updatePrewhereInfo(const PrewhereInfoPtr & prewhere_info_value) { - info = updateFormatPrewhereInfo(info, prewhere_info_value); + info = updateFormatPrewhereInfo(info, query_info.row_level_filter, prewhere_info_value); query_info.prewhere_info = prewhere_info_value; - prewhere_info = prewhere_info_value; output_header = std::make_shared(info.source_header); } @@ -1195,10 +1194,10 @@ void IStorageURLBase::read( /*supports_tuple_elements=*/ supports_prewhere, PrepareReadingFromFormatHiveParams {file_columns, hive_partition_columns_to_read_from_file_path.getNameToTypeMap()}); - if (query_info.prewhere_info) - read_from_format_info = updateFormatPrewhereInfo(read_from_format_info, query_info.prewhere_info); + if (query_info.prewhere_info || query_info.row_level_filter) + read_from_format_info = updateFormatPrewhereInfo(read_from_format_info, query_info.row_level_filter, query_info.prewhere_info); - bool need_only_count = (query_info.optimize_trivial_count || (read_from_format_info.requested_columns.empty() && !read_from_format_info.prewhere_info)) + bool need_only_count = (query_info.optimize_trivial_count || (read_from_format_info.requested_columns.empty() && !read_from_format_info.prewhere_info && !read_from_format_info.row_level_filter)) && local_context->getSettingsRef()[Setting::optimize_count_from_files]; auto read_post_data_callback = getReadPOSTDataCallback( @@ -1311,8 +1310,7 @@ void ReadFromURL::initializePipeline(QueryPipelineBuilder & pipeline, const Buil pipes.reserve(num_streams); auto parser_shared_resources = std::make_shared(settings, num_streams); - auto format_filter_info = std::make_shared(filter_actions_dag, context, nullptr); - format_filter_info->prewhere_info = prewhere_info; + auto format_filter_info = std::make_shared(filter_actions_dag, context, nullptr, query_info.row_level_filter, query_info.prewhere_info); for (size_t i = 0; i < num_streams; ++i) { @@ -1376,10 +1374,10 @@ void StorageURLWithFailover::read( /*supports_tuple_elements=*/ supports_prewhere, PrepareReadingFromFormatHiveParams {file_columns, hive_partition_columns_to_read_from_file_path.getNameToTypeMap()}); - if (query_info.prewhere_info) - read_from_format_info = updateFormatPrewhereInfo(read_from_format_info, query_info.prewhere_info); + if (query_info.prewhere_info || query_info.row_level_filter) + read_from_format_info = updateFormatPrewhereInfo(read_from_format_info, query_info.row_level_filter, query_info.prewhere_info); - bool need_only_count = (query_info.optimize_trivial_count || (read_from_format_info.requested_columns.empty() && !read_from_format_info.prewhere_info)) + bool need_only_count = (query_info.optimize_trivial_count || (read_from_format_info.requested_columns.empty() && !read_from_format_info.prewhere_info && !read_from_format_info.row_level_filter)) && local_context->getSettingsRef()[Setting::optimize_count_from_files]; auto read_post_data_callback = getReadPOSTDataCallback( diff --git a/src/Storages/prepareReadingFromFormat.cpp b/src/Storages/prepareReadingFromFormat.cpp index e84005511bb4..9a3dd83db417 100644 --- a/src/Storages/prepareReadingFromFormat.cpp +++ b/src/Storages/prepareReadingFromFormat.cpp @@ -246,11 +246,11 @@ ReadFromFormatInfo prepareReadingFromFormat( return info; } -ReadFromFormatInfo updateFormatPrewhereInfo(const ReadFromFormatInfo & info, const PrewhereInfoPtr & prewhere_info) +ReadFromFormatInfo updateFormatPrewhereInfo(const ReadFromFormatInfo & info, const FilterDAGInfoPtr & row_level_filter, const PrewhereInfoPtr & prewhere_info) { chassert(prewhere_info); - if (info.prewhere_info) + if (info.prewhere_info || info.row_level_filter) throw Exception(ErrorCodes::LOGICAL_ERROR, "updateFormatPrewhereInfo called more than once"); ReadFromFormatInfo new_info; @@ -258,7 +258,7 @@ ReadFromFormatInfo updateFormatPrewhereInfo(const ReadFromFormatInfo & info, con /// Removes columns that are only used as prewhere input. /// Adds prewhere result column if !remove_prewhere_column. - new_info.format_header = SourceStepWithFilter::applyPrewhereActions(info.format_header, prewhere_info); + new_info.format_header = SourceStepWithFilter::applyPrewhereActions(info.format_header, row_level_filter, prewhere_info); /// We assume that any format that supports prewhere also supports subset of subcolumns, so we /// don't need to replace subcolumns with their nested columns etc. @@ -367,7 +367,7 @@ ReadFromFormatInfo ReadFromFormatInfo::deserialize(IQueryPlanStep::Deserializati bool has_prewhere_info; readBinary(has_prewhere_info, ctx.in); if (has_prewhere_info) - result.prewhere_info = PrewhereInfo::deserialize(ctx); + result.prewhere_info = std::make_shared(PrewhereInfo::deserialize(ctx)); ctx.in >> "\n"; diff --git a/src/Storages/prepareReadingFromFormat.h b/src/Storages/prepareReadingFromFormat.h index 27951f1da8dd..19f766e45ddf 100644 --- a/src/Storages/prepareReadingFromFormat.h +++ b/src/Storages/prepareReadingFromFormat.h @@ -10,6 +10,9 @@ namespace DB struct PrewhereInfo; using PrewhereInfoPtr = std::shared_ptr; + struct FilterDAGInfo; + using FilterDAGInfoPtr = std::shared_ptr; + struct ReadFromFormatInfo { /// Header that will return Source from storage. @@ -40,6 +43,7 @@ namespace DB /// The list of hive partition columns. It shall be read from the path regardless if it is present in the file NamesAndTypesList hive_partition_columns_to_read_from_file_path; PrewhereInfoPtr prewhere_info; + FilterDAGInfoPtr row_level_filter; }; struct PrepareReadingFromFormatHiveParams @@ -70,7 +74,7 @@ namespace DB bool supports_tuple_elements = false, const PrepareReadingFromFormatHiveParams & hive_parameters = {}); - ReadFromFormatInfo updateFormatPrewhereInfo(const ReadFromFormatInfo & info, const PrewhereInfoPtr & prewhere_info); + ReadFromFormatInfo updateFormatPrewhereInfo(const ReadFromFormatInfo & info, const FilterDAGInfoPtr & row_level_filter, const PrewhereInfoPtr & prewhere_info); /// Returns the serialization hints from the insertion table (if it's set in the Context). SerializationInfoByName getSerializationHintsForFileLikeStorage(const StorageMetadataPtr & metadata_snapshot, const ContextPtr & context); diff --git a/tests/queries/0_stateless/02679_explain_merge_tree_prewhere_row_policy.reference b/tests/queries/0_stateless/02679_explain_merge_tree_prewhere_row_policy.reference index 4a4e338438b2..669d0fe248d3 100644 --- a/tests/queries/0_stateless/02679_explain_merge_tree_prewhere_row_policy.reference +++ b/tests/queries/0_stateless/02679_explain_merge_tree_prewhere_row_policy.reference @@ -19,7 +19,7 @@ Positions: 0 1 FUNCTION equals(id : 0, 5 :: 1) -> equals(id, 5) UInt8 : 2 Positions: 2 0 Row level filter - Row level filter column: greaterOrEquals(id, 5) + Row level filter column: greaterOrEquals(id, 5) (removed) Actions: INPUT : 0 -> id UInt64 : 0 COLUMN Const(UInt8) -> 5 UInt8 : 1 FUNCTION greaterOrEquals(id : 0, 5 :: 1) -> greaterOrEquals(id, 5) UInt8 : 2 @@ -49,7 +49,7 @@ Positions: 1 2 FUNCTION equals(id : 0, 5_UInt8 :: 1) -> equals(id, 5_UInt8) UInt8 : 2 Positions: 2 0 Row level filter - Row level filter column: greaterOrEquals(id, 5_UInt8) + Row level filter column: greaterOrEquals(id, 5_UInt8) (removed) Actions: INPUT : 0 -> id UInt64 : 0 COLUMN Const(UInt8) -> 5_UInt8 UInt8 : 1 FUNCTION greaterOrEquals(id : 0, 5_UInt8 :: 1) -> greaterOrEquals(id, 5_UInt8) UInt8 : 2 diff --git a/tests/queries/0_stateless/03591_optimize_prewhere_row_policy.reference b/tests/queries/0_stateless/03591_optimize_prewhere_row_policy.reference new file mode 100644 index 000000000000..ec5934037622 --- /dev/null +++ b/tests/queries/0_stateless/03591_optimize_prewhere_row_policy.reference @@ -0,0 +1,144 @@ +-- {echoOn} + +SET use_query_condition_cache = 0; +SET enable_parallel_replicas = 0; +DROP TABLE IF EXISTS 03591_test; +DROP ROW POLICY IF EXISTS 03591_rp ON 03591_test; +CREATE TABLE 03591_test (a Int32, b Int32) ENGINE=MergeTree ORDER BY tuple(); +INSERT INTO 03591_test VALUES (3, 1), (2, 2), (3, 2); +SELECT * FROM 03591_test; +3 1 +2 2 +3 2 +SELECT * FROM 03591_test WHERE throwIf(b=1, 'Should throw') SETTINGS optimize_move_to_prewhere = 1; -- {serverError FUNCTION_THROW_IF_VALUE_IS_NON_ZERO} +CREATE ROW POLICY 03591_rp ON 03591_test USING b=2 TO CURRENT_USER; +SELECT * FROM 03591_test; +2 2 +3 2 +-- Print plan with actions to make sure both a > 0 and b=2 are present in the prewhere section +EXPLAIN PLAN actions=1 SELECT * FROM 03591_test WHERE a > 0 SETTINGS optimize_move_to_prewhere = 1, allow_experimental_analyzer = 1; +Expression ((Project names + Projection)) +Actions: INPUT : 0 -> __table1.b Int32 : 0 + INPUT : 1 -> __table1.a Int32 : 1 + ALIAS __table1.b :: 0 -> b Int32 : 2 + ALIAS __table1.a :: 1 -> a Int32 : 0 +Positions: 0 2 + Expression ((WHERE + Change column names to column identifiers)) + Actions: INPUT : 0 -> b Int32 : 0 + INPUT : 1 -> a Int32 : 1 + ALIAS b :: 0 -> __table1.b Int32 : 2 + ALIAS a :: 1 -> __table1.a Int32 : 0 + Positions: 2 0 + ReadFromMergeTree (default.03591_test) + ReadType: Default + Parts: 1 + Granules: 1 + Prewhere info + Need filter: 1 + Prewhere filter + Prewhere filter column: greater(__table1.a, 0_UInt8) (removed) + Actions: INPUT : 0 -> a Int32 : 0 + COLUMN Const(UInt8) -> 0_UInt8 UInt8 : 1 + FUNCTION greater(a : 0, 0_UInt8 :: 1) -> greater(__table1.a, 0_UInt8) UInt8 : 2 + Positions: 0 2 + Row level filter + Row level filter column: equals(b, 2_UInt8) (removed) + Actions: INPUT : 0 -> b Int32 : 0 + COLUMN Const(UInt8) -> 2_UInt8 UInt8 : 1 + FUNCTION equals(b : 0, 2_UInt8 :: 1) -> equals(b, 2_UInt8) UInt8 : 2 + Positions: 2 0 +EXPLAIN PLAN actions=1 SELECT * FROM 03591_test WHERE a > 0 SETTINGS optimize_move_to_prewhere = 1, allow_experimental_analyzer = 0; +Expression ((Projection + Before ORDER BY)) +Actions: INPUT :: 0 -> a Int32 : 0 + INPUT :: 1 -> b Int32 : 1 +Positions: 0 1 + Expression (WHERE) + Actions: INPUT :: 0 -> a Int32 : 0 + INPUT :: 1 -> b Int32 : 1 + Positions: 0 1 + ReadFromMergeTree (default.03591_test) + ReadType: Default + Parts: 1 + Granules: 1 + Prewhere info + Need filter: 1 + Prewhere filter + Prewhere filter column: greater(a, 0) (removed) + Actions: INPUT : 0 -> a Int32 : 0 + COLUMN Const(UInt8) -> 0 UInt8 : 1 + FUNCTION greater(a : 0, 0 :: 1) -> greater(a, 0) UInt8 : 2 + Positions: 0 2 + Row level filter + Row level filter column: equals(b, 2) (removed) + Actions: INPUT : 0 -> b Int32 : 0 + COLUMN Const(UInt8) -> 2 UInt8 : 1 + FUNCTION equals(b : 0, 2 :: 1) -> equals(b, 2) UInt8 : 2 + Positions: 2 0 +SELECT * FROM 03591_test WHERE throwIf(b=1, 'Should not throw because b=1 is not visible to this user due to the b=2 row policy') SETTINGS optimize_move_to_prewhere = 1; +-- Print plan with actions to make sure a > 0, b = 2 and a = 3 are present in the prewhere section +EXPLAIN PLAN actions=1 SELECT * FROM 03591_test WHERE a > 0 SETTINGS optimize_move_to_prewhere = 1, additional_table_filters={'03591_test': 'a=3'}, allow_experimental_analyzer = 1; +Expression ((Project names + Projection)) +Actions: INPUT : 0 -> __table1.a Int32 : 0 + INPUT : 1 -> __table1.b Int32 : 1 + ALIAS __table1.a :: 0 -> a Int32 : 2 + ALIAS __table1.b :: 1 -> b Int32 : 0 +Positions: 2 0 + Expression (((WHERE + Change column names to column identifiers) + additional filter)) + Actions: INPUT : 1 -> b Int32 : 0 + INPUT : 0 -> a Int32 : 1 + ALIAS b :: 0 -> __table1.b Int32 : 2 + ALIAS a :: 1 -> __table1.a Int32 : 0 + Positions: 0 2 + ReadFromMergeTree (default.03591_test) + ReadType: Default + Parts: 1 + Granules: 1 + Prewhere info + Need filter: 1 + Prewhere filter + Prewhere filter column: and(equals(a, 3_UInt8), greater(__table1.a, 0_UInt8)) (removed) + Actions: INPUT : 0 -> a Int32 : 0 + COLUMN Const(UInt8) -> 3_UInt8 UInt8 : 1 + COLUMN Const(UInt8) -> 0_UInt8 UInt8 : 2 + FUNCTION equals(a : 0, 3_UInt8 :: 1) -> equals(a, 3_UInt8) UInt8 : 3 + FUNCTION greater(a : 0, 0_UInt8 :: 2) -> greater(__table1.a, 0_UInt8) UInt8 : 1 + FUNCTION and(equals(a, 3_UInt8) :: 3, greater(__table1.a, 0_UInt8) :: 1) -> and(equals(a, 3_UInt8), greater(__table1.a, 0_UInt8)) UInt8 : 2 + Positions: 0 2 + Row level filter + Row level filter column: equals(b, 2_UInt8) (removed) + Actions: INPUT : 0 -> b Int32 : 0 + COLUMN Const(UInt8) -> 2_UInt8 UInt8 : 1 + FUNCTION equals(b : 0, 2_UInt8 :: 1) -> equals(b, 2_UInt8) UInt8 : 2 + Positions: 2 0 +EXPLAIN PLAN actions=1 SELECT * FROM 03591_test WHERE a > 0 SETTINGS optimize_move_to_prewhere = 1, additional_table_filters={'03591_test': 'a=3'}, allow_experimental_analyzer = 0; +Expression ((Projection + Before ORDER BY)) +Actions: INPUT :: 0 -> a Int32 : 0 + INPUT :: 1 -> b Int32 : 1 +Positions: 0 1 + Expression ((WHERE + Additional filter)) + Actions: INPUT :: 0 -> a Int32 : 0 + INPUT :: 1 -> b Int32 : 1 + Positions: 0 1 + ReadFromMergeTree (default.03591_test) + ReadType: Default + Parts: 1 + Granules: 1 + Prewhere info + Need filter: 1 + Prewhere filter + Prewhere filter column: and(equals(a, 3), greater(a, 0)) (removed) + Actions: INPUT : 0 -> a Int32 : 0 + COLUMN Const(UInt8) -> 3 UInt8 : 1 + COLUMN Const(UInt8) -> 0 UInt8 : 2 + FUNCTION equals(a : 0, 3 :: 1) -> equals(a, 3) UInt8 : 3 + FUNCTION greater(a : 0, 0 :: 2) -> greater(a, 0) UInt8 : 1 + FUNCTION and(equals(a, 3) :: 3, greater(a, 0) :: 1) -> and(equals(a, 3), greater(a, 0)) UInt8 : 2 + Positions: 0 2 + Row level filter + Row level filter column: equals(b, 2) (removed) + Actions: INPUT : 0 -> b Int32 : 0 + COLUMN Const(UInt8) -> 2 UInt8 : 1 + FUNCTION equals(b : 0, 2 :: 1) -> equals(b, 2) UInt8 : 2 + Positions: 2 0 +DROP ROW POLICY 03591_rp ON 03591_test; +SELECT * FROM 03591_test WHERE throwIf(b=2, 'Should throw') SETTINGS optimize_move_to_prewhere = 1; -- {serverError FUNCTION_THROW_IF_VALUE_IS_NON_ZERO} diff --git a/tests/queries/0_stateless/03591_optimize_prewhere_row_policy.sql b/tests/queries/0_stateless/03591_optimize_prewhere_row_policy.sql new file mode 100644 index 000000000000..b010fadff97b --- /dev/null +++ b/tests/queries/0_stateless/03591_optimize_prewhere_row_policy.sql @@ -0,0 +1,34 @@ +-- {echoOn} + +SET use_query_condition_cache = 0; +SET enable_parallel_replicas = 0; + +DROP TABLE IF EXISTS 03591_test; + +DROP ROW POLICY IF EXISTS 03591_rp ON 03591_test; + +CREATE TABLE 03591_test (a Int32, b Int32) ENGINE=MergeTree ORDER BY tuple(); + +INSERT INTO 03591_test VALUES (3, 1), (2, 2), (3, 2); + +SELECT * FROM 03591_test; + +SELECT * FROM 03591_test WHERE throwIf(b=1, 'Should throw') SETTINGS optimize_move_to_prewhere = 1; -- {serverError FUNCTION_THROW_IF_VALUE_IS_NON_ZERO} + +CREATE ROW POLICY 03591_rp ON 03591_test USING b=2 TO CURRENT_USER; + +SELECT * FROM 03591_test; + +-- Print plan with actions to make sure both a > 0 and b=2 are present in the prewhere section +EXPLAIN PLAN actions=1 SELECT * FROM 03591_test WHERE a > 0 SETTINGS optimize_move_to_prewhere = 1, allow_experimental_analyzer = 1; +EXPLAIN PLAN actions=1 SELECT * FROM 03591_test WHERE a > 0 SETTINGS optimize_move_to_prewhere = 1, allow_experimental_analyzer = 0; + +SELECT * FROM 03591_test WHERE throwIf(b=1, 'Should not throw because b=1 is not visible to this user due to the b=2 row policy') SETTINGS optimize_move_to_prewhere = 1; + +-- Print plan with actions to make sure a > 0, b = 2 and a = 3 are present in the prewhere section +EXPLAIN PLAN actions=1 SELECT * FROM 03591_test WHERE a > 0 SETTINGS optimize_move_to_prewhere = 1, additional_table_filters={'03591_test': 'a=3'}, allow_experimental_analyzer = 1; +EXPLAIN PLAN actions=1 SELECT * FROM 03591_test WHERE a > 0 SETTINGS optimize_move_to_prewhere = 1, additional_table_filters={'03591_test': 'a=3'}, allow_experimental_analyzer = 0; + +DROP ROW POLICY 03591_rp ON 03591_test; + +SELECT * FROM 03591_test WHERE throwIf(b=2, 'Should throw') SETTINGS optimize_move_to_prewhere = 1; -- {serverError FUNCTION_THROW_IF_VALUE_IS_NON_ZERO} diff --git a/tests/queries/0_stateless/03641_analyzer_issue_85834.reference b/tests/queries/0_stateless/03641_analyzer_issue_85834.reference new file mode 100644 index 000000000000..d81cc0710eb6 --- /dev/null +++ b/tests/queries/0_stateless/03641_analyzer_issue_85834.reference @@ -0,0 +1 @@ +42 diff --git a/tests/queries/0_stateless/03641_analyzer_issue_85834.sql b/tests/queries/0_stateless/03641_analyzer_issue_85834.sql new file mode 100644 index 000000000000..a8903f82508b --- /dev/null +++ b/tests/queries/0_stateless/03641_analyzer_issue_85834.sql @@ -0,0 +1,14 @@ +-- https://github.com/ClickHouse/ClickHouse/issues/85834 + +DROP TABLE IF EXISTS test_generic_events_all; + +CREATE TABLE test_generic_events_all (APIKey UInt8, SessionType UInt8) ENGINE = MergeTree() PARTITION BY APIKey ORDER BY tuple(); +INSERT INTO test_generic_events_all VALUES( 42, 42 ); +ALTER TABLE test_generic_events_all ADD COLUMN OperatingSystem UInt64 DEFAULT 42; + +CREATE ROW POLICY rp ON test_generic_events_all USING APIKey>35 TO CURRENT_USER; + +SELECT OperatingSystem +FROM test_generic_events_all +PREWHERE APIKey = 42 +SETTINGS additional_table_filters = {'test_generic_events_all':'APIKey > 40'}; From 6a87a33976ca05a76761cca4605aa732b427daae Mon Sep 17 00:00:00 2001 From: Mikhail Koviazin Date: Wed, 5 Aug 2026 12:23:41 +0200 Subject: [PATCH 2/6] Port upstream follow-ups to the row-level-filter split (#87303), and fix two-step PREWHERE in the Parquet v3 reader PR #1345 backported upstream #87303, which lifts row-level security out of PrewhereInfo into SelectQueryInfo::row_level_filter. Upstream then had to fix several places that were left reading row-level security off prewhere_info. None of those fixes are in #1345, so backport them here, and fix a crash that #87303 makes reachable on this branch. Parquet::Reader::applyPrewhere could not run two filtering steps. Every block it assembles holds rows_pass rows - the count surviving all previous steps - but the function got that wrong in two ways: * Columns were materialized lazily per step via formOutputColumn, which for a primitive column takes the decoded subchunk. The decoders only saw the filter as it stood before any step ran, so that subchunk still holds the pre-filter row count, while the per-step filtering only shrinks what is already in row_subgroup.output - pending subchunks are never touched. A column first needed by the second step therefore arrived one filter generation behind and tripped chassert(filter.size() == row_subgroup.filter.rows_pass). Materialize every step's inputs before running any step so they are filtered in lockstep. formOutputColumn moves out of the subchunk, so this only shifts ownership earlier and does not change peak memory. * addDummyColumnWithRowCount was passed rows_total, and it asserts that every column already in the block has exactly that many rows. That held only because of the bug above, which left the second step's column unfiltered; with the columns correctly at rows_pass it fails instead. Pass rows_pass, which is the row count the block actually has at every step. Planner: #87303 also replaced the pre-existing add_filter gate (canMoveConditionsToPrewhere && optimize_move_to_prewhere && supportedPrewhereColumns->contains(...) && !has_table_virtual_column) with a bare supportsPrewhere() for the row policy, dropping the supportedPrewhereColumns() check. StorageFile, IStorageURLBase and StorageObjectStorage all restrict prewhere to physical columns, while a row policy's filter column is usually an expression name, so before #87303 those storages routed the policy to WHERE. Upstream release branches do not notice the loss because input_format_parquet_use_native_reader_v3 defaults to false there, making supportsPrewhere() false for Parquet anyway; this branch enables that reader by default. Restore the check for the row policy only. MergeTree returns nullopt and is unaffected, so the move-to-prewhere fix that motivates the backport is preserved, and its prewhere steps are executed by MergeTreeRangeReader rather than by the code above. The two fixes cover different shapes. The guard diverts expression-valued policies, which is the common case. A policy whose condition is a bare column (USING flag) is named after a physical column, passes the guard, and still reaches the reader: with an explicit PREWHERE that shape aborted on this branch even before #87303, and after #87303 a plain WHERE moved into prewhere aborts too, so the reader fix is needed as well. The applyPrewhere limitation is present on every upstream release branch carrying #87303 and was only fixed on master, by the multistage-prewhere redesign (#93542); the fix here is local to this branch and worth offering upstream separately. updateFormatPrewhereInfo, two upstream commits that must go together: * 8ddee54956a, "Fix exception in updateFormatPrewhereInfo when only row_level_filter is set": the assertion still required prewhere_info, but every caller now invokes the function when either filter is set, so a row policy without PREWHERE on an object storage / File / URL table tripped it. row_level_filter was also never stored into the new ReadFromFormatInfo and got lost. * 774b56b47d8, "Fix updateFormatPrewhereInfo called more than once when row policy and prewhere are both active": storing row_level_filter (above) makes the duplicate-call guard reject a legitimate second call. When a table has a row policy and the optimizer later pushes WHERE into PREWHERE, the function runs twice - once from read() for the row_level_filter, once from updatePrewhereInfo() for both. Guard only against duplicate prewhere_info, and skip re-applying a row_level_filter that a previous call already applied. * 92b0d17a3, "Consider row level filter for read in order optimization": the row-level filter expression was no longer appended to the sorting DAG, so its fixed columns were not recognised and read-in-order was skipped; the limit was also no longer reset despite filtering being present. * 25c22b71d, 6b35e2761, "Fix row policy filter error when using projections" / "Fix for NOT_FOUND_COLUMN_IN_BLOCK when selecting from projections": projection_query_info kept row_level_filter while projectionsCommon already folds it into the projection prewhere, so the filter was applied twice and failed on the projection's block layout. Tests come from the upstream commits verbatim, except: * 04490_row_policy_parquet_v3_two_prewhere_steps is new and specific to this branch: it covers both shapes above on a File(Parquet) table with the v3 reader, and pins the routing guard for expression-valued policies. * 03800_projection_row_policy_filter_column.reference: its EXPLAIN indexes=1 output has "Ranges: 1" indented two spaces deeper on this branch, because ReadFromMergeTree::describeIndexes still prints it with an extra indent level here. Regenerated against this branch; no other byte differs. Co-Authored-By: Claude Opus 5 (1M context) --- src/Planner/PlannerJoinTree.cpp | 11 +++- .../Formats/Impl/Parquet/Reader.cpp | 23 +++++++- .../Optimizations/optimizeReadInOrder.cpp | 10 ++++ .../optimizeUseAggregateProjection.cpp | 2 + .../optimizeUseNormalProjection.cpp | 5 ++ src/Storages/prepareReadingFromFormat.cpp | 10 +++- ...jection_row_policy_filter_column.reference | 17 ++++++ ...00_projection_row_policy_filter_column.sql | 48 ++++++++++++++++ ...er_in_read_in_order_optimization.reference | 4 ++ ...w_filter_in_read_in_order_optimization.sql | 39 +++++++++++++ .../04053_row_policy_object_storage.reference | 7 +++ .../04053_row_policy_object_storage.sh | 46 +++++++++++++++ ...cy_parquet_v3_two_prewhere_steps.reference | 20 +++++++ ...w_policy_parquet_v3_two_prewhere_steps.sql | 56 +++++++++++++++++++ ...licy_projection_not_found_column.reference | 4 ++ ...row_policy_projection_not_found_column.sql | 37 ++++++++++++ 16 files changed, 333 insertions(+), 6 deletions(-) create mode 100644 tests/queries/0_stateless/03800_projection_row_policy_filter_column.reference create mode 100644 tests/queries/0_stateless/03800_projection_row_policy_filter_column.sql create mode 100644 tests/queries/0_stateless/03927_row_filter_in_read_in_order_optimization.reference create mode 100644 tests/queries/0_stateless/03927_row_filter_in_read_in_order_optimization.sql create mode 100644 tests/queries/0_stateless/04053_row_policy_object_storage.reference create mode 100755 tests/queries/0_stateless/04053_row_policy_object_storage.sh create mode 100644 tests/queries/0_stateless/04490_row_policy_parquet_v3_two_prewhere_steps.reference create mode 100644 tests/queries/0_stateless/04490_row_policy_parquet_v3_two_prewhere_steps.sql create mode 100644 tests/queries/0_stateless/4060_row_policy_projection_not_found_column.reference create mode 100644 tests/queries/0_stateless/4060_row_policy_projection_not_found_column.sql diff --git a/src/Planner/PlannerJoinTree.cpp b/src/Planner/PlannerJoinTree.cpp index 1488d1ea5118..20c9d6df6f16 100644 --- a/src/Planner/PlannerJoinTree.cpp +++ b/src/Planner/PlannerJoinTree.cpp @@ -986,7 +986,16 @@ JoinTreeQueryPlan buildQueryPlanForTableExpression(QueryTreeNodePtr table_expres { table_expression_data.setRowLevelFilterActions(row_policy_filter_info->actions.clone()); /// TODO: Never put row-level security filter in WHERE clause for storages that do not support PREWHERE to avoid merging of filters. - if (storage->supportsPrewhere()) + /// Upstream checks only supportsPrewhere() here, having dropped the supportedPrewhereColumns() + /// check that the pre-#87303 add_filter lambda applied to every filter. Keep that check for the + /// row policy: File/URL/ObjectStorage restrict prewhere to physical columns, and a row policy's + /// filter column is an expression name, so those storages sent the policy to WHERE before. The + /// Parquet v3 reader cannot run two filtering steps (Reader::applyPrewhere asserts + /// filter.size() == rows_pass, but forms a later step's column from its unfiltered subchunk), and + /// antalya enables that reader by default, so routing the policy into prewhere there aborts. + auto supported_prewhere_columns = storage->supportedPrewhereColumns(); + if (storage->supportsPrewhere() + && (!supported_prewhere_columns || supported_prewhere_columns->contains(row_policy_filter_info->column_name))) row_level_filter = std::make_shared(std::move(*row_policy_filter_info)); else where_filters.emplace_back(std::move(*row_policy_filter_info), "Row-level security filter"); diff --git a/src/Processors/Formats/Impl/Parquet/Reader.cpp b/src/Processors/Formats/Impl/Parquet/Reader.cpp index 33b822f9e1f6..7c1672e8e517 100644 --- a/src/Processors/Formats/Impl/Parquet/Reader.cpp +++ b/src/Processors/Formats/Impl/Parquet/Reader.cpp @@ -2087,6 +2087,24 @@ MutableColumnPtr Reader::formOutputColumn(RowSubgroup & row_subgroup, size_t out void Reader::applyPrewhere(RowSubgroup & row_subgroup, const RowGroup & row_group) { + /// Materialize the input columns of every step before running any step. formOutputColumn takes + /// the decoded subchunk, which still has rows_total rows because the decoders only saw the + /// filter as it stood before any step ran. A column first formed by a later step would + /// therefore disagree with the columns an earlier step already filtered down to rows_pass, and + /// the filtering below only shrinks what is already in row_subgroup.output - pending subchunks + /// are never touched. Forming them all up front keeps every step's block at rows_pass. + /// + /// rows_pass, not rows_total, is the row count every block here has: the decoded subchunks hold + /// rows_pass rows (see decodePrimitiveColumn), and each step shrinks them further. + for (const PrewhereStep & step : prewhere_steps) + for (size_t output_idx : step.input_column_idxs) + { + const auto & output_info = output_columns.at(output_idx); + auto & col = row_subgroup.output.at(output_info.idx_in_output_block.value()); + if (!col) + col = formOutputColumn(row_subgroup, output_idx, row_subgroup.filter.rows_pass); + } + for (size_t step_idx = 0; step_idx < prewhere_steps.size(); ++step_idx) { const PrewhereStep & step = prewhere_steps.at(step_idx); @@ -2096,11 +2114,12 @@ void Reader::applyPrewhere(RowSubgroup & row_subgroup, const RowGroup & row_grou { const auto & output_info = output_columns.at(output_idx); auto & col = row_subgroup.output.at(output_info.idx_in_output_block.value()); + /// Unreachable: the loop above materialized every step's inputs. if (!col) - col = formOutputColumn(row_subgroup, output_idx, row_subgroup.filter.rows_total); + col = formOutputColumn(row_subgroup, output_idx, row_subgroup.filter.rows_pass); block.insert({col, output_info.type, output_info.name}); } - addDummyColumnWithRowCount(block, row_subgroup.filter.rows_total); + addDummyColumnWithRowCount(block, row_subgroup.filter.rows_pass); step.actions.execute(block); diff --git a/src/Processors/QueryPlan/Optimizations/optimizeReadInOrder.cpp b/src/Processors/QueryPlan/Optimizations/optimizeReadInOrder.cpp index c0e92c2989aa..d7e4fc757120 100644 --- a/src/Processors/QueryPlan/Optimizations/optimizeReadInOrder.cpp +++ b/src/Processors/QueryPlan/Optimizations/optimizeReadInOrder.cpp @@ -189,6 +189,16 @@ void buildSortingDAG(QueryPlan::Node & node, std::optional & dag, Fi if (const auto * filter_expression = dag->tryFindInOutputs(prewhere_info->prewhere_column_name)) appendFixedColumnsFromFilterExpression(*filter_expression, fixed_columns); + } + if (const auto row_level_filter = reading->getRowLevelFilter()) + { + /// Should ignore limit if there is filtering. + limit = 0; + + appendExpression(dag, row_level_filter->actions); + if (const auto * filter_expression = dag->tryFindInOutputs(row_level_filter->column_name)) + appendFixedColumnsFromFilterExpression(*filter_expression, fixed_columns); + } return; } diff --git a/src/Processors/QueryPlan/Optimizations/optimizeUseAggregateProjection.cpp b/src/Processors/QueryPlan/Optimizations/optimizeUseAggregateProjection.cpp index c746501fa145..38ed8d9b7905 100644 --- a/src/Processors/QueryPlan/Optimizations/optimizeUseAggregateProjection.cpp +++ b/src/Processors/QueryPlan/Optimizations/optimizeUseAggregateProjection.cpp @@ -626,6 +626,7 @@ std::optional optimizeUseAggregateProjections( auto projection_query_info = query_info; projection_query_info.prewhere_info = nullptr; + projection_query_info.row_level_filter = nullptr; projection_query_info.filter_actions_dag = std::make_unique(candidate.dag.clone()); bool analyzed = analyzeProjectionCandidate( @@ -788,6 +789,7 @@ std::optional optimizeUseAggregateProjections( auto proj_snapshot = std::make_shared(storage_snapshot->storage, best_candidate->projection->metadata); auto projection_query_info = query_info; projection_query_info.prewhere_info = nullptr; + projection_query_info.row_level_filter = nullptr; projection_query_info.filter_actions_dag = nullptr; projection_reading = reader.readFromParts( diff --git a/src/Processors/QueryPlan/Optimizations/optimizeUseNormalProjection.cpp b/src/Processors/QueryPlan/Optimizations/optimizeUseNormalProjection.cpp index fab5f1a26b7b..a7f6dab797e8 100644 --- a/src/Processors/QueryPlan/Optimizations/optimizeUseNormalProjection.cpp +++ b/src/Processors/QueryPlan/Optimizations/optimizeUseNormalProjection.cpp @@ -216,6 +216,11 @@ std::optional optimizeUseNormalProjections( bool optimize_use_projection_filtering = context->getSettingsRef()[Setting::optimize_use_projection_filtering]; auto projection_query_info = query_info; projection_query_info.prewhere_info = nullptr; + /// Clear row_level_filter - it will be included in the prewhere_info below via splitAndFillPrewhereInfo. + /// The QueryDAG::build() already collected the RLS filter into query.dag/filter_nodes. + /// Keeping row_level_filter here would cause it to be applied twice and fail because + /// the projection's sample block doesn't necessarily match the original table's layout. + projection_query_info.row_level_filter = nullptr; if (query.dag) projection_query_info.filter_actions_dag = std::make_unique(query.dag->clone()); auto empty_mutations_snapshot = reading->getMutationsSnapshot()->cloneEmpty(); diff --git a/src/Storages/prepareReadingFromFormat.cpp b/src/Storages/prepareReadingFromFormat.cpp index 9a3dd83db417..0522509a5244 100644 --- a/src/Storages/prepareReadingFromFormat.cpp +++ b/src/Storages/prepareReadingFromFormat.cpp @@ -248,17 +248,21 @@ ReadFromFormatInfo prepareReadingFromFormat( ReadFromFormatInfo updateFormatPrewhereInfo(const ReadFromFormatInfo & info, const FilterDAGInfoPtr & row_level_filter, const PrewhereInfoPtr & prewhere_info) { - chassert(prewhere_info); + chassert(prewhere_info || row_level_filter); - if (info.prewhere_info || info.row_level_filter) + if (info.prewhere_info) throw Exception(ErrorCodes::LOGICAL_ERROR, "updateFormatPrewhereInfo called more than once"); ReadFromFormatInfo new_info; new_info.prewhere_info = prewhere_info; + new_info.row_level_filter = row_level_filter; /// Removes columns that are only used as prewhere input. /// Adds prewhere result column if !remove_prewhere_column. - new_info.format_header = SourceStepWithFilter::applyPrewhereActions(info.format_header, row_level_filter, prewhere_info); + /// If row_level_filter was already applied in a previous call, don't re-apply it; + /// only apply the new prewhere_info on top. + new_info.format_header = SourceStepWithFilter::applyPrewhereActions( + info.format_header, info.row_level_filter ? nullptr : row_level_filter, prewhere_info); /// We assume that any format that supports prewhere also supports subset of subcolumns, so we /// don't need to replace subcolumns with their nested columns etc. diff --git a/tests/queries/0_stateless/03800_projection_row_policy_filter_column.reference b/tests/queries/0_stateless/03800_projection_row_policy_filter_column.reference new file mode 100644 index 000000000000..1fbd655702ca --- /dev/null +++ b/tests/queries/0_stateless/03800_projection_row_policy_filter_column.reference @@ -0,0 +1,17 @@ +Expression (Project names) + Sorting (Sorting for ORDER BY) + Expression ((Before ORDER BY + Projection)) + Aggregating + Expression + ReadFromMergeTree (proj_by_data) + Indexes: + PrimaryKey + Keys: + tenant_id + Condition: (tenant_id in [\'tenant_A\', \'tenant_A\']) + Parts: 1/1 + Granules: 1/1 + Search Algorithm: binary search + Ranges: 1 +item_1 2 +item_2 1 diff --git a/tests/queries/0_stateless/03800_projection_row_policy_filter_column.sql b/tests/queries/0_stateless/03800_projection_row_policy_filter_column.sql new file mode 100644 index 000000000000..5a763295fb5c --- /dev/null +++ b/tests/queries/0_stateless/03800_projection_row_policy_filter_column.sql @@ -0,0 +1,48 @@ +DROP TABLE IF EXISTS test_rls_projection; +DROP ROW POLICY IF EXISTS rls_policy ON test_rls_projection; + +CREATE TABLE test_rls_projection +( + id UInt64, + tenant_id String, + data String, + + PROJECTION proj_by_data + ( + SELECT * ORDER BY tenant_id, data, id + ) +) +ENGINE = MergeTree() +ORDER BY (tenant_id, id); + +INSERT INTO test_rls_projection VALUES + (1, 'tenant_A', 'item_1'), + (2, 'tenant_A', 'item_2'), + (3, 'tenant_A', 'item_1'), + (4, 'tenant_B', 'item_1'), + (5, 'tenant_B', 'item_3'); + +ALTER TABLE test_rls_projection MATERIALIZE PROJECTION proj_by_data; + +CREATE ROW POLICY rls_policy ON test_rls_projection +FOR SELECT +USING tenant_id = 'tenant_A' +TO default; + +-- Verify projection is used +EXPLAIN indexes = 1 +SELECT data, count() as cnt +FROM test_rls_projection +GROUP BY data +ORDER BY data +SETTINGS force_optimize_projection = 1; + +-- Query without tenant_id in SELECT - should work with RLS filter applied via projection +SELECT data, count() as cnt +FROM test_rls_projection +GROUP BY data +ORDER BY data +SETTINGS force_optimize_projection = 1; + +DROP ROW POLICY rls_policy ON test_rls_projection; +DROP TABLE test_rls_projection; diff --git a/tests/queries/0_stateless/03927_row_filter_in_read_in_order_optimization.reference b/tests/queries/0_stateless/03927_row_filter_in_read_in_order_optimization.reference new file mode 100644 index 000000000000..86daf281a8f7 --- /dev/null +++ b/tests/queries/0_stateless/03927_row_filter_in_read_in_order_optimization.reference @@ -0,0 +1,4 @@ +-- Row policy only (no explicit key1 in WHERE) +ReadType: InOrder +-- Explicit key1 in WHERE clause +ReadType: InOrder diff --git a/tests/queries/0_stateless/03927_row_filter_in_read_in_order_optimization.sql b/tests/queries/0_stateless/03927_row_filter_in_read_in_order_optimization.sql new file mode 100644 index 000000000000..c47de047d82c --- /dev/null +++ b/tests/queries/0_stateless/03927_row_filter_in_read_in_order_optimization.sql @@ -0,0 +1,39 @@ +SET optimize_read_in_order = 1; +DROP TABLE IF EXISTS t_row_policy_rio; + +CREATE TABLE t_row_policy_rio +( + key1 UInt64, + key2 UInt64 +) +ENGINE = MergeTree +ORDER BY (key1, key2); + +INSERT INTO t_row_policy_rio +SELECT c1 % 10, c2 +FROM generateRandom('c1 UInt64, c2 UInt64', 42) +LIMIT 1000; + +CREATE ROW POLICY test_policy ON t_row_policy_rio USING key1 = 0 TO ALL; + +SELECT '-- Row policy only (no explicit key1 in WHERE)'; + +SELECT trim(explain) FROM ( +EXPLAIN actions = 1, indexes = 1 +SELECT key2 FROM t_row_policy_rio +ORDER BY key2 +LIMIT 100 +) WHERE explain LIKE '%ReadType%'; + +SELECT '-- Explicit key1 in WHERE clause'; + +SELECT trim(explain) FROM ( +EXPLAIN actions = 1, indexes = 1 +SELECT key2 FROM t_row_policy_rio +WHERE key1 = 0 +ORDER BY key2 +LIMIT 100 +) WHERE explain LIKE '%ReadType%'; + +DROP ROW POLICY test_policy ON t_row_policy_rio; +DROP TABLE t_row_policy_rio; diff --git a/tests/queries/0_stateless/04053_row_policy_object_storage.reference b/tests/queries/0_stateless/04053_row_policy_object_storage.reference new file mode 100644 index 000000000000..7f4f9bef5bd9 --- /dev/null +++ b/tests/queries/0_stateless/04053_row_policy_object_storage.reference @@ -0,0 +1,7 @@ +--- Row policy filters URL Parquet table --- +1 a +2 b +--- Row policy with WHERE on URL Parquet table --- +1 a +--- Row policy count on URL Parquet table --- +2 diff --git a/tests/queries/0_stateless/04053_row_policy_object_storage.sh b/tests/queries/0_stateless/04053_row_policy_object_storage.sh new file mode 100755 index 000000000000..e66ff60e8ad6 --- /dev/null +++ b/tests/queries/0_stateless/04053_row_policy_object_storage.sh @@ -0,0 +1,46 @@ +#!/usr/bin/env bash +# Tags: no-replicated-database, no-fasttest + +# Regression test: row policy on a URL table with Parquet format caused +# "Logical error: 'prewhere_info'" because updateFormatPrewhereInfo asserted +# prewhere_info was non-null, but only row_level_filter was set. + +CURDIR=$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd) +# shellcheck source=../shell_config.sh +. "$CURDIR"/../shell_config.sh + +user="user04053_${CLICKHOUSE_DATABASE}_$RANDOM" +db=${CLICKHOUSE_DATABASE} + +${CLICKHOUSE_CLIENT} < Date: Wed, 5 Aug 2026 15:32:45 +0200 Subject: [PATCH 3/6] fix test references after the changes --- .../04099_row_policy_trivial_limit_threads.reference | 4 ++-- .../04201_trivial_count_with_additional_filter.reference | 2 ++ 2 files changed, 4 insertions(+), 2 deletions(-) diff --git a/tests/queries/0_stateless/04099_row_policy_trivial_limit_threads.reference b/tests/queries/0_stateless/04099_row_policy_trivial_limit_threads.reference index e6482e7b382f..ac70543519a7 100644 --- a/tests/queries/0_stateless/04099_row_policy_trivial_limit_threads.reference +++ b/tests/queries/0_stateless/04099_row_policy_trivial_limit_threads.reference @@ -18,5 +18,5 @@ row_policy_always_true additional_filter (Limit) Limit 4 → 4 - (ReadFromMergeTree) - MergeTreeSelect(pool: ReadPool, algorithm: Thread) × 4 0 → 1 + (ReadFromMergeTree) + MergeTreeSelect(pool: ReadPool, algorithm: Thread) × 4 0 → 1 diff --git a/tests/queries/0_stateless/04201_trivial_count_with_additional_filter.reference b/tests/queries/0_stateless/04201_trivial_count_with_additional_filter.reference index e3d9c95e5455..3dc11440b1e2 100644 --- a/tests/queries/0_stateless/04201_trivial_count_with_additional_filter.reference +++ b/tests/queries/0_stateless/04201_trivial_count_with_additional_filter.reference @@ -2,6 +2,8 @@ baseline (ReadFromPreparedSource) SourceFromSingleChunk 0 → 1 filtered_this_table +(Filter) +FilterTransform (ReadFromMergeTree) MergeTreeSelect(pool: ReadPoolInOrder, algorithm: InOrder) 0 → 1 filter_other_table From 46d1b1fdc1f90a54ec43058278700363461844b1 Mon Sep 17 00:00:00 2001 From: Mikhail Koviazin Date: Thu, 6 Aug 2026 08:03:24 +0200 Subject: [PATCH 4/6] fixed a test reference --- ...04098_row_policy_disjunction_optimization.reference | 4 ++-- .../04098_row_policy_disjunction_optimization.sql | 10 +++++----- 2 files changed, 7 insertions(+), 7 deletions(-) diff --git a/tests/queries/0_stateless/04098_row_policy_disjunction_optimization.reference b/tests/queries/0_stateless/04098_row_policy_disjunction_optimization.reference index 30678f7b7969..67d82bc6da81 100644 --- a/tests/queries/0_stateless/04098_row_policy_disjunction_optimization.reference +++ b/tests/queries/0_stateless/04098_row_policy_disjunction_optimization.reference @@ -1,7 +1,7 @@ old analyzer -Filter column: in(id, (1, 3, 5)) (removed) + Row level filter column: in(id, (1, 3, 5)) (removed) new analyzer -Filter column: in(id, __set_UInt64_) (removed) + Row level filter column: in(id, __set_UInt64_) (removed) result 1 3 diff --git a/tests/queries/0_stateless/04098_row_policy_disjunction_optimization.sql b/tests/queries/0_stateless/04098_row_policy_disjunction_optimization.sql index 2e8f5a9b76bc..e97074ad7079 100644 --- a/tests/queries/0_stateless/04098_row_policy_disjunction_optimization.sql +++ b/tests/queries/0_stateless/04098_row_policy_disjunction_optimization.sql @@ -18,14 +18,14 @@ CREATE ROW POLICY 04098_p3 ON t_row_policy_or USING id = 5 AS permissive TO ALL; SET enable_analyzer = 0; SELECT 'old analyzer'; -SELECT trim(BOTH ' ' FROM explain) FROM (EXPLAIN actions = 1 SELECT id FROM t_row_policy_or SETTINGS optimize_move_to_prewhere = 0) -WHERE explain LIKE '%Filter column: in(%'; +SELECT explain FROM (EXPLAIN actions = 1 SELECT id FROM t_row_policy_or) +WHERE explain LIKE '%Row level filter column:%'; SET enable_analyzer = 1; SELECT 'new analyzer'; -SELECT trim(BOTH ' ' FROM replaceRegexpOne(explain, '__set_UInt64_\\d+_\\d+', '__set_UInt64_')) -FROM (EXPLAIN actions = 1 SELECT id FROM t_row_policy_or SETTINGS optimize_move_to_prewhere = 0) -WHERE explain LIKE '%Filter column: in(%'; +SELECT replaceRegexpOne(explain, '__set_UInt64_\\d+_\\d+', '__set_UInt64_') +FROM (EXPLAIN actions = 1 SELECT id FROM t_row_policy_or) +WHERE explain LIKE '%Row level filter column:%'; SELECT 'result'; SELECT id FROM t_row_policy_or ORDER BY id; From cb8ac0e265af19f059c8eaf92de9b25dc040d6b2 Mon Sep 17 00:00:00 2001 From: Mikhail Koviazin Date: Thu, 6 Aug 2026 12:00:21 +0200 Subject: [PATCH 5/6] Take the current upstream versions of 03800 and 03927 Both tests were imported from the commits that introduced them (25c22b71d, 92b0d17a3), but upstream hardened them afterwards, in both cases because of the failures we hit: * 6251342e07d, "Disable parallel replicas for test" (same day 03927 landed): adds SET enable_parallel_replicas = 0. clickhouse-test randomizes parallel_replicas_local_plan, and with no local plan there is no local ReadFromMergeTree, so the ReadType lines the test greps for disappear. * c70a81c5558, "Fix flaky 03800 RLS+projection test under ParallelReplicas", plus 5c5e975c99d, 2954a15b8c8 and 8473072527d: disables parallel replicas on the EXPLAIN queries for the same reason, pins index_granularity because the EXPLAIN indexes section asserts an exact granule count, adds a baseline query without the row policy so the result demonstrably changes once the policy applies, and adds two assertions that do not depend on plan indentation - a count() > 0 check that the projection was read, and an extract() of the equals(tenant_id, ...) predicate showing the policy is applied as a prewhere filter on the projection. Both files are upstream/master verbatim except for SET explain_query_plan_default, which selects between the legacy and pretty EXPLAIN plan formats and does not exist on this branch - it was added upstream on 2026-05-20, and only the legacy format exists here. 03800's reference is regenerated against this branch: it differs from upstream only in "Ranges: 1" being indented two spaces deeper, because ReadFromMergeTree::describeIndexes still prints it with an extra indent level here. 03927's reference is byte-identical to upstream. Co-Authored-By: Claude Opus 5 (1M context) --- ...jection_row_policy_filter_column.reference | 5 ++ ...00_projection_row_policy_filter_column.sql | 53 ++++++++++++++++--- ...w_filter_in_read_in_order_optimization.sql | 3 +- 3 files changed, 53 insertions(+), 8 deletions(-) diff --git a/tests/queries/0_stateless/03800_projection_row_policy_filter_column.reference b/tests/queries/0_stateless/03800_projection_row_policy_filter_column.reference index 1fbd655702ca..5364507fb25b 100644 --- a/tests/queries/0_stateless/03800_projection_row_policy_filter_column.reference +++ b/tests/queries/0_stateless/03800_projection_row_policy_filter_column.reference @@ -1,3 +1,6 @@ +item_1 3 +item_2 1 +item_3 1 Expression (Project names) Sorting (Sorting for ORDER BY) Expression ((Before ORDER BY + Projection)) @@ -13,5 +16,7 @@ Expression (Project names) Granules: 1/1 Search Algorithm: binary search Ranges: 1 +1 +equals(tenant_id, \'tenant_A\'_String) item_1 2 item_2 1 diff --git a/tests/queries/0_stateless/03800_projection_row_policy_filter_column.sql b/tests/queries/0_stateless/03800_projection_row_policy_filter_column.sql index 5a763295fb5c..5f23ee1d9d98 100644 --- a/tests/queries/0_stateless/03800_projection_row_policy_filter_column.sql +++ b/tests/queries/0_stateless/03800_projection_row_policy_filter_column.sql @@ -6,16 +6,18 @@ CREATE TABLE test_rls_projection id UInt64, tenant_id String, data String, - + PROJECTION proj_by_data ( SELECT * ORDER BY tenant_id, data, id ) ) ENGINE = MergeTree() -ORDER BY (tenant_id, id); +ORDER BY (tenant_id, id) +-- Pin index_granularity: the EXPLAIN indexes section below asserts an exact granule count. +SETTINGS index_granularity = 8192; -INSERT INTO test_rls_projection VALUES +INSERT INTO test_rls_projection VALUES (1, 'tenant_A', 'item_1'), (2, 'tenant_A', 'item_2'), (3, 'tenant_A', 'item_1'), @@ -24,25 +26,62 @@ INSERT INTO test_rls_projection VALUES ALTER TABLE test_rls_projection MATERIALIZE PROJECTION proj_by_data; + +-- Baseline without any row policy: all rows are visible, so 'item_1' is counted 3 times +-- (twice for tenant_A and once for tenant_B) and 'item_3' is present. +SELECT data, count() as cnt +FROM test_rls_projection +GROUP BY data +ORDER BY data +SETTINGS force_optimize_projection = 1, optimize_use_projections = 1; + CREATE ROW POLICY rls_policy ON test_rls_projection FOR SELECT USING tenant_id = 'tenant_A' TO default; --- Verify projection is used +-- Verify the projection is used and the row policy filter is pushed down into the projection read. +-- Parallel replicas are disabled so that the EXPLAIN plan is deterministic: with parallel replicas +-- the read is wrapped into MergingAggregated/Union/ReadFromRemoteParallelReplicas and the plan text diverges. EXPLAIN indexes = 1 SELECT data, count() as cnt FROM test_rls_projection GROUP BY data ORDER BY data -SETTINGS force_optimize_projection = 1; +SETTINGS force_optimize_projection = 1, optimize_use_projections = 1, enable_analyzer = 1, enable_parallel_replicas = 0; + +-- Same plan via EXPLAIN actions = 1, showing that the row policy predicate on tenant_id is +-- applied as a PREWHERE filter on the projection read even though tenant_id is not selected. +SELECT count() > 0 AS projection_used +FROM ( + EXPLAIN actions = 1 + SELECT data, count() as cnt + FROM test_rls_projection + GROUP BY data + ORDER BY data +) +WHERE explain ILIKE '%ReadFromMergeTree (proj_by_data)%' +SETTINGS force_optimize_projection = 1, optimize_use_projections = 1, enable_analyzer = 1, enable_parallel_replicas = 0; + +SELECT DISTINCT extract(explain, 'equals\(tenant_id, ''[^'']+''_String\)') AS prewhere_filter +FROM ( + EXPLAIN actions = 1 + SELECT data, count() as cnt + FROM test_rls_projection + GROUP BY data + ORDER BY data +) +WHERE explain ILIKE '%-> equals(tenant_id,%' +SETTINGS force_optimize_projection = 1, optimize_use_projections = 1, enable_analyzer = 1, enable_parallel_replicas = 0; --- Query without tenant_id in SELECT - should work with RLS filter applied via projection +-- Query without tenant_id in SELECT - should work with the RLS filter applied via the projection. +-- The result must differ from the baseline: 'item_1' is now counted only twice (tenant_A rows) +-- and 'item_3' (a tenant_B row) is filtered out, proving the row policy filter is applied. SELECT data, count() as cnt FROM test_rls_projection GROUP BY data ORDER BY data -SETTINGS force_optimize_projection = 1; +SETTINGS force_optimize_projection = 1, optimize_use_projections = 1; DROP ROW POLICY rls_policy ON test_rls_projection; DROP TABLE test_rls_projection; diff --git a/tests/queries/0_stateless/03927_row_filter_in_read_in_order_optimization.sql b/tests/queries/0_stateless/03927_row_filter_in_read_in_order_optimization.sql index c47de047d82c..0d3a742ae4b6 100644 --- a/tests/queries/0_stateless/03927_row_filter_in_read_in_order_optimization.sql +++ b/tests/queries/0_stateless/03927_row_filter_in_read_in_order_optimization.sql @@ -1,3 +1,4 @@ +SET enable_parallel_replicas = 0; SET optimize_read_in_order = 1; DROP TABLE IF EXISTS t_row_policy_rio; @@ -36,4 +37,4 @@ LIMIT 100 ) WHERE explain LIKE '%ReadType%'; DROP ROW POLICY test_policy ON t_row_policy_rio; -DROP TABLE t_row_policy_rio; +DROP TABLE t_row_policy_rio; \ No newline at end of file From 21a65c30e01781d3a36787e6b78bde0c151c588e Mon Sep 17 00:00:00 2001 From: Mikhail Koviazin Date: Thu, 6 Aug 2026 14:01:26 +0200 Subject: [PATCH 6/6] Blacklist 03800 for parallel replicas 03800_projection_row_policy_filter_column's two data queries fail with PROJECTION_NOT_USED under the ParallelReplicas variant. Upstream's enable_parallel_replicas = 0 pins cover only its three EXPLAIN queries. projectionsCommon.cpp reports projection support for the initiator only when parallel_replicas_local_plan is set, so with it 0 optimizeUseNormalProjection skips projection reading on remote replicas and force_optimize_projection = 1 throws. That logic is identical upstream, and so is the randomization of the setting in clickhouse-test - but upstream additionally forces it back to 1 (its clickhouse-test has that override, ours does not). The test itself is therefore not at fault, so blacklist it rather than diverging the file from upstream, which the previous commit had just converged. Note this leaves the underlying gap in place: any other test relying on projections under parallel replicas will hit the same randomization. Co-Authored-By: Claude Opus 5 (1M context) --- tests/parallel_replicas_blacklist.txt | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/tests/parallel_replicas_blacklist.txt b/tests/parallel_replicas_blacklist.txt index 666cb12a3af9..6d03d7f718bf 100644 --- a/tests/parallel_replicas_blacklist.txt +++ b/tests/parallel_replicas_blacklist.txt @@ -528,3 +528,8 @@ 00443_optimize_final_vertical_merge 01660_join_or_any 01801_s3_cluster + + PROJECTION_NOT_USED with force_optimize_projection: projectionsCommon.cpp only + reports projection support for the initiator when parallel_replicas_local_plan is + set, which clickhouse-test randomizes here (upstream forces it back to 1) +03800_projection_row_policy_filter_column