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

jackietien pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/master by this push:
     new cdab0f68666 [Table Model Subquery] Support correlated scalar subquery
cdab0f68666 is described below

commit cdab0f68666b111bd1dc644999239e7228137f6a
Author: Liao Lanyu <[email protected]>
AuthorDate: Wed May 14 14:24:54 2025 +0800

    [Table Model Subquery] Support correlated scalar subquery
---
 .../IoTDBCorrelatedExistsSubqueryIT.java           |   2 +-
 ...T.java => IoTDBCorrelatedScalarSubqueryIT.java} | 283 ++++++++++-----------
 .../operator/process/FilterAndProjectOperator.java |   3 +
 .../relational/ColumnTransformerBuilder.java       |   7 +
 .../relational/metadata/TableMetadataImpl.java     |   3 +
 .../rule/TransformCorrelatedScalarSubquery.java    | 194 ++++++++++++++
 .../rule/TransformExistsApplyToCorrelatedJoin.java |   1 +
 .../optimizations/LogicalOptimizeFactory.java      |   5 +-
 .../optimizations/PlanNodeDecorrelator.java        |   2 +-
 .../dag/column/FailFunctionColumnTransformer.java  |  62 +++++
 .../unary/ArithmeticNegationColumnTransformer.java |   3 +-
 11 files changed, 411 insertions(+), 154 deletions(-)

diff --git 
a/integration-test/src/test/java/org/apache/iotdb/relational/it/query/recent/subquery/correlated/IoTDBCorrelatedExistsSubqueryIT.java
 
b/integration-test/src/test/java/org/apache/iotdb/relational/it/query/recent/subquery/correlated/IoTDBCorrelatedExistsSubqueryIT.java
index fd80453c4df..495b93d3345 100644
--- 
a/integration-test/src/test/java/org/apache/iotdb/relational/it/query/recent/subquery/correlated/IoTDBCorrelatedExistsSubqueryIT.java
+++ 
b/integration-test/src/test/java/org/apache/iotdb/relational/it/query/recent/subquery/correlated/IoTDBCorrelatedExistsSubqueryIT.java
@@ -358,7 +358,7 @@ public class IoTDBCorrelatedExistsSubqueryIT {
   }
 
   @Test
