This is an automated email from the ASF dual-hosted git repository.

tkalkirill pushed a commit to branch ignite-29003
in repository https://gitbox.apache.org/repos/asf/ignite.git

commit 94b4b85a472704af77170cc1d7f1355da2c6facd
Author: Kirill Tkalenko <[email protected]>
AuthorDate: Fri Aug 21 15:23:46 2026 +0300

    IGNITE-29003 Wip
---
 .../query/calcite/prepare/PlannerPhase.java        |   4 +-
 .../rule/TableFunctionScanScalarSubQueryRule.java  | 178 +++++++++++++++++++++
 .../UserDefinedFunctionsIntegrationTest.java       |  52 ++++++
 3 files changed, 233 insertions(+), 1 deletion(-)

diff --git 
a/modules/calcite/src/main/java/org/apache/ignite/internal/processors/query/calcite/prepare/PlannerPhase.java
 
b/modules/calcite/src/main/java/org/apache/ignite/internal/processors/query/calcite/prepare/PlannerPhase.java
index e7f790625ec..71753b173f2 100644
--- 
a/modules/calcite/src/main/java/org/apache/ignite/internal/processors/query/calcite/prepare/PlannerPhase.java
+++ 
b/modules/calcite/src/main/java/org/apache/ignite/internal/processors/query/calcite/prepare/PlannerPhase.java
@@ -66,6 +66,7 @@ import 
org.apache.ignite.internal.processors.query.calcite.rule.SetOpConverterRu
 import 
org.apache.ignite.internal.processors.query.calcite.rule.SortAggregateConverterRule;
 import 
org.apache.ignite.internal.processors.query.calcite.rule.SortConverterRule;
 import 
org.apache.ignite.internal.processors.query.calcite.rule.TableFunctionScanConverterRule;
+import 
org.apache.ignite.internal.processors.query.calcite.rule.TableFunctionScanScalarSubQueryRule;
 import 
org.apache.ignite.internal.processors.query.calcite.rule.TableModifyDistributedConverterRule;
 import 
org.apache.ignite.internal.processors.query.calcite.rule.TableModifySingleNodeConverterRule;
 import 
org.apache.ignite.internal.processors.query.calcite.rule.UncollectConverterRule;
@@ -94,7 +95,8 @@ public enum PlannerPhase {
                 RuleSets.ofList(
                     CoreRules.FILTER_SUB_QUERY_TO_CORRELATE,
                     CoreRules.PROJECT_SUB_QUERY_TO_CORRELATE,
-                    CoreRules.JOIN_SUB_QUERY_TO_CORRELATE
+                    CoreRules.JOIN_SUB_QUERY_TO_CORRELATE,
+                    TableFunctionScanScalarSubQueryRule.INSTANCE
                 )
             );
         }
diff --git 
a/modules/calcite/src/main/java/org/apache/ignite/internal/processors/query/calcite/rule/TableFunctionScanScalarSubQueryRule.java
 
