Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
22 commits
Select commit Hold shift + click to select a range
c66ee3b
tests: stage all Altinity#1844 alias-marker tests for 26.6 triage
mkmkme Sep 11, 2026
cbebdb3
Fix MULTIPLE_EXPRESSIONS_FOR_ALIAS from __aliasMarker on distributed …
mkmkme Sep 15, 2026
1796c7f
Fix ALIAS columns read through Merge over Distributed
mkmkme Sep 15, 2026
714a78b
Materialize __aliasMarker ids after the shard tree is renumbered
mkmkme Sep 15, 2026
629e2c7
tests: assert results where the expected errors no longer happen
mkmkme Sep 15, 2026
9c167de
Finalize __aliasMarker inside lambdas; drop a no-op converting helper
mkmkme Sep 19, 2026
a2ea20f
Keep the __aliasMarker id out of the shipped GLOBAL JOIN subquery
mkmkme Sep 19, 2026
63e9fed
Revert two StorageMerge changes that do nothing
mkmkme Sep 19, 2026
9cdf271
docs: describe __aliasMarker as it actually behaves
mkmkme Sep 19, 2026
7f91cb2
Merge branch 'antalya-26.6' into mkmkme/26.6/alias-marker
mkmkme Sep 19, 2026
6a6848d
Merge branch 'antalya-26.6' into mkmkme/26.6/alias-marker
mkmkme Sep 22, 2026
232982f
Merge branch 'antalya-26.6' into mkmkme/26.6/alias-marker
mkmkme Sep 26, 2026
eb13e8b
tests: skip `05060` in the fast test
mkmkme Sep 26, 2026
cceecd5
Fix provenance of injected `__aliasMarker` in distributed queries
mkmkme Sep 26, 2026
1d7a7de
tests: assert parallel-replica execution in `03931`
mkmkme Sep 26, 2026
bcea758
Fix quoted `Merge` alias identifiers and row-policy inputs
mkmkme Sep 26, 2026
3093350
Preserve bare row-policy inputs in `StorageMerge`
mkmkme Sep 26, 2026
685d5d3
Consolidate alias-marker regression coverage
mkmkme Sep 29, 2026
ede8dad
Keep pending-marker local test off parallel replicas
mkmkme Sep 29, 2026
ea3ed2c
Merge branch 'antalya-26.6' into mkmkme/26.6/alias-marker
mkmkme Sep 29, 2026
9a0a8fb
Pin execution modes for parallel-replica alias test
mkmkme Sep 29, 2026
9a2a424
Preserve `LowCardinality` payload types in `__aliasMarker`
mkmkme Sep 30, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
77 changes: 77 additions & 0 deletions src/Analyzer/Utils.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@

#include <Functions/FunctionHelpers.h>
#include <Functions/FunctionFactory.h>
#include <Functions/identity.h>

#include <Storages/IStorage.h>

Expand All @@ -55,6 +56,7 @@

#include <functional>
#include <ranges>
#include <vector>