-  public void testUnCorrelatedExistsSubqueryInSelectClause() {
+  public void testCorrelatedExistsSubqueryInSelectClause() {
     String sql;
     String[] expectedHeader;
     String[] retArray;
diff --git 
a/integration-test/src/test/java/org/apache/iotdb/relational/it/query/recent/subquery/correlated/IoTDBCorrelatedExistsSubqueryIT.java
 
b/integration-test/src/test/java/org/apache/iotdb/relational/it/query/recent/subquery/correlated/IoTDBCorrelatedScalarSubqueryIT.java
similarity index 52%
copy from 
integration-test/src/test/java/org/apache/iotdb/relational/it/query/recent/subquery/correlated/IoTDBCorrelatedExistsSubqueryIT.java
copy to 
integration-test/src/test/java/org/apache/iotdb/relational/it/query/recent/subquery/correlated/IoTDBCorrelatedScalarSubqueryIT.java
index fd80453c4df..ef0921d9f38 100644
--- 
a/integration-test/src/test/java/org/apache/iotdb/relational/it/query/recent/subquery/correlated/IoTDBCorrelatedExistsSubqueryIT.java
+++ 
b/integration-test/src/test/java/org/apache/iotdb/relational/it/query/recent/subquery/correlated/IoTDBCorrelatedScalarSubqueryIT.java
@@ -23,7 +23,6 @@ import org.apache.iotdb.it.env.EnvFactory;
 import org.apache.iotdb.it.framework.IoTDBTestRunner;
 import org.apache.iotdb.itbase.category.TableClusterIT;
 import org.apache.iotdb.itbase.category.TableLocalStandaloneIT;
-import org.apache.iotdb.rpc.TSStatusCode;
 
 import org.junit.AfterClass;
 import org.junit.BeforeClass;
@@ -34,14 +33,14 @@ import org.junit.runner.RunWith;
 import static org.apache.iotdb.db.it.utils.TestUtils.prepareTableData;
 import static org.apache.iotdb.db.it.utils.TestUtils.tableAssertTestFail;
 import static org.apache.iotdb.db.it.utils.TestUtils.tableResultSetEqualTest;
-import static 
org.apache.iotdb.db.queryengine.plan.relational.planner.optimizations.JoinUtils.ONLY_SUPPORT_EQUI_JOIN;
 import static 
org.apache.iotdb.relational.it.query.recent.subquery.SubqueryDataSetUtils.CREATE_SQLS;
 import static 
org.apache.iotdb.relational.it.query.recent.subquery.SubqueryDataSetUtils.DATABASE_NAME;
 import static 
org.apache.iotdb.relational.it.query.recent.subquery.SubqueryDataSetUtils.NUMERIC_MEASUREMENTS;
 
 @RunWith(IoTDBTestRunner.class)
 @Category({TableLocalStandaloneIT.class, TableClusterIT.class})
-public class IoTDBCorrelatedExistsSubqueryIT {
+public class IoTDBCorrelatedScalarSubqueryIT {
+
   @BeforeClass
   public static void setUp() throws Exception {
     EnvFactory.getEnv().getConfig().getCommonConfig().setSortBufferSize(128 * 
1024);
@@ -56,78 +55,62 @@ public class IoTDBCorrelatedExistsSubqueryIT {
   }
 
   @Test
-  public void testCorrelatedExistsSubqueryInWhereClause() {
+  public void testCorrelatedScalarSubqueryInWhereClause() {
     String sql;
     String[] expectedHeader;
     String[] retArray;
 
-    // Test case: exists and other filter
-    sql =
-        "SELECT cast(%s AS INT32) as %s FROM table1 t1 WHERE device_id = 'd01' 
and exists(SELECT (%s) from table3 t3 WHERE device_id = 'd01' and t1.%s = 
t3.%s)";
-    retArray = new String[] {"30,", "40,"};
-    for (String measurement : NUMERIC_MEASUREMENTS) {
-      expectedHeader = new String[] {measurement};
-      tableResultSetEqualTest(
-          String.format(sql, measurement, measurement, measurement, 
measurement, measurement),
-          expectedHeader,
-          retArray,
-          DATABASE_NAME);
-    }
-
-    // Test case: only exists
-    sql = "SELECT s1 FROM table3 t3 WHERE exists(SELECT s1 from table1 t1 
WHERE t1.s1 = t3.s1)";
-    retArray = new String[] {"30,", "30,", "40,"};
-    expectedHeader = new String[] {"s1"};
-    tableResultSetEqualTest(sql, expectedHeader, retArray, DATABASE_NAME);
-
-    // Test case: exists with distinct
+    // Test case: Aggregation with correlated filter in scalar subquery
     sql =
-        "SELECT distinct cast(%s AS INT32) as %s FROM table1 t1 WHERE 
exists(SELECT (%s) from table3 t3 WHERE device_id = 'd01' and t1.%s = t3.%s)";
+        "SELECT cast(%s AS INT32) as %s FROM table1 t1 WHERE device_id = 'd01' 
and %s >= (SELECT max(%s) from table3 t3 where t1.%s = t3.%s)";
     retArray = new String[] {"30,", "40,"};
     for (String measurement : NUMERIC_MEASUREMENTS) {
       expectedHeader = new String[] {measurement};
       tableResultSetEqualTest(
-          String.format(sql, measurement, measurement, measurement, 
measurement, measurement),
-          expectedHeader,
-          retArray,
-          DATABASE_NAME);
-    }
-
-    // Test case: aggregation in exists
-    sql =
-        "SELECT cast(min(%s) as INT32) as %s FROM table1 t1 WHERE 
exists(SELECT avg(%s) from table3 t3 WHERE device_id = 'd01' and t1.%s = 
t3.%s)";
-    retArray = new String[] {"30,"};
-    for (String measurement : NUMERIC_MEASUREMENTS) {
-      expectedHeader = new String[] {measurement};
-      tableResultSetEqualTest(
-          String.format(sql, measurement, measurement, measurement, 
measurement, measurement),
+          String.format(
+              sql, measurement, measurement, measurement, measurement, 
measurement, measurement),
           expectedHeader,
           retArray,
           DATABASE_NAME);
     }
 
-    // Test case: aggregation with group by in exists(subquery returns empty 
result with having
-    // clause)
+    // Test case: Non-Aggregation with correlated filter in scalar subquery
     sql =
-        "SELECT distinct cast(%s AS INT32) as %s FROM table1 t1 WHERE 
exists(SELECT count(*) from table3 t3 where t1.%s = t3.%s group by device_id 
having count(*) > 5)";
-    retArray = new String[] {};
+        "SELECT cast(%s AS INT32) as %s FROM table1 t1 WHERE device_id = 'd01' 
and %s >= (SELECT distinct %s from table3 t3 where t1.%s = t3.%s and %s > 30)";
+    retArray = new String[] {"40,"};
     for (String measurement : NUMERIC_MEASUREMENTS) {
       expectedHeader = new String[] {measurement};
       tableResultSetEqualTest(
-          String.format(sql, measurement, measurement, measurement, 
measurement, measurement),
+          String.format(
+              sql,
+              measurement,
+              measurement,
+              measurement,
+              measurement,
+              measurement,
+              measurement,
+              measurement),
           expectedHeader,
           retArray,
           DATABASE_NAME);
     }
 
-    // Test case: limit 1 in exists
+    // Test case: limit 1 in scalar subquery
     sql =
-        "SELECT cast(%s AS INT32) as %s FROM table1 t1 WHERE device_id = 'd01' 
and exists(SELECT (%s) from table3 t3 WHERE device_id = 'd01' and t1.%s = t3.%s 
limit 1)";
-    retArray = new String[] {"30,", "40,"};
+        "SELECT cast(%s AS INT32) as %s FROM table1 t1 WHERE device_id = 'd01' 
and %s >= (SELECT  %s from table3 t3 where t1.%s = t3.%s and %s > 30 limit 1)";
+    retArray = new String[] {"40,"};
     for (String measurement : NUMERIC_MEASUREMENTS) {
       expectedHeader = new String[] {measurement};
       tableResultSetEqualTest(
-          String.format(sql, measurement, measurement, measurement, 
measurement, measurement),
+          String.format(
+              sql,
+              measurement,
+              measurement,
+              measurement,
+              measurement,
+              measurement,
+              measurement,
+              measurement),
           expectedHeader,
           retArray,
           DATABASE_NAME);
@@ -135,56 +118,70 @@ public class IoTDBCorrelatedExistsSubqueryIT {
   }
 
   @Test
-  public void testNestedExistsSubquery() {
+  public void testNestedCorrelatedScalarSubquery() {
     String sql;
     String[] expectedHeader;
     String[] retArray;
 
     // Test case: Nested exists
     sql =
-        "select distinct s1 from table1 t1 where exists(select s1 from table3 
t3 where t1.s1 = t3.s1 and exists(select s1 from table2 t2 where t2.s1 = t3.s1 
- 25))";
-    retArray = new String[] {"30,"};
+        "select distinct s1 from table1 t1 where s1 >= (select max(s1) from 
table3 t3 where t1.s1 = t3.s1 and s1 = (select max(s1) from table1 t1_2 where 
t1_2.s1 = t3.s1))";
+    retArray = new String[] {"30,", "40,"};
     expectedHeader = new String[] {"s1"};
     tableResultSetEqualTest(sql, expectedHeader, retArray, DATABASE_NAME);
-
-    // Test case: Nested exists with not(64 - 4 = 60 rows remain)
-    sql =
-        "select count(*) as cnt from table1 t1 where not exists(select s1 from 
table3 t3 where t1.s1 = t3.s1 and exists(select s1 from table2 t2 where t2.s1 = 
t3.s1 - 25))";
-    retArray = new String[] {"60,"};
-    expectedHeader = new String[] {"cnt"};
-    tableResultSetEqualTest(sql, expectedHeader, retArray, DATABASE_NAME);
   }
 
   @Test
-  public void testMultipleExistsSubquery() {
+  public void testMultipleScalarSubquery() {
     String sql;
     String[] expectedHeader;
     String[] retArray;
 
-    // Test case: multiple exists
+    // Test case: multiple scalar subquery
     sql =
-        "select distinct s1 from table1 t1 where exists(select s1 from table3 
t3 where t1.s1 = t3.s1) and exists(select s1 from table2 t2 where t2.s1 = t1.s1 
- 25)";
-    retArray = new String[] {"30,"};
+        "select distinct s1 from table1 t1 where s1 = (select max(s1) from 
table3 t3 where t1.s1 = t3.s1) and s1 = (select min(s1) from table3 t3 where 
t1.s1 = t3.s1)";
+    retArray = new String[] {"30,", "40,"};
     expectedHeader = new String[] {"s1"};
     tableResultSetEqualTest(sql, expectedHeader, retArray, DATABASE_NAME);
+  }
 
-    // Test case: multiple exists with not
+  @Test
+  public void 
testCorrelatedScalarSubqueryInWhereClauseWithOtherCorrelatedSubquery() {
+    String sql;
+    String[] expectedHeader;
+    String[] retArray;
+    // Test case: with Exists
     sql =
-        "select distinct s1 from table1 t1 where exists(select s1 from table3 
t3 where t1.s1 = t3.s1) and not exists(select s1 from table2 t2 where t2.s1 = 
t1.s1)";
+        "SELECT distinct cast(%s AS INT32) as %s FROM table1 t1 WHERE %s = 
(SELECT max(%s) from table3 t3 WHERE device_id = 'd01' and t1.%s = t3.%s) and 
exists(select s1 from table3 t3 where t1.s1 = t3.s1)";
     retArray = new String[] {"30,", "40,"};
-    expectedHeader = new String[] {"s1"};
-    tableResultSetEqualTest(sql, expectedHeader, retArray, DATABASE_NAME);
+    for (String measurement : NUMERIC_MEASUREMENTS) {
+      expectedHeader = new String[] {measurement};
+      tableResultSetEqualTest(
+          String.format(
+              sql,
+              measurement,
+              measurement,
+              measurement,
+              measurement,
+              measurement,
+              measurement,
+              measurement,
+              measurement),
+          expectedHeader,
+          retArray,
+          DATABASE_NAME);
+    }
   }
 
   @Test
-  public void 
testCorrelatedExistsSubqueryInWhereClauseWithOtherUncorrelatedSubquery() {
+  public void 
testCorrelatedScalarSubqueryInWhereClauseWithOtherUncorrelatedSubquery() {
     String sql;
     String[] expectedHeader;
     String[] retArray;
 
-    // Test case: with InPredicate
+    // Test case: with In Predicate
     sql =
-        "SELECT distinct cast(%s AS INT32) as %s FROM table1 t1 WHERE 
exists(SELECT (%s) from table3 t3 WHERE device_id = 'd01' and t1.%s = t3.%s) 
and %s in (select %s from table3)";
+        "SELECT distinct cast(%s AS INT32) as %s FROM table1 t1 WHERE %s = 
(SELECT max(%s) from table3 t3 WHERE device_id = 'd01' and t1.%s = t3.%s) and 
%s in (select %s from table3)";
     retArray = new String[] {"30,", "40,"};
     for (String measurement : NUMERIC_MEASUREMENTS) {
       expectedHeader = new String[] {measurement};
@@ -197,15 +194,16 @@ public class IoTDBCorrelatedExistsSubqueryIT {
               measurement,
               measurement,
               measurement,
+              measurement,
               measurement),
           expectedHeader,
           retArray,
           DATABASE_NAME);
     }
 
-    // Test case: with InPredicate in exists
+    // Test case: with InPredicate in scalar subquery
     sql =
-        "SELECT distinct cast(%s AS INT32) as %s FROM table1 t1 WHERE 
device_id = 'd01' and exists(SELECT (%s) from table3 t3 WHERE device_id = 'd01' 
and t1.%s = t3.%s and %s not in (select %s from table2 where %s is not null))";
+        "SELECT distinct cast(%s AS INT32) as %s FROM table1 t1 WHERE 
device_id = 'd01' and %s = (SELECT max(%s) from table3 t3 WHERE device_id = 
'd01' and t1.%s = t3.%s and %s not in (select %s from table2 where %s is not 
null))";
     retArray = new String[] {"30,", "40,"};
     for (String measurement : NUMERIC_MEASUREMENTS) {
       expectedHeader = new String[] {measurement};
@@ -219,6 +217,7 @@ public class IoTDBCorrelatedExistsSubqueryIT {
               measurement,
               measurement,
               measurement,
+              measurement,
               measurement),
           expectedHeader,
           retArray,
@@ -227,27 +226,41 @@ public class IoTDBCorrelatedExistsSubqueryIT {
 
     // Test case: with Scalar Subquery
     sql =
-        "SELECT cast(%s AS INT32) as %s FROM table1 t1 WHERE device_id = 'd01' 
and exists(SELECT (%s) from table3 t3 WHERE device_id = 'd01' and t1.%s = 
t3.%s) and s1 > (select min(%s) from table3)";
+        "SELECT cast(%s AS INT32) as %s FROM table1 t1 WHERE device_id = 'd01' 
and %s = (SELECT max(%s) from table3 t3 WHERE device_id = 'd01' and t1.%s = 
t3.%s) and s1 > (select min(%s) from table3)";
     retArray = new String[] {"40,"};
     for (String measurement : NUMERIC_MEASUREMENTS) {
       expectedHeader = new String[] {measurement};
       tableResultSetEqualTest(
           String.format(
-              sql, measurement, measurement, measurement, measurement, 
measurement, measurement),
+              sql,
+              measurement,
+              measurement,
+              measurement,
+              measurement,
+              measurement,
+              measurement,
+              measurement),
           expectedHeader,
           retArray,
           DATABASE_NAME);
     }
 
-    // Test case: with Scalar Subquery in exists
+    // Test case: with nested Scalar Subquery
     sql =
-        "SELECT distinct cast(%s AS INT32) as %s FROM table1 t1 WHERE 
device_id = 'd01' and exists(SELECT (%s) from table3 t3 WHERE device_id = 'd01' 
and t1.%s = t3.%s and s1 = (select min(%s) from table3))";
+        "SELECT distinct cast(%s AS INT32) as %s FROM table1 t1 WHERE 
device_id = 'd01' and %s = (SELECT (%s) from table3 t3 WHERE device_id = 'd01' 
and t1.%s = t3.%s and s1 = (select min(%s) from table3))";
     retArray = new String[] {"30,"};
     for (String measurement : NUMERIC_MEASUREMENTS) {
       expectedHeader = new String[] {measurement};
       tableResultSetEqualTest(
           String.format(
-              sql, measurement, measurement, measurement, measurement, 
measurement, measurement),
+              sql,
+              measurement,
+              measurement,
+              measurement,
+              measurement,
+              measurement,
+              measurement,
+              measurement),
           expectedHeader,
           retArray,
           DATABASE_NAME);
@@ -255,21 +268,28 @@ public class IoTDBCorrelatedExistsSubqueryIT {
 
     // Test case: with Quantified Comparison
     sql =
-        "SELECT distinct cast(%s AS INT32) as %s FROM table1 t1 WHERE 
device_id = 'd01' and exists(SELECT (%s) from table3 t3 WHERE device_id = 'd01' 
and t1.%s = t3.%s) and s1 != any(select %s from table2)";
+        "SELECT distinct cast(%s AS INT32) as %s FROM table1 t1 WHERE 
device_id = 'd01' and %s = (SELECT max(%s) from table3 t3 WHERE device_id = 
'd01' and t1.%s = t3.%s) and s1 != any(select %s from table2)";
     retArray = new String[] {"30,", "40,"};
     for (String measurement : NUMERIC_MEASUREMENTS) {
       expectedHeader = new String[] {measurement};
       tableResultSetEqualTest(
           String.format(
-              sql, measurement, measurement, measurement, measurement, 
measurement, measurement),
+              sql,
+              measurement,
+              measurement,
+              measurement,
+              measurement,
+              measurement,
+              measurement,
+              measurement),
           expectedHeader,
           retArray,
           DATABASE_NAME);
     }
 
-    // Test case: with Quantified Comparison in exists
+    // Test case: with Quantified Comparison in scalar subquery
     sql =
-        "SELECT distinct cast(%s AS INT32) as %s FROM table1 t1 WHERE 
device_id = 'd01' and exists(SELECT (%s) from table3 t3 WHERE device_id = 'd01' 
and t1.%s = t3.%s and %s = any (select %s from table3))";
+        "SELECT distinct cast(%s AS INT32) as %s FROM table1 t1 WHERE 
device_id = 'd01' and %s = (SELECT max(%s) from table3 t3 WHERE device_id = 
'd01' and t1.%s = t3.%s and %s = any (select %s from table3))";
     retArray = new String[] {"30,", "40,"};
     for (String measurement : NUMERIC_MEASUREMENTS) {
       expectedHeader = new String[] {measurement};
@@ -282,6 +302,7 @@ public class IoTDBCorrelatedExistsSubqueryIT {
               measurement,
               measurement,
               measurement,
+              measurement,
               measurement),
           expectedHeader,
           retArray,
@@ -290,31 +311,20 @@ public class IoTDBCorrelatedExistsSubqueryIT {
   }
 
   @Test
-  public void testCorrelatedExistsSubqueryWithMultipleCorrelation() {
+  public void testCorrelatedScalarSubqueryWithMultipleCorrelation() {
     String sql;
     String[] expectedHeader;
     String[] retArray;
 
     // Test case: Multiple correlation in exists
     sql =
-        "SELECT cast(%s AS INT32) as %s FROM table1 t1 WHERE device_id = 'd01' 
and exists(SELECT (%s) from table3 t3 WHERE t1.%s = t3.%s and t1.s1 = t3.s1 and 
t1.s2 = t3.s2)";
+        "SELECT cast(%s AS INT32) as %s FROM table1 t1 WHERE device_id = 'd01' 
and %s = (SELECT max(%s) from table3 t3 WHERE t1.%s = t3.%s and t1.s1 = t3.s1 
and t1.s2 = t3.s2)";
     retArray = new String[] {"30,", "40,"};
     for (String measurement : NUMERIC_MEASUREMENTS) {
       expectedHeader = new String[] {measurement};
       tableResultSetEqualTest(
-          String.format(sql, measurement, measurement, measurement, 
measurement, measurement),
-          expectedHeader,
-          retArray,
-          DATABASE_NAME);
-    }
-
-    sql =
-        "SELECT cast(%s AS INT32) as %s FROM table1 t1 WHERE device_id = 'd01' 
and not exists(SELECT (%s) from table3 t3 WHERE t1.%s = t3.%s and t1.s1 = t3.s1 
and t1.s2 = t3.s2)";
-    retArray = new String[] {"50,", "60,", "70,"};
-    for (String measurement : NUMERIC_MEASUREMENTS) {
-      expectedHeader = new String[] {measurement};
-      tableResultSetEqualTest(
-          String.format(sql, measurement, measurement, measurement, 
measurement, measurement),
+          String.format(
+              sql, measurement, measurement, measurement, measurement, 
measurement, measurement),
           expectedHeader,
           retArray,
           DATABASE_NAME);
@@ -322,111 +332,86 @@ public class IoTDBCorrelatedExistsSubqueryIT {
   }
 
   @Test
-  public void testCorrelatedExistsSubqueryInHavingClause() {
+  public void testCorrelatedScalarSubqueryInHavingClause() {
     String sql;
     String[] expectedHeader;
     String[] retArray;
 
-    // Test case: exists in having
+    // Test case: scalar subquery in having
     sql =
-        "SELECT device_id, count(*) from table1 t1 group by device_id having 
exists(SELECT 1 from table3 t3 where t3.device_id = t1.device_id)";
+        "SELECT device_id, count(*) from table1 t1 group by device_id having 
count(*) + 35 = (SELECT max(s1) from table3 t3 where t3.device_id = 
t1.device_id)";
     expectedHeader = new String[] {"device_id", "_col1"};
     retArray = new String[] {"d01,5,"};
     for (String measurement : NUMERIC_MEASUREMENTS) {
       tableResultSetEqualTest(
           String.format(sql, measurement), expectedHeader, retArray, 
DATABASE_NAME);
     }
-
-    sql =
-        "SELECT device_id, count(*) from table1 t1 group by device_id having 
exists(SELECT 1 from table3 t3 where t3.device_id != 'd01' and t3.device_id = 
t1.device_id)";
-    expectedHeader = new String[] {"device_id", "_col1"};
-    retArray = new String[] {};
-    for (String measurement : NUMERIC_MEASUREMENTS) {
-      tableResultSetEqualTest(
-          String.format(sql, measurement), expectedHeader, retArray, 
DATABASE_NAME);
-    }
-
-    // Test case: not exists in having
-    sql =
-        "SELECT device_id, count(*) from table1 t1 group by device_id having 
count(*) = 5 and not exists(SELECT 1 from table3 t3 where t3.device_id = 
t1.device_id)";
-    expectedHeader = new String[] {"device_id", "_col1"};
-    retArray = new String[] {"d03,5,", "d05,5,", "d07,5,", "d09,5,", "d11,5,", 
"d13,5,", "d15,5,"};
-    for (String measurement : NUMERIC_MEASUREMENTS) {
-      tableResultSetEqualTest(
-          String.format(sql, measurement), expectedHeader, retArray, 
DATABASE_NAME);
-    }
   }
 
   @Test
-  public void testUnCorrelatedExistsSubqueryInSelectClause() {
+  public void testCorrelatedScalarSubqueryInSelectClause() {
     String sql;
     String[] expectedHeader;
     String[] retArray;
 
     // Test case: exists in Select clause
     sql =
-        "select exists(select s1 from table1 t1 where t1.s1 = t3.s1) from 
table3 t3 where exists(select s1 from table2 t2 where t2.s1 = t3.s1 - 25)";
+        "select s1 = (select max(s1) from table1 t1 where t1.s1 = t3.s1) from 
table3 t3 where exists(select s1 from table2 t2 where t2.s1 = t3.s1 - 25)";
     retArray = new String[] {"true,", "true,"};
     expectedHeader = new String[] {"_col0"};
     tableResultSetEqualTest(sql, expectedHeader, retArray, DATABASE_NAME);
 
-    sql = "select exists(select s1 from table1 t1 where t1.s1 = t3.s1) from 
table3 t3";
-    retArray = new String[] {"true,", "true,", "true,", "false,"};
-    expectedHeader = new String[] {"_col0"};
-    tableResultSetEqualTest(sql, expectedHeader, retArray, DATABASE_NAME);
-
-    sql =
-        "select not exists(select s1 from table1 t1 where t1.s1 = t3.s1) from 
table3 t3 where exists(select s1 from table2 t2 where t2.s1 = t3.s1 - 25)";
-    retArray = new String[] {"false,", "false,"};
-    expectedHeader = new String[] {"_col0"};
-    tableResultSetEqualTest(sql, expectedHeader, retArray, DATABASE_NAME);
-
-    sql = "select not exists(select s1 from table1 t1 where t1.s1 = t3.s1) 
from table3 t3";
-    retArray = new String[] {"false,", "false,", "false,", "true,"};
+    sql = "select (select max(s1) from table1 t1 where t1.s1 = t3.s1) from 
table3 t3";
+    retArray = new String[] {"30,", "30,", "40,", "null,"};
     expectedHeader = new String[] {"_col0"};
     tableResultSetEqualTest(sql, expectedHeader, retArray, DATABASE_NAME);
   }
 
   @Test
-  public void testNonComparisonFilterInCorrelatedExistsSubquery() {
-    String errMsg = TSStatusCode.SEMANTIC_ERROR.getStatusCode() + ": " + 
ONLY_SUPPORT_EQUI_JOIN;
+  public void testNonEqualityComparisonFilterInCorrelatedScalarSubquery() {
     // Legality check: Correlated subquery with Non-equality comparison is not 
support for now.
     tableAssertTestFail(
-        "select s1 from table1 t1 where exists(select s1 from table3 t3 where 
t1.s1 > t3.s1)",
-        errMsg,
+        "select s1 from table1 t1 where s1 > (select max(s1) from table3 t3 
where t1.s1 > t3.s1)",
+        "For now, FullOuterJoin and LeftJoin only support EquiJoinClauses",
         DATABASE_NAME);
     tableAssertTestFail(
-        "select s1 from table1 t1 where exists(select s1 from table3 t3 where 
t1.s1 >= t3.s1)",
-        errMsg,
+        "select s1 from table1 t1 where s1 > (select max(s1) from table3 t3 
where t1.s1 >= t3.s1)",
+        "For now, FullOuterJoin and LeftJoin only support EquiJoinClauses",
         DATABASE_NAME);
     tableAssertTestFail(
-        "select s1 from table1 t1 where exists(select s1 from table3 t3 where 
t1.s1 < t3.s1)",
-        errMsg,
+        "select s1 from table1 t1 where s1 > (select max(s1) from table3 t3 
where t1.s1 < t3.s1)",
+        "For now, FullOuterJoin and LeftJoin only support EquiJoinClauses",
         DATABASE_NAME);
     tableAssertTestFail(
-        "select s1 from table1 t1 where exists(select s1 from table3 t3 where 
t1.s1 <= t3.s1)",
-        errMsg,
+        "select s1 from table1 t1 where s1 > (select max(s1) from table3 t3 
where t1.s1 <= t3.s1)",
+        "For now, FullOuterJoin and LeftJoin only support EquiJoinClauses",
         DATABASE_NAME);
     tableAssertTestFail(
-        "select s1 from table1 t1 where exists(select s1 from table3 t3 where 
t1.s1 != t3.s1)",
-        errMsg,
+        "select s1 from table1 t1 where s1 > (select max(s1) from table3 t3 
where t1.s1 != t3.s1)",
+        "For now, FullOuterJoin and LeftJoin only support EquiJoinClauses",
         DATABASE_NAME);
   }
 
   @Test
-  public void testCorrelatedExistsLegalityCheck() {
+  public void testCorrelatedScalarSubqueryLegalityCheck() {
     // Legality check: Correlated subqueries can only access columns from the 
immediately outer
     // scope and cannot access columns from the further outer queries.
     tableAssertTestFail(
-        "select s1 from table1 t1 where exists(select s1 from table3 t3 where 
t1.s1 = t3.s1 and exists(select s1 from table2 t2 where t2.s1 = t1.s1))",
+        "select s1 from table1 t1 where s1 > (select s1 from table3 t3 where 
t1.s1 = t3.s1 and s1 = (select s1 from table2 t2 where t2.s1 = t1.s1 limit 1) 
limit 1)",
         "701: Given correlated subquery is not supported",
         DATABASE_NAME);
 
     // Legality check: Correlated subqueries with limit clause and limit count 
greater than 1 is not
-    // supported for now
+    // supported for now.
+    tableAssertTestFail(
+        "select s1 from table3 t3 where 30 = t3.s1 and s1 > (select max(s1) 
from table2 t2 where t2.s1 = t3.s1 limit 2)",
+        "701: Given correlated subquery is not supported",
+        DATABASE_NAME);
+
+    // Legality check: Scalar subquery should return only one row.
     tableAssertTestFail(
-        "select s1 from table3 t3 where 30 = t3.s1 and exists(select s1 from 
table2 t2 where t2.s1 = t3.s1 limit 2)",
-        "Decorrelation for LIMIT with row count greater than 1 is not 
supported yet",
+        "select s1 from table1 t1 where s1 >= (select s1 from table3 t3 where 
t3.s1 = t1.s1)",
+        "701: Scalar sub-query has returned multiple rows",
         DATABASE_NAME);
   }
 }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/FilterAndProjectOperator.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/FilterAndProjectOperator.java