b/modules/calcite/src/main/java/org/apache/ignite/internal/processors/query/calcite/rule/TableFunctionScanScalarSubQueryRule.java
new file mode 100644
index 00000000000..cee66c922d4
--- /dev/null
+++ 
b/modules/calcite/src/main/java/org/apache/ignite/internal/processors/query/calcite/rule/TableFunctionScanScalarSubQueryRule.java
@@ -0,0 +1,178 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *      http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.ignite.internal.processors.query.calcite.rule;
+
+import java.util.stream.Collectors;
+import java.util.stream.IntStream;
+import com.google.common.collect.ImmutableList;
+import org.apache.calcite.plan.RelOptRuleCall;
+import org.apache.calcite.plan.RelRule;
+import org.apache.calcite.rel.RelNode;
+import org.apache.calcite.rel.core.CorrelationId;
+import org.apache.calcite.rel.core.JoinRelType;
+import org.apache.calcite.rel.core.TableFunctionScan;
+import org.apache.calcite.rel.logical.LogicalCorrelate;
+import org.apache.calcite.rel.rules.TransformationRule;
+import org.apache.calcite.rex.RexCorrelVariable;
+import org.apache.calcite.rex.RexNode;
+import org.apache.calcite.rex.RexShuttle;
+import org.apache.calcite.rex.RexSubQuery;
+import org.apache.calcite.rex.RexVisitorImpl;
+import org.apache.calcite.sql.SqlKind;
+import org.apache.calcite.sql.fun.SqlStdOperatorTable;
+import org.apache.calcite.tools.RelBuilder;
+import org.apache.calcite.util.ImmutableBitSet;
+import org.apache.calcite.util.Util;
+import org.immutables.value.Value;
+
+import static java.util.Objects.requireNonNull;
+
+/**
+ * Rewrites scalar subqueries in table function arguments to correlates.
+ *
+ * <p>This is a temporary backport of
+ * <a 
href="https://issues.apache.org/jira/browse/CALCITE-7688";>CALCITE-7688</a>.
+ * Remove this rule and use Calcite's {@code 
CoreRules.TABLE_FUNCTION_SCAN_SCALAR_QUERY_TO_CORRELATE}
+ * after upgrading to Calcite 1.43.
+ */
[email protected]
+public class TableFunctionScanScalarSubQueryRule
+    extends RelRule<TableFunctionScanScalarSubQueryRule.Config> implements 
TransformationRule {
+    /** */
+    public static final TableFunctionScanScalarSubQueryRule INSTANCE = 
Config.DEFAULT.toRule();
+
+    /** */
+    private TableFunctionScanScalarSubQueryRule(Config cfg) {
+        super(cfg);
+    }
+
+    /** {@inheritDoc} */
+    @Override public void onMatch(RelOptRuleCall call) {
+        TableFunctionScan scan = call.rel(0);
+        RexSubQuery subQry = 
requireNonNull(findScalarSubQuery(scan.getCall()));
+        RelBuilder builder = call.builder();
+
+        builder.push(subQry.rel);
+        builder.aggregate(builder.groupKey(),
+            builder.aggregateCall(SqlStdOperatorTable.SINGLE_VALUE, 
builder.field(0)));
+
+        RelNode scalarVal = builder.build();
+        CorrelationId correlationId = scan.getCluster().createCorrel();
+        RexCorrelVariable correlationVar = 
(RexCorrelVariable)scan.getCluster().getRexBuilder()
+            .makeCorrel(scalarVal.getRowType(), correlationId);
+        RexNode target = 
scan.getCluster().getRexBuilder().makeFieldAccess(correlationVar, 0);
+        RexNode newCall = scan.getCall().accept(new 
ReplaceSubQueryShuttle(subQry, target));
+        TableFunctionScan newScan = (TableFunctionScan)scan.copy(
+            scan.getTraitSet(),
+            scan.getInputs(),
+            newCall,
+            scan.getElementType(),
+            scan.getRowType(),
+            scan.getColumnMappings()
+        ).withHints(scan.getHints());
+
+        RelNode correlate = LogicalCorrelate.create(
+            scalarVal,
+            newScan,
+            ImmutableList.of(),
+            correlationId,
+            ImmutableBitSet.of(0),
+            JoinRelType.INNER
+        );
+
+        builder.push(correlate);
+
+        int scalarFieldCnt = scalarVal.getRowType().getFieldCount();
+
+        builder.project(
+            IntStream.range(0, scan.getRowType().getFieldCount())
+                .mapToObj(i -> builder.field(scalarFieldCnt + i))
+                .collect(Collectors.toList()),
+            scan.getRowType().getFieldNames()
+        );
+
+        call.transformTo(builder.build());
+    }
+
+    /** Finds the first scalar subquery in the expression. */
+    private static RexSubQuery findScalarSubQuery(RexNode node) {
+        try {
+            node.accept(ScalarSubQueryFinder.INSTANCE);
+
+            return null;
+        }
+        catch (Util.FoundOne e) {
+            return (RexSubQuery)e.getNode();
+        }
+    }
+
+    /** Replaces one scalar subquery with a reference to the aggregate result. 
*/
+    private static class ReplaceSubQueryShuttle extends RexShuttle {
+        /** Subquery to replace. */
+        private final RexSubQuery subQry;
+
+        /** Replacement expression. */
+        private final RexNode replacement;
+
+        /** */
+        private ReplaceSubQueryShuttle(RexSubQuery subQry, RexNode 
replacement) {
+            this.subQry = subQry;
+            this.replacement = replacement;
+        }
+
+        /** {@inheritDoc} */
+        @Override public RexNode visitSubQuery(RexSubQuery subQry) {
+            return subQry.equals(this.subQry) ? replacement : subQry;
+        }
+    }
+
+    /** Finds scalar subqueries without matching other subquery kinds. */
+    private static class ScalarSubQueryFinder extends RexVisitorImpl<Void> {
+        /** */
+        private static final ScalarSubQueryFinder INSTANCE = new 
ScalarSubQueryFinder();
+
+        /** */
+        private ScalarSubQueryFinder() {
+            super(true);
+        }
+
+        /** {@inheritDoc} */
+        @Override public Void visitSubQuery(RexSubQuery subQry) {
+            if (subQry.getKind() == SqlKind.SCALAR_QUERY)
+                throw new Util.FoundOne(subQry);
+
+            return super.visitSubQuery(subQry);
+        }
+    }
+
+    /** Rule configuration. */
+    @Value.Immutable
+    public interface Config extends RelRule.Config {
+        /** */
+        Config DEFAULT = 
ImmutableTableFunctionScanScalarSubQueryRule.Config.of()
+            .withDescription("TableFunctionScanScalarSubQueryRule")
+            .withOperandSupplier(b -> b.operand(TableFunctionScan.class)
+                .predicate(scan -> findScalarSubQuery(scan.getCall()) != null)
+                .anyInputs());
+
+        /** {@inheritDoc} */
+        @Override default TableFunctionScanScalarSubQueryRule toRule() {
+            return new TableFunctionScanScalarSubQueryRule(this);
+        }
+    }
+}
diff --git 
a/modules/calcite/src/test/java/org/apache/ignite/internal/processors/query/calcite/integration/UserDefinedFunctionsIntegrationTest.java
 