namespace DB
{
Expand Down Expand Up @@ -1017,6 +1019,81 @@ void resolveAggregateFunctionNodeByName(FunctionNode & function_node, const Stri
function_node.resolveAsAggregateFunction(std::move(aggregate_function));
}

namespace
{

class FinalizeAliasMarkersVisitor : public InDepthQueryTreeVisitor<FinalizeAliasMarkersVisitor>
{
public:
explicit FinalizeAliasMarkersVisitor(ContextPtr context_) : context(std::move(context_)) {}

/// Visit children first, so a nested marker chain is materialized from the inside out.
bool shouldTraverseTopToBottom() const { return false; }

void visitImpl(QueryTreeNodePtr & node)
{
auto * function_node = node->as<FunctionNode>();
if (!function_node || function_node->getFunctionName() != "__aliasMarker")
return;

auto & arguments = function_node->getArguments().getNodes();
if (arguments.size() != 3)
return;

if (!arguments[0] || !arguments[1] || !arguments[2])
throw Exception(ErrorCodes::BAD_ARGUMENTS, "Invalid internal pending __aliasMarker arguments");

const auto * token = arguments[2]->as<ConstantNode>();
if (!token || !isString(token->getResultType()) || token->getValue().safeGet<String>() != AliasMarkerName::pending_token)
throw Exception(ErrorCodes::BAD_ARGUMENTS, "Invalid internal pending __aliasMarker token");

const auto * column_node = arguments[1]->as<ColumnNode>();
if (!column_node)
throw Exception(ErrorCodes::BAD_ARGUMENTS, "Internal pending __aliasMarker requires a column id");

const auto & column_source = column_node->getColumnSourceOrNull();
if (!column_source || column_source->getNodeType() == QueryTreeNodeType::LAMBDA || !column_source->hasAlias())
throw Exception(ErrorCodes::BAD_ARGUMENTS, "Internal pending __aliasMarker id has no table source alias");

auto alias_id = column_source->getAlias() + "." + column_node->getColumnName();
arguments[1] = std::make_shared<ConstantNode>(std::move(alias_id), std::make_shared<DataTypeString>());
arguments.pop_back();
resolveOrdinaryFunctionNodeByName(*function_node, "__aliasMarker", context);
}

private:
ContextPtr context;
};

}

void finalizeAliasMarkersForDistributedSerialization(QueryTreeNodePtr & node, const ContextPtr & context)
{
FinalizeAliasMarkersVisitor visitor(context);
visitor.visit(node);
}

void assertNoPendingAliasMarkersForDistributedSerialization(const QueryTreeNodePtr & node)
{
if (!node)
return;

std::vector<const IQueryTreeNode *> nodes_to_visit{node.get()};
while (!nodes_to_visit.empty())
{
const auto * current = nodes_to_visit.back();
nodes_to_visit.pop_back();

if (const auto * function_node = current->as<FunctionNode>();
function_node && function_node->getFunctionName() == "__aliasMarker" && function_node->getArguments().getNodes().size() == 3)
throw Exception(ErrorCodes::BAD_ARGUMENTS, "Internal pending __aliasMarker cannot be shipped before finalization");

for (const auto & child : current->getChildren())
if (child)
nodes_to_visit.push_back(child.get());
}
}

std::pair<QueryTreeNodePtr, bool> getExpressionSource(const QueryTreeNodePtr & node)
{
if (const auto * column = node->as<ColumnNode>())
Expand Down
16 changes: 16 additions & 0 deletions src/Analyzer/Utils.h
Original file line number Diff line number Diff line change
Expand Up @@ -157,6 +157,22 @@ void resolveOrdinaryFunctionNodeByName(FunctionNode & function_node, const Strin
/// Arguments and parameters are taken from the node.
void resolveAggregateFunctionNodeByName(FunctionNode & function_node, const String & function_name);

/** Materialize the id of each injected, three-argument pending `__aliasMarker` in the tree. Replace its `ColumnNode`
* with the `String` constant the shard will read, and remove its pending-form token. Leave two-argument calls alone.
*
* Call this immediately before the tree is rendered as SQL for a shard, and no earlier. The id is the column's
* analyzer identifier, and `createUniqueAliasesIfNecessary` -- which runs inside `buildQueryTreeForShard` -- is what
* settles the `__tableN` part of it. An id frozen before that point names a table alias that no longer exists by the
* time the shard header is built.
*
* A marker already materialized on an earlier hop has two arguments. A hand-written two-argument marker is left
* alone too, even when its id is a `ColumnNode`.
*/
void finalizeAliasMarkersForDistributedSerialization(QueryTreeNodePtr & node, const ContextPtr & context);

/// Reject a pending `__aliasMarker` before converting a query tree to SQL or building a plan to ship.
void assertNoPendingAliasMarkersForDistributedSerialization(const QueryTreeNodePtr & node);

/// Returns single source of expression node.
/// First element of pair is source node, can be nullptr if there are no sources or multiple sources.
/// Second element of pair is true if there is at most one source, false if there are multiple sources.
Expand Down
8 changes: 5 additions & 3 deletions src/Functions/identity.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -38,12 +38,14 @@ REGISTER_FUNCTION(AliasMarker)
{
factory.registerFunction<FunctionAliasMarker>(FunctionDocumentation{
.description = R"(
Internal function that marks ALIAS column expressions for the analyzer. Not intended for direct use.
Internal function. Returns its first argument unchanged. The second argument records which ALIAS column the
expression was inlined from, so that a shard names the result the way the initiator expects. Not intended for
direct use: an explicitly supplied string id controls the planner's action name.
)",
.syntax = {"__aliasMarker(expr, alias_name)"},
.syntax = {"__aliasMarker(expr, alias_id)"},
.arguments = {
{"expr", "Expression to mark.", {"Any"}},
{"alias_name", "Alias name attached to the expression.", {"String"}},
{"alias_id", "Identity of the ALIAS column the expression came from.", {"Any"}},
},
.returned_value = {"Returns expr unchanged.", {"Any"}},
.introduced_in = {25, 8},
Expand Down
56 changes: 48 additions & 8 deletions src/Functions/identity.h
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
#pragma once
#include <Columns/ColumnConst.h>
#include <DataTypes/IDataType.h>
#include <Functions/IFunction.h>
#include <Interpreters/Context_fwd.h>
Expand Down Expand Up @@ -106,29 +107,68 @@ class FunctionActionName final : public FunctionIdentityBase
struct AliasMarkerName
{
static constexpr auto name = "__aliasMarker";
static constexpr auto pending_token = "__aliasMarker_pending_v1";
};

/** `__aliasMarker(expr, id)` is an internal pass-through identity. It returns `expr` untouched; `id` exists only to
* give the expression a stable name in the planner's `ActionsDAG`.
*
* It is injected when an `ALIAS` column is inlined into its defining expression for transport to a shard. The
* initiator still sees the un-inlined column, so without the marker the shard would name the output after the
* expression (`multiply(__table1.value, 2)`) while the initiator expects the column (`__table1.computed`), and the
* two headers could not be matched by name.
*
* The marker carries the low-level column identity, not the user's SQL alias. A SQL alias cannot do this job: it
* participates in user-visible query semantics, it can collide with names the user chose, and in the
* mergeable-state path the projection step that would normally apply it is skipped.
*
* This is also why the marker is not `__actionName`. `__actionName` survives into the `ActionsDAG` as a function node
* with a forced result name; `__aliasMarker` is consumed into an alias on top of its child, which is what keeps the
* expression behaving like a distinct logical column.
*
* An injected marker has a third, constant token argument until it is shipped. The second argument is a `ColumnNode`
* so analyzer passes, in particular `createUniqueAliasesIfNecessary`, can assign the final `__tableN` alias.
* `finalizeAliasMarkersForDistributedSerialization` replaces it with a `String` id and removes the token before shipping.
* The token distinguishes this pending form from a hand-written two-argument call; it is not authentication.
*
* A hand-written two-argument `__aliasMarker(x, x)` is not finalized. An explicitly supplied `String` id still
* controls action naming, so direct use of this internal function is not guaranteed to be harmless.
*
* This is a bridge for as long as distributed transport still goes through SQL text. Once query plan serialization
* replaces that boundary, the marker should become unnecessary.
*/
class FunctionAliasMarker : public IFunction
{
public:
static constexpr auto name = AliasMarkerName::name;
static FunctionPtr create(ContextPtr) { return std::make_shared<FunctionAliasMarker>(); }

String getName() const override { return name; }
size_t getNumberOfArguments() const override { return 2; }
ColumnNumbers getArgumentsThatAreAlwaysConstant() const override { return {1}; }
size_t getNumberOfArguments() const override { return 0; }
bool isVariadic() const override { return true; }
/// Index 2 exists only on the pending form; the executable and DAG builder ignore absent constant arguments.
ColumnNumbers getArgumentsThatAreAlwaysConstant() const override { return {2}; }
bool isSuitableForConstantFolding() const override { return false; }
bool isSuitableForShortCircuitArgumentsExecution(const DataTypesWithConstInfo & /*arguments*/) const override { return false; }
/// Validate the token even if an argument is `NULL` or has type `Nothing`.
bool useDefaultImplementationForNulls() const override { return false; }
bool useDefaultImplementationForNothing() const override { return false; }
/// Preserve the payload type: the planner replaces the marker with its child without conversion.
bool useDefaultImplementationForLowCardinalityColumns() const override { return false; }

DataTypePtr getReturnTypeImpl(const DataTypes & arguments) const override
DataTypePtr getReturnTypeImpl(const ColumnsWithTypeAndName & arguments) const override
{
if (arguments.size() != 2)
throw Exception(ErrorCodes::BAD_ARGUMENTS, "Function __aliasMarker expects 2 arguments");
if (arguments.size() != 2 && arguments.size() != 3)
throw Exception(ErrorCodes::BAD_ARGUMENTS, "Function __aliasMarker expects 2 arguments, or 3 for its internal pending form");

if (!WhichDataType(arguments[1]).isString())
throw Exception(ErrorCodes::BAD_ARGUMENTS, "Function __aliasMarker is internal and should not be used directly");
if (arguments.size() == 3)
{
const auto * token = arguments[2].column ? checkAndGetColumn<ColumnConst>(arguments[2].column.get()) : nullptr;
if (!WhichDataType(arguments[2].type).isString() || !token || token->getValue<String>() != AliasMarkerName::pending_token)
throw Exception(ErrorCodes::BAD_ARGUMENTS, "Invalid internal pending __aliasMarker token");
}

return arguments.front();
return arguments.front().type;
}

ColumnPtr executeImpl(const ColumnsWithTypeAndName & arguments, const DataTypePtr &, size_t /*input_rows_count*/) const override
Expand Down
3 changes: 3 additions & 0 deletions src/Interpreters/ClusterProxy/executeQuery.cpp
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
#include <Analyzer/QueryNode.h>
#include <Analyzer/UnionNode.h>
#include <Analyzer/Utils.h>
#include <Columns/ColumnConst.h>
#include <Common/ProfileEvents.h>
#include <Core/QueryProcessingStage.h>
Expand Down Expand Up @@ -1169,6 +1170,8 @@ std::optional<QueryPipeline> executeInsertSelectWithParallelReplicas(
{
InterpreterSelectQueryAnalyzer analyzer(query_ast.select, new_context, {});
const auto & query_tree = analyzer.getQueryTree();
/// This parallel `INSERT SELECT` path renders directly, bypassing `queryNodeToDistributedSelectQuery`.
assertNoPendingAliasMarkersForDistributedSerialization(query_tree);
auto select_ast = query_tree->toAST();

auto new_query_ast = query_ast.clone();
Expand Down
45 changes: 28 additions & 17 deletions src/Planner/PlannerActionsVisitor.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -95,6 +95,16 @@ String calculateActionNodeNameWithCastIfNeeded(const ConstantNode & constant_nod
return buffer.str();
}

/// Return a two-argument `__aliasMarker`'s string id, or an empty string when it has none.
/// A three-argument pending marker always resolves to its payload in an intermediate plan.
String tryExtractAliasMarkerId(const QueryTreeNodePtr & id_argument)
{
if (const auto * id_node = id_argument->as<ConstantNode>(); id_node && isString(id_node->getResultType()))
return id_node->getValue().safeGet<String>();

return {};
}

class ActionNodeNameHelper
{
public:
Expand Down Expand Up @@ -198,18 +208,19 @@ class ActionNodeNameHelper
const auto & function_node = node->as<FunctionNode &>();
if (function_node.getFunctionName() == "__aliasMarker")
{
/// Perform sanity check, because user may call this function with unexpected arguments
const auto & function_argument_nodes = function_node.getArguments().getNodes();
if (function_argument_nodes.size() != 2 && function_argument_nodes.size() != 3)
throw Exception(ErrorCodes::BAD_ARGUMENTS, "Function __aliasMarker expects 2 or 3 arguments");

if (function_argument_nodes.size() == 2)
{
if (const auto * second_argument = function_argument_nodes.at(1)->as<ConstantNode>())
{
if (isString(second_argument->getResultType()))
result = second_argument->getValue().safeGet<String>();
}
}
result = tryExtractAliasMarkerId(function_argument_nodes.at(1));

/// Empty node name is not allowed and leads to logical errors
/// A pending marker, or a hand-written one without a string id, is named after its payload.
/// The pending form can reach an intermediate planner, but must be finalized before shipping.
if (result.empty())
result = calculateActionNodeName(function_argument_nodes.at(0));

/// An empty node name is not allowed and leads to logical errors.
if (result.empty())
throw Exception(ErrorCodes::BAD_ARGUMENTS, "Function __aliasMarker is internal and should not be used directly");
break;
Expand Down Expand Up @@ -1256,18 +1267,18 @@ PlannerActionsVisitorImpl::NodeNameAndNodeMinLevel PlannerActionsVisitorImpl::vi
if (function_node.getFunctionName() == "__aliasMarker")
{
const auto & function_arguments = function_node.getArguments().getNodes();
if (function_arguments.size() != 2)
throw Exception(ErrorCodes::BAD_ARGUMENTS, "Function __aliasMarker expects 2 arguments");
if (function_arguments.size() != 2 && function_arguments.size() != 3)
throw Exception(ErrorCodes::BAD_ARGUMENTS, "Function __aliasMarker expects 2 or 3 arguments");

const auto * alias_id_node = function_arguments.at(1)->as<ConstantNode>();
if (!alias_id_node || !isString(alias_id_node->getResultType()))
throw Exception(ErrorCodes::BAD_ARGUMENTS, "Function __aliasMarker is internal and should not be used directly");
auto [child_name, levels] = visitImpl(function_arguments.at(0));

const auto & alias_id = alias_id_node->getValue().safeGet<String>();
/// An intermediate pending marker and a hand-written marker without a string id resolve to their payload.
String alias_id;
if (function_arguments.size() == 2)
alias_id = tryExtractAliasMarkerId(function_arguments.at(1));
if (alias_id.empty())
throw Exception(ErrorCodes::BAD_ARGUMENTS, "Function __aliasMarker is internal and should not be used directly");
alias_id = child_name;

auto [child_name, levels] = visitImpl(function_arguments.at(0));
if (alias_id == child_name)
return {child_name, levels};

Expand Down
3 changes: 3 additions & 0 deletions src/Planner/Utils.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -317,6 +317,9 @@ void normalizeAliasMarkersInQueryTree(QueryTreeNodePtr & node)

ASTPtr queryNodeToDistributedSelectQuery(const QueryTreeNodePtr & query_node)
{
/// Check before normalization: it can discard a nested marker and hide an unfinalized pending id.
assertNoPendingAliasMarkersForDistributedSerialization(query_node);

/// Remove CTEs information from distributed queries.
/// Now, if cte_name is set for subquery node, AST -> String serialization will only print cte name.
/// But CTE is defined only for top-level query part, so may not be sent.
Expand Down
4 changes: 4 additions & 0 deletions src/Processors/QueryPlan/ParallelReplicasLocalPlan.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
#include <Common/FailPoint.h>
#include <Analyzer/QueryNode.h>
#include <Analyzer/UnionNode.h>
#include <Analyzer/Utils.h>
#include <Core/Settings.h>
#include <Interpreters/Context.h>
#include <Interpreters/IJoin.h>
Expand Down Expand Up @@ -116,6 +117,9 @@ std::shared_ptr<const QueryPlan> createRemotePlanForParallelReplicas(
{
checkStackSize();

/// A pending marker planned here would resolve to its payload and disappear before the plan was shipped.
assertNoPendingAliasMarkersForDistributedSerialization(query_tree);

auto new_context = Context::createCopy(context);

auto select_query_options = SelectQueryOptions(processed_stage);
Expand Down
5 changes: 5 additions & 0 deletions src/Storages/IStorageCluster.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -437,6 +437,11 @@ void IStorageCluster::read(

auto cluster_name_from_settings = getClusterName(context);
const auto & settings = context->getSettingsRef();
/// The ALLOW path can forward the original AST without calling `queryNodeToDistributedSelectQuery`.
if (settings[Setting::allow_experimental_analyzer]
&& (!cluster_name_from_settings.empty() || settings[Setting::object_storage_remote_initiator]))
assertNoPendingAliasMarkersForDistributedSerialization(query_info.query_tree);

ASTPtr query_to_send = query_info.query;

if (cluster_name_from_settings.empty())
Expand Down
Loading
Loading