index 06ca6d6d2ce..7c32d2c20a8 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/FilterAndProjectOperator.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/FilterAndProjectOperator.java
@@ -24,6 +24,7 @@ import 
org.apache.iotdb.db.queryengine.execution.operator.Operator;
 import org.apache.iotdb.db.queryengine.execution.operator.OperatorContext;
 import 
org.apache.iotdb.db.queryengine.transformation.dag.column.AbstractCaseWhenThenColumnTransformer;
 import 
org.apache.iotdb.db.queryengine.transformation.dag.column.ColumnTransformer;
+import 
org.apache.iotdb.db.queryengine.transformation.dag.column.FailFunctionColumnTransformer;
 import 
org.apache.iotdb.db.queryengine.transformation.dag.column.binary.BinaryColumnTransformer;
 import 
org.apache.iotdb.db.queryengine.transformation.dag.column.leaf.IdentityColumnTransformer;
 import 
org.apache.iotdb.db.queryengine.transformation.dag.column.leaf.LeafColumnTransformer;
@@ -422,6 +423,8 @@ public class FilterAndProjectOperator implements 
ProcessOperator {
                       .getElseTransformer()));
       childMaxLevel = Math.max(childMaxLevel, childCount + 2);
       return childMaxLevel;
+    } else if (columnTransformer instanceof FailFunctionColumnTransformer) {
+      return 0;
     } else {
       throw new UnsupportedOperationException("Unsupported ColumnTransformer");
     }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/relational/ColumnTransformerBuilder.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/relational/ColumnTransformerBuilder.java