b/modules/calcite/src/test/java/org/apache/ignite/internal/processors/query/calcite/integration/UserDefinedFunctionsIntegrationTest.java
index 278b3bac52f..0217b6359bb 100644
--- 
a/modules/calcite/src/test/java/org/apache/ignite/internal/processors/query/calcite/integration/UserDefinedFunctionsIntegrationTest.java
+++ 
b/modules/calcite/src/test/java/org/apache/ignite/internal/processors/query/calcite/integration/UserDefinedFunctionsIntegrationTest.java
@@ -344,6 +344,51 @@ public class UserDefinedFunctionsIntegrationTest extends 
AbstractBasicIntegratio
             .check();
     }
 
+    /** */
+    @Test
+    public void testScalarSubqueriesInTableFunctionArguments() throws 
Exception {
+        IgniteCache<Integer, Employer> emp = client.getOrCreateCache(new 
CacheConfiguration<Integer, Employer>("emp")
+            .setSqlSchema("PUBLIC")
+            .setSqlFunctionClasses(TableFunctionsLibrary.class)
+            .setQueryEntities(F.asList(new QueryEntity(Integer.class, 
Employer.class).setTableName("emp")))
+        );
+
+        emp.put(1, new Employer("Igor1", 1d));
+        emp.put(2, new Employer("Roman1", 2d));
+
+        awaitPartitionMapExchange();
+
+        assertQuery("SELECT * FROM TABLE(scalarQueryArguments((SELECT 10), 
20))")
+            .returns(10, 20)
+            .check();
+
+        assertQuery("SELECT * FROM TABLE(scalarQueryArguments((SELECT 10), 
(SELECT 20)))")
+            .returns(10, 20)
+            .check();
+
+        assertQuery("SELECT * FROM TABLE(scalarQueryArguments((SELECT 4) + 
(SELECT 6), 20))")
+            .returns(10, 20)
+            .check();
+
+        assertQuery("SELECT * FROM TABLE(scalarQueryArguments(" +
+            "(SELECT _KEY FROM emp WHERE _KEY < 0), 20))")
+            .returns(null, 20)
+            .check();
+
+        assertQuery("SELECT e._KEY, (SELECT f.SCALAR_VALUE FROM TABLE(" +
+            "scalarQueryArguments((SELECT e._KEY + 1), e._KEY)) f) " +
+            "FROM emp e ORDER BY e._KEY")
+            .returns(1, 2)
+            .returns(2, 3)
+            .check();
+
+        assertThrows(
+            "SELECT * FROM TABLE(scalarQueryArguments((SELECT _KEY FROM emp), 
20))",
+            IllegalArgumentException.class,
+            "Subquery returned more than 1 value."
+        );
+    }
+
     /** */
     @Test
     public void testIncorrectTableFunctions() throws Exception {
@@ -430,6 +475,13 @@ public class UserDefinedFunctionsIntegrationTest extends 
AbstractBasicIntegratio
             );
         }
 
+        /** Returns a single row containing the function arguments. */
+        @QuerySqlTableFunction(columnTypes = {Integer.class, Integer.class},
+            columnNames = {"SCALAR_VALUE", "LITERAL_VALUE"})
+        public static Iterable<Collection<?>> scalarQueryArguments(Integer 
scalarVal, int literalVal) {
+            return List.of(Arrays.asList(scalarVal, literalVal));
+        }
+
         /** Overrides. */
         @QuerySqlTableFunction(columnTypes = {int.class, int.class, 
int.class}, columnNames = {"COL_1", "COL_2", "COL_3"})
         public static Collection<Collection<?>> collectionRow(int x, int y, 
int z) {

Reply via email to