index 472899e1d0c..d0b7e35e94c 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/relational/ColumnTransformerBuilder.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/relational/ColumnTransformerBuilder.java
@@ -72,6 +72,7 @@ import 
org.apache.iotdb.db.queryengine.plan.relational.type.InternalTypeManager;
 import 
org.apache.iotdb.db.queryengine.plan.relational.type.TypeNotFoundException;
 import org.apache.iotdb.db.queryengine.plan.udf.TableUDFUtils;
 import 
org.apache.iotdb.db.queryengine.transformation.dag.column.ColumnTransformer;
+import 
org.apache.iotdb.db.queryengine.transformation.dag.column.FailFunctionColumnTransformer;
 import 
org.apache.iotdb.db.queryengine.transformation.dag.column.TableCaseWhenThenColumnTransformer;
 import 
org.apache.iotdb.db.queryengine.transformation.dag.column.binary.ArithmeticColumnTransformerApi;
 import 
org.apache.iotdb.db.queryengine.transformation.dag.column.binary.CompareEqualToColumnTransformer;
@@ -187,10 +188,12 @@ import java.util.Optional;
 import java.util.Set;
 import java.util.stream.Collectors;
 
+import static com.google.common.base.Preconditions.checkArgument;
 import static 
org.apache.iotdb.db.queryengine.plan.expression.unary.LikeExpression.getEscapeCharacter;
 import static 
org.apache.iotdb.db.queryengine.plan.relational.analyzer.predicate.PredicatePushIntoMetadataChecker.isStringLiteral;
 import static 
org.apache.iotdb.db.queryengine.plan.relational.type.InternalTypeManager.getTSDataType;
 import static 
org.apache.iotdb.db.queryengine.plan.relational.type.TypeSignatureTranslator.toTypeSignature;
+import static 
org.apache.iotdb.db.queryengine.transformation.dag.column.FailFunctionColumnTransformer.FAIL_FUNCTION_NAME;
 import static org.apache.tsfile.read.common.type.BlobType.BLOB;
 import static org.apache.tsfile.read.common.type.BooleanType.BOOLEAN;
 import static org.apache.tsfile.read.common.type.DoubleType.DOUBLE;
@@ -978,6 +981,10 @@ public class ColumnTransformerBuilder
       }
       return new FormatColumnTransformer(
           STRING, columnTransformers, context.sessionInfo.getZoneId());
+    } else if (FAIL_FUNCTION_NAME.equalsIgnoreCase(functionName)) {
+      checkArgument(children.size() == 1 && children.get(0) instanceof 
StringLiteral);
+      return new FailFunctionColumnTransformer(
+          STRING, ((StringLiteral) children.get(0)).getValue());
     } else if (TableBuiltinScalarFunction.GREATEST
         .getFunctionName()
         .equalsIgnoreCase(functionName)) {
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/metadata/TableMetadataImpl.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/metadata/TableMetadataImpl.java
index 369522254dd..742f980e2e8 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/metadata/TableMetadataImpl.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/metadata/TableMetadataImpl.java
@@ -75,6 +75,7 @@ import java.util.Map;
 import java.util.Optional;
 import java.util.stream.Collectors;
 
+import static 
org.apache.iotdb.db.queryengine.transformation.dag.column.FailFunctionColumnTransformer.FAIL_FUNCTION_NAME;
 import static org.apache.tsfile.read.common.type.BinaryType.TEXT;
 import static org.apache.tsfile.read.common.type.BooleanType.BOOLEAN;
 import static org.apache.tsfile.read.common.type.DateType.DATE;
@@ -554,6 +555,8 @@ public class TableMetadataImpl implements Metadata {
                 + " must have at least two arguments, and first argument 
pattern must be TEXT or STRING type.");
       }
       return STRING;
+    } else if (FAIL_FUNCTION_NAME.equalsIgnoreCase(functionName)) {
+      return UNKNOWN;
     } else if 
(TableBuiltinScalarFunction.GREATEST.getFunctionName().equalsIgnoreCase(functionName)
         || 
TableBuiltinScalarFunction.LEAST.getFunctionName().equalsIgnoreCase(functionName))
 {
       if (argumentTypes.size() < 2 || 
!areAllTypesSameAndComparable(argumentTypes)) {
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/iterative/rule/TransformCorrelatedScalarSubquery.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/iterative/rule/TransformCorrelatedScalarSubquery.java
new file mode 100644
index 00000000000..3dcb42f5ddd
--- /dev/null
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/iterative/rule/TransformCorrelatedScalarSubquery.java
@@ -0,0 +1,194 @@
+/*
+ * 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.iotdb.db.queryengine.plan.relational.planner.iterative.rule;
+
+import org.apache.iotdb.db.queryengine.plan.planner.plan.node.PlanNode;
+import org.apache.iotdb.db.queryengine.plan.relational.metadata.Metadata;
+import org.apache.iotdb.db.queryengine.plan.relational.planner.Assignments;
+import org.apache.iotdb.db.queryengine.plan.relational.planner.Symbol;
+import org.apache.iotdb.db.queryengine.plan.relational.planner.iterative.Rule;
+import 
org.apache.iotdb.db.queryengine.plan.relational.planner.node.AssignUniqueId;
+import 
org.apache.iotdb.db.queryengine.plan.relational.planner.node.CorrelatedJoinNode;
+import 
org.apache.iotdb.db.queryengine.plan.relational.planner.node.EnforceSingleRowNode;
+import org.apache.iotdb.db.queryengine.plan.relational.planner.node.FilterNode;
+import 
org.apache.iotdb.db.queryengine.plan.relational.planner.node.MarkDistinctNode;
+import 
org.apache.iotdb.db.queryengine.plan.relational.planner.node.ProjectNode;
+import 
org.apache.iotdb.db.queryengine.plan.relational.planner.optimizations.Cardinality;
+import org.apache.iotdb.db.queryengine.plan.relational.sql.ast.Cast;
+import org.apache.iotdb.db.queryengine.plan.relational.sql.ast.FunctionCall;
+import org.apache.iotdb.db.queryengine.plan.relational.sql.ast.QualifiedName;
+import 
org.apache.iotdb.db.queryengine.plan.relational.sql.ast.SimpleCaseExpression;
+import org.apache.iotdb.db.queryengine.plan.relational.sql.ast.StringLiteral;
+import org.apache.iotdb.db.queryengine.plan.relational.sql.ast.WhenClause;
+import org.apache.iotdb.db.queryengine.plan.relational.utils.matching.Captures;
+import org.apache.iotdb.db.queryengine.plan.relational.utils.matching.Pattern;
+
+import com.google.common.collect.ImmutableList;
+
+import java.util.Optional;
+
+import static com.google.common.base.Preconditions.checkArgument;
+import static java.util.Objects.requireNonNull;
+import static 
org.apache.iotdb.db.queryengine.plan.relational.planner.PlanNodeSearcher.searchFrom;
+import static 
org.apache.iotdb.db.queryengine.plan.relational.planner.node.JoinNode.JoinType.INNER;
+import static 
org.apache.iotdb.db.queryengine.plan.relational.planner.node.JoinNode.JoinType.LEFT;
+import static 
org.apache.iotdb.db.queryengine.plan.relational.planner.node.Patterns.CorrelatedJoin.correlation;
+import static 
org.apache.iotdb.db.queryengine.plan.relational.planner.node.Patterns.CorrelatedJoin.filter;
+import static 
org.apache.iotdb.db.queryengine.plan.relational.planner.node.Patterns.correlatedJoin;
+import static 
org.apache.iotdb.db.queryengine.plan.relational.planner.optimizations.QueryCardinalityUtil.extractCardinality;
+import static 
org.apache.iotdb.db.queryengine.plan.relational.sql.ast.BooleanLiteral.TRUE_LITERAL;
+import static 
org.apache.iotdb.db.queryengine.plan.relational.type.TypeSignatureTranslator.toSqlType;
+import static 
org.apache.iotdb.db.queryengine.plan.relational.utils.matching.Pattern.nonEmpty;
+import static 
org.apache.iotdb.db.queryengine.transformation.dag.column.FailFunctionColumnTransformer.FAIL_FUNCTION_NAME;
+import static org.apache.tsfile.read.common.type.BooleanType.BOOLEAN;
+import static org.apache.tsfile.read.common.type.LongType.INT64;
+
+/**
+ * Scalar filter scan query is something like:
+ *
+ * <pre>
+ *     SELECT a,b,c FROM rel WHERE a = correlated1 AND b = correlated2
+ * </pre>
+ *
+ * <p>This optimizer can rewrite to mark distinct and filter over a left outer 
join:
+ *
+ * <p>From:
+ *
+ * <pre>
+ * - CorrelatedJoin (with correlation list: [C])
+ *   - (input) plan which produces symbols: [A, B, C]
+ *   - (scalar subquery) Project F
+ *     - Filter(D = C AND E > 5)
+ *       - plan which produces symbols: [D, E, F]
+ * </pre>
+ *
+ * to:
+ *
+ * <pre>
+ * - Filter(CASE isDistinct WHEN true THEN true ELSE fail('Scalar sub-query 
has returned multiple rows'))
+ *   - MarkDistinct(isDistinct)
+ *     - CorrelatedJoin (with correlation list: [C])
+ *       - AssignUniqueId(adds symbol U)
+ *         - (input) plan which produces symbols: [A, B, C]
+ *       - non scalar subquery
+ * </pre>
+ *
+ * <p>This must be run after aggregation decorrelation rules.
+ *
+ * <p>This rule is used to support non-aggregation scalar subquery.
+ */
+public class TransformCorrelatedScalarSubquery implements 
Rule<CorrelatedJoinNode> {
+  private static final Pattern<CorrelatedJoinNode> PATTERN =
+      
correlatedJoin().with(nonEmpty(correlation())).with(filter().equalTo(TRUE_LITERAL));
+
+  private final Metadata metadata;
+
+  public TransformCorrelatedScalarSubquery(Metadata metadata) {
+    this.metadata = requireNonNull(metadata, "metadata is null");
+  }
+
+  @Override
+  public Pattern<CorrelatedJoinNode> getPattern() {
+    return PATTERN;
+  }
+
+  @Override
+  public Result apply(CorrelatedJoinNode correlatedJoinNode, Captures 
captures, Context context) {
+    // lateral references are only allowed for INNER or LEFT correlated join
+    checkArgument(
+        correlatedJoinNode.getJoinType() == INNER || 
correlatedJoinNode.getJoinType() == LEFT,
+        "unexpected correlated join type: %s",
+        correlatedJoinNode.getJoinType());
+    PlanNode subquery = 
context.getLookup().resolve(correlatedJoinNode.getSubquery());
+
+    if (!searchFrom(subquery, context.getLookup())
+        .where(EnforceSingleRowNode.class::isInstance)
+        .recurseOnlyWhen(ProjectNode.class::isInstance)
+        .matches()) {
+      return Result.empty();
+    }
+
+    PlanNode rewrittenSubquery =
+        searchFrom(subquery, context.getLookup())
+            .where(EnforceSingleRowNode.class::isInstance)
+            .recurseOnlyWhen(ProjectNode.class::isInstance)
+            .removeFirst();
+
+    Cardinality subqueryCardinality = extractCardinality(rewrittenSubquery, 
context.getLookup());
+    boolean producesAtMostOneRow = subqueryCardinality.isAtMostScalar();
+    if (producesAtMostOneRow) {
+      boolean producesSingleRow = subqueryCardinality.isScalar();
+      return Result.ofPlanNode(
+          new CorrelatedJoinNode(
+              context.getIdAllocator().genPlanNodeId(),
+              correlatedJoinNode.getInput(),
+              rewrittenSubquery,
+              correlatedJoinNode.getCorrelation(),
+              // EnforceSingleRowNode guarantees that exactly single matching 
row is produced
+              // for every input row (independently of correlated join type). 
Decorrelated plan
+              // must preserve this semantics.
+              producesSingleRow ? INNER : LEFT,
+              correlatedJoinNode.getFilter(),
+              correlatedJoinNode.getOriginSubquery()));
+    }
+
+    Symbol unique = context.getSymbolAllocator().newSymbol("unique", INT64);
+
+    CorrelatedJoinNode rewrittenCorrelatedJoinNode =
+        new CorrelatedJoinNode(
+            context.getIdAllocator().genPlanNodeId(),
+            new AssignUniqueId(
+                context.getIdAllocator().genPlanNodeId(), 
correlatedJoinNode.getInput(), unique),
+            rewrittenSubquery,
+            correlatedJoinNode.getCorrelation(),
+            LEFT,
+            correlatedJoinNode.getFilter(),
+            correlatedJoinNode.getOriginSubquery());
+
+    Symbol isDistinct = context.getSymbolAllocator().newSymbol("is_distinct", 
BOOLEAN);
+    MarkDistinctNode markDistinctNode =
+        new MarkDistinctNode(
+            context.getIdAllocator().genPlanNodeId(),
+            rewrittenCorrelatedJoinNode,
+            isDistinct,
+            rewrittenCorrelatedJoinNode.getInput().getOutputSymbols(),
+            Optional.empty());
+
+    FilterNode filterNode =
+        new FilterNode(
+            context.getIdAllocator().genPlanNodeId(),
+            markDistinctNode,
+            new SimpleCaseExpression(
+                isDistinct.toSymbolReference(),
+                ImmutableList.of(new WhenClause(TRUE_LITERAL, TRUE_LITERAL)),
+                new Cast(
+                    new FunctionCall(
+                        QualifiedName.of(FAIL_FUNCTION_NAME),
+                        ImmutableList.of(
+                            new StringLiteral("Scalar sub-query has returned 
multiple rows."))),
+                    toSqlType(BOOLEAN))));
+
+    return Result.ofPlanNode(
+        new ProjectNode(
+            context.getIdAllocator().genPlanNodeId(),
+            filterNode,
+            Assignments.identity(correlatedJoinNode.getOutputSymbols())));
+  }
+}
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/iterative/rule/TransformExistsApplyToCorrelatedJoin.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/iterative/rule/TransformExistsApplyToCorrelatedJoin.java
index 4362620119e..0e864eed386 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/iterative/rule/TransformExistsApplyToCorrelatedJoin.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/iterative/rule/TransformExistsApplyToCorrelatedJoin.java
@@ -116,6 +116,7 @@ public class TransformExistsApplyToCorrelatedJoin 
implements Rule<ApplyNode> {
     To support the latter case, the ApplyNode with empty correlation list is 
rewritten to default
     aggregation, which is inefficient in the rare case of uncorrelated EXISTS 
subquery,
     but currently allows to successfully decorrelate a correlated EXISTS 
subquery.
+
     Perhaps we can remove this condition when exploratory optimizer is 
implemented or support for decorrelating joins is implemented in 
PlanNodeDecorrelator
     */
     if (parent.getCorrelation().isEmpty()) {
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/optimizations/LogicalOptimizeFactory.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/optimizations/LogicalOptimizeFactory.java
index 59df50bd861..ba04ab80a7e 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/optimizations/LogicalOptimizeFactory.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/optimizations/LogicalOptimizeFactory.java
@@ -78,6 +78,7 @@ import 
org.apache.iotdb.db.queryengine.plan.relational.planner.iterative.rule.Tr
 import 
org.apache.iotdb.db.queryengine.plan.relational.planner.iterative.rule.TransformCorrelatedGroupedAggregationWithProjection;
 import 
org.apache.iotdb.db.queryengine.plan.relational.planner.iterative.rule.TransformCorrelatedGroupedAggregationWithoutProjection;
 import 
org.apache.iotdb.db.queryengine.plan.relational.planner.iterative.rule.TransformCorrelatedJoinToJoin;
+import 
org.apache.iotdb.db.queryengine.plan.relational.planner.iterative.rule.TransformCorrelatedScalarSubquery;
 import 
org.apache.iotdb.db.queryengine.plan.relational.planner.iterative.rule.TransformExistsApplyToCorrelatedJoin;
 import 
org.apache.iotdb.db.queryengine.plan.relational.planner.iterative.rule.TransformUncorrelatedInPredicateSubqueryToSemiJoin;
 import 
org.apache.iotdb.db.queryengine.plan.relational.planner.iterative.rule.TransformUncorrelatedSubqueryToJoin;
@@ -265,8 +266,8 @@ public class LogicalOptimizeFactory {
                 new RemoveUnreferencedScalarApplyNodes(),
                 //                            new 
TransformCorrelatedInPredicateToJoin(metadata), //
                 // must be run after columnPruningOptimizer
-                //                            new 
TransformCorrelatedScalarSubquery(metadata), //
-                // must be run after TransformCorrelatedAggregation rules
+                new TransformCorrelatedScalarSubquery(
+                    metadata), // must be run after 
TransformCorrelatedAggregation rules
                 new TransformCorrelatedJoinToJoin(plannerContext))),
         new IterativeOptimizer(
             plannerContext,
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/optimizations/PlanNodeDecorrelator.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/optimizations/PlanNodeDecorrelator.java
index 17effd38343..4011370adc7 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/optimizations/PlanNodeDecorrelator.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/optimizations/PlanNodeDecorrelator.java
@@ -210,7 +210,7 @@ public class PlanNodeDecorrelator {
 
     // Limit (1) could be decorrelated by the method 
rewriteLimitWithRowCountGreaterThanOne()
     // as well.
-    // The current decorrelation method for Limit (1) cannot deal with 
subqueries outputting other
+    // The current decorrelation method for Limit (1) can not deal with 
subqueries outputting other
     // symbols
     // than constants.
     //
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/transformation/dag/column/FailFunctionColumnTransformer.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/transformation/dag/column/FailFunctionColumnTransformer.java
new file mode 100644
index 00000000000..41632d7be1a
--- /dev/null
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/transformation/dag/column/FailFunctionColumnTransformer.java
@@ -0,0 +1,62 @@
+/*
+ * 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.iotdb.db.queryengine.transformation.dag.column;
+
+import org.apache.iotdb.db.exception.sql.SemanticException;
+
+import org.apache.tsfile.read.common.block.column.NullColumn;
+import org.apache.tsfile.read.common.type.Type;
+
+/**
+ * when evaluate of FailFunctionColumnTransformer is called, throw Exception 
to inform the user with
+ * specified errorMsg
+ */
+public class FailFunctionColumnTransformer extends ColumnTransformer {
+  public static final String FAIL_FUNCTION_NAME = "fail";
+
+  private final String errorMsg;
+
+  public FailFunctionColumnTransformer(Type returnType, String errorMsg) {
+    super(returnType);
+    this.errorMsg = errorMsg;
+  }
+
+  @Override
+  protected void evaluate() {
+    throw new SemanticException(errorMsg);
+  }
+
+  @Override
+  public void evaluateWithSelection(boolean[] selection) {
+    // if there is true value in selection, throw exception
+    for (boolean b : selection) {
+      if (b) {
+        throw new SemanticException(errorMsg);
+      }
+    }
+    // Result of fail function should be ignored, we just fake the output here.
+    initializeColumnCache(new NullColumn(1));
+  }
+
+  @Override
+  protected void checkType() {
+    // do nothing
+  }
+}
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/transformation/dag/column/unary/ArithmeticNegationColumnTransformer.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/transformation/dag/column/unary/ArithmeticNegationColumnTransformer.java
index b5dc095380b..81819dddf81 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/transformation/dag/column/unary/ArithmeticNegationColumnTransformer.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/transformation/dag/column/unary/ArithmeticNegationColumnTransformer.java
@@ -58,7 +58,8 @@ public class ArithmeticNegationColumnTransformer extends 
UnaryColumnTransformer
   @Override
   protected final void checkType() {
     if (!childColumnTransformer.isReturnTypeNumeric()) {
-      throw new UnsupportedOperationException("Unsupported Type: " + 
returnType.toString());
+      throw new UnsupportedOperationException(
+          "Unsupported Type: " + childColumnTransformer.getType().toString());
     }
   }
 }

Reply via email to