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

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


The following commit(s) were added to refs/heads/master by this push:
     new 28577df9f6d [fix](fe) Reject unsafe casts in storage MIN/MAX pushdown 
(#67979)
28577df9f6d is described below

commit 28577df9f6df3446f7fd1d0c90032fb7d85c7733
Author: Gabriel <[email protected]>
AuthorDate: Tue Sep 15 17:52:08 2026 +0800

    [fix](fe) Reject unsafe casts in storage MIN/MAX pushdown (#67979)
    
    ### What problem does this PR solve?
    
    Storage MIN/MAX pushdown can discard values needed after a numeric cast.
    With BIGINT values `-2147483649`, `0`, and `2147483648`, reducing a file
    or OLAP segment to its source endpoints makes MIN/MAX(CAST(value AS
    INT)) return NULL instead of 0. Floating casts also require care: footer
    extrema may omit NaN, and DOUBLE/DECIMAL-to-FLOAT underflow can change
    the signed-zero representative.
    
    Apply the cast safety check to both file and OLAP scans, after resolving
    projection aliases. Reject casts that introduce NULL, casts from
    floating-point sources, and DECIMAL-to-FLOAT casts. Check cast
    nullability independently of source nullability so safe widening casts
    remain eligible. COUNT-only behavior is unchanged.
    
    ### Release note
    
    Fix incorrect MIN/MAX results from metadata pushdown with overflow-prone
    casts, floating-source casts, and DECIMAL-to-FLOAT casts.
    
    ### Check List (For Author)
    
    - Test:
    - [x] FE unit tests: `PhysicalStorageLayerAggregateTest` and
    `AggregateStrategiesTest` (18 passed, 3 existing skips). The new
    floating-cast and OLAP tests failed before the fix and pass after it.
    Coverage includes direct/projected arguments, nullable columns, strict
    casts, and safe widening controls.
        - [x] FE Checkstyle.
    - [x] Regression coverage: Paimon pushdown-on/off comparisons for
    overflow, signed zero, and NaN; native routing is required with
    `paimonNativeReadSplits=1/1`. A separate OLAP suite covers
    direct/projected overflow queries and safe controls.
    - [x] Both regression files pass Groovy syntax parsing. Local execution
    was attempted, but the FE connection was refused before SQL execution;
    end-to-end integration remains unverified. An unrelated existing
    `IvmNormalizeMTMVJoinTest` constructor mismatch was temporarily excluded
    from unit-test compilation and restored afterward.
    - Behavior changed:
    - [x] Yes. Unsafe MIN/MAX casts evaluate rows before aggregation;
    supported safe casts retain metadata pushdown.
    - Does this need documentation?
        - [x] No.
---
 .../rules/implementation/AggregateStrategies.java  |  14 +++
 .../rewrite/PhysicalStorageLayerAggregateTest.java | 125 ++++++++++++++++++++-
 .../paimon/test_paimon_minmax_cast.groovy          | 105 +++++++++++++++++
 .../nereids_p0/test_minmax_cast_pushdown.groovy    |  64 +++++++++++
 4 files changed, 302 insertions(+), 6 deletions(-)

diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/implementation/AggregateStrategies.java
 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/implementation/AggregateStrategies.java
index 906ba73052b..bc6438698b8 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/implementation/AggregateStrategies.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/implementation/AggregateStrategies.java
@@ -693,6 +693,20 @@ public class AggregateStrategies implements 
ImplementationRuleFactory {
             return canNotPush;
         }
 
+        // File footers and OLAP zone maps retain only source endpoints. Casts 
that introduce NULL
+        // can discard a valid interior value. Check the cast independently of 
source nullability
+        // so safe widening casts over nullable columns remain eligible. 
Floating sources may have
+        // NaNs omitted by file statistics; DOUBLE/DECIMAL-to-FLOAT can also 
underflow to signed
+        // zero and change the MIN/MAX representative even without introducing 
NULL.
+        if ((functionClasses.contains(Min.class) || 
functionClasses.contains(Max.class))
+                && argumentsOfAggregateFunction.stream().anyMatch(argument -> 
argument instanceof Cast
+                        && (Cast.castNullable(false, 
argument.child(0).getDataType(), argument.getDataType())
+                                || 
argument.child(0).getDataType().isFloatLikeType()
+                                || 
(argument.child(0).getDataType().isDecimalLikeType()
+                                        && 
argument.getDataType().isFloatType())))) {
+            return canNotPush;
+        }
+
         Set<PushDownAggOp> pushDownAggOps = functionClasses.stream()
                 .map(supportedAgg::get)
                 .collect(Collectors.toSet());
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/nereids/rules/rewrite/PhysicalStorageLayerAggregateTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/nereids/rules/rewrite/PhysicalStorageLayerAggregateTest.java
index 74708104e35..31ffd59d27a 100644
--- 
a/fe/fe-core/src/test/java/org/apache/doris/nereids/rules/rewrite/PhysicalStorageLayerAggregateTest.java
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/nereids/rules/rewrite/PhysicalStorageLayerAggregateTest.java
@@ -31,11 +31,14 @@ import org.apache.doris.nereids.rules.RulePromise;
 import org.apache.doris.nereids.rules.RuleType;
 import org.apache.doris.nereids.rules.implementation.AggregateStrategies;
 import org.apache.doris.nereids.trees.expressions.Alias;
+import org.apache.doris.nereids.trees.expressions.Cast;
+import org.apache.doris.nereids.trees.expressions.Expression;
 import org.apache.doris.nereids.trees.expressions.IsNull;
 import org.apache.doris.nereids.trees.expressions.functions.agg.Count;
 import org.apache.doris.nereids.trees.expressions.functions.agg.Max;
 import org.apache.doris.nereids.trees.expressions.functions.agg.Min;
 import org.apache.doris.nereids.trees.expressions.functions.scalar.Ln;
+import org.apache.doris.nereids.trees.plans.Plan;
 import org.apache.doris.nereids.trees.plans.RelationId;
 import org.apache.doris.nereids.trees.plans.logical.LogicalAggregate;
 import org.apache.doris.nereids.trees.plans.logical.LogicalFileScan;
@@ -43,7 +46,14 @@ import 
org.apache.doris.nereids.trees.plans.logical.LogicalFileScan.SelectedPart
 import org.apache.doris.nereids.trees.plans.logical.LogicalFilter;
 import org.apache.doris.nereids.trees.plans.logical.LogicalOlapScan;
 import org.apache.doris.nereids.trees.plans.logical.LogicalProject;
+import org.apache.doris.nereids.trees.plans.logical.LogicalRelation;
 import 
org.apache.doris.nereids.trees.plans.physical.PhysicalStorageLayerAggregate.PushDownAggOp;
+import org.apache.doris.nereids.types.BigIntType;
+import org.apache.doris.nereids.types.DataType;
+import org.apache.doris.nereids.types.DecimalV3Type;
+import org.apache.doris.nereids.types.DoubleType;
+import org.apache.doris.nereids.types.FloatType;
+import org.apache.doris.nereids.types.IntegerType;
 import org.apache.doris.nereids.util.MemoPatternMatchSupported;
 import org.apache.doris.nereids.util.MemoTestUtils;
 import org.apache.doris.nereids.util.PlanChecker;
@@ -170,7 +180,15 @@ public class PhysicalStorageLayerAggregateTest implements 
MemoPatternMatchSuppor
     }
 
     private LogicalAggregate<LogicalFileScan> newNullableFileCountAggregate() {
-        Column nullableColumn = new Column("value", Type.INT, true);
+        LogicalFileScan fileScan = newFileScan(Type.INT, true);
+        return new LogicalAggregate<>(
+                Collections.emptyList(),
+                ImmutableList.of(new Alias(new 
Count(fileScan.getOutput().get(0)), "count")),
+                true, Optional.empty(), fileScan);
+    }
+
+    private LogicalFileScan newFileScan(Type type, boolean nullable) {
+        Column nullableColumn = new Column("value", type, nullable);
         PluginDrivenExternalTable table = 
Mockito.mock(PluginDrivenExternalTable.class);
         Mockito.when(table.initSelectedPartitions(Mockito.any()))
                 .thenReturn(SelectedPartitions.NOT_PRUNED);
@@ -186,13 +204,108 @@ public class PhysicalStorageLayerAggregateTest 
implements MemoPatternMatchSuppor
         Mockito.when(database.getCatalog()).thenReturn(catalog);
         Mockito.when(database.getFullName()).thenReturn("db");
         Mockito.when(table.getDatabase()).thenReturn(database);
-        LogicalFileScan fileScan = new LogicalFileScan(new RelationId(1), 
table,
+        return new LogicalFileScan(new RelationId(1), table,
                 ImmutableList.of("catalog", "db"), Collections.emptyList(),
                 Optional.empty(), Optional.empty(), Optional.empty(), 
Optional.empty());
-        return new LogicalAggregate<>(
-                Collections.emptyList(),
-                ImmutableList.of(new Alias(new 
Count(fileScan.getOutput().get(0)), "count")),
-                true, Optional.empty(), fileScan);
+    }
+
+    @Test
+    public void testFileMinMaxUnsafeCast() {
+        for (boolean projected : new boolean[] {false, true}) {
+            for (boolean nullable : new boolean[] {false, true}) {
+                for (boolean strict : new boolean[] {false, true}) {
+                    checkFileMinMaxCast(Type.BIGINT, IntegerType.INSTANCE, 
nullable, projected, strict, false);
+                    checkFileMinMaxCast(Type.DOUBLE, IntegerType.INSTANCE, 
nullable, projected, strict, false);
+                    checkFileMinMaxCast(DecimalV3Type.createDecimalV3Type(3, 
2).toCatalogDataType(),
+                            DecimalV3Type.createDecimalV3Type(2, 1), nullable, 
projected, strict, false);
+                }
+            }
+        }
+    }
+
+    @Test
+    public void testFileMinMaxSafeCast() {
+        for (boolean projected : new boolean[] {false, true}) {
+            for (boolean nullable : new boolean[] {false, true}) {
+                checkFileMinMaxCast(Type.INT, BigIntType.INSTANCE, nullable, 
projected, false, true);
+                checkFileMinMaxCast(Type.INT, DoubleType.INSTANCE, nullable, 
projected, false, true);
+                checkFileMinMaxCast(DecimalV3Type.createDecimalV3Type(3, 
2).toCatalogDataType(),
+                        DecimalV3Type.createDecimalV3Type(4, 2), nullable, 
projected, false, true);
+            }
+        }
+    }
+
+    @Test
+    public void testFileMinMaxFloatingCast() {
+        for (boolean projected : new boolean[] {false, true}) {
+            for (boolean nullable : new boolean[] {false, true}) {
+                for (boolean strict : new boolean[] {false, true}) {
+                    checkFileMinMaxCast(Type.DOUBLE, FloatType.INSTANCE, 
nullable, projected, strict, false);
+                    checkFileMinMaxCast(Type.FLOAT, DoubleType.INSTANCE, 
nullable, projected, strict, false);
+                    
checkFileMinMaxCast(DecimalV3Type.createDecimalV3TypeNotCheck256(76, 
60).toCatalogDataType(),
+                            FloatType.INSTANCE, nullable, projected, strict, 
false);
+                }
+            }
+        }
+    }
+
+    @Test
+    public void testOlapMinMaxCast() {
+        for (boolean projected : new boolean[] {false, true}) {
+            for (boolean nullable : new boolean[] {false, true}) {
+                for (boolean strict : new boolean[] {false, true}) {
+                    checkOlapMinMaxCast(Type.BIGINT, IntegerType.INSTANCE, 
nullable, projected, strict, false);
+                    checkOlapMinMaxCast(Type.DOUBLE, FloatType.INSTANCE, 
nullable, projected, strict, false);
+                    checkOlapMinMaxCast(Type.INT, BigIntType.INSTANCE, 
nullable, projected, strict, true);
+                }
+            }
+        }
+    }
+
+    private void checkOlapMinMaxCast(Type sourceType, DataType targetType, 
boolean nullable,
+            boolean projected, boolean strict, boolean expectedPushdown) {
+        LogicalOlapScan scan = PlanConstructor.newLogicalOlapScan(1, 
"cast_table", 0);
+        scan.getTable().getFullSchema().get(0).setType(sourceType);
+        scan.getTable().getFullSchema().get(0).setIsAllowNull(nullable);
+        checkMinMaxCast(scan, targetType, projected, strict, expectedPushdown);
+    }
+
+    private void checkFileMinMaxCast(Type sourceType, DataType targetType, 
boolean nullable,
+            boolean projected, boolean strict, boolean expectedPushdown) {
+        checkMinMaxCast(newFileScan(sourceType, nullable), targetType, 
projected, strict, expectedPushdown);
+    }
+
+    private void checkMinMaxCast(LogicalRelation scan, DataType targetType,
+            boolean projected, boolean strict, boolean expectedPushdown) {
+        Expression argument = new Cast(scan.getOutput().get(0), targetType, 
true, strict);
+        Plan child = scan;
+        RuleType ruleType = scan instanceof LogicalFileScan
+                ? 
RuleType.STORAGE_LAYER_AGGREGATE_WITHOUT_PROJECT_FOR_FILE_SCAN
+                : RuleType.STORAGE_LAYER_AGGREGATE_WITHOUT_PROJECT;
+        if (projected) {
+            Alias alias = new Alias(argument, "cast_value");
+            child = new LogicalProject<>(ImmutableList.of(alias), scan);
+            argument = alias.toSlot();
+            ruleType = scan instanceof LogicalFileScan
+                    ? 
RuleType.STORAGE_LAYER_AGGREGATE_WITH_PROJECT_FOR_FILE_SCAN
+                    : RuleType.STORAGE_LAYER_AGGREGATE_WITH_PROJECT;
+        }
+        LogicalAggregate<Plan> aggregate = new 
LogicalAggregate<>(Collections.emptyList(),
+                ImmutableList.of(new Alias(new Min(argument), "min"), new 
Alias(new Max(argument), "max")),
+                true, Optional.empty(), child);
+        CascadesContext context = 
MemoTestUtils.createCascadesContext(aggregate);
+        context.getConnectContext().getSessionVariable().enableStrictCast = 
strict;
+        RuleType selectedRuleType = ruleType;
+        Rule rule = new AggregateStrategies().buildRules().stream()
+                .filter(candidate -> candidate.getRuleType() == 
selectedRuleType).findFirst().get();
+        PlanChecker checker = 
PlanChecker.from(context).applyImplementation(rule);
+        if (expectedPushdown) {
+            checker.matches(projected
+                    ? 
logicalAggregate(logicalProject(physicalStorageLayerAggregate()))
+                    : logicalAggregate(physicalStorageLayerAggregate()));
+        } else {
+            checker.nonMatch(physicalStorageLayerAggregate());
+        }
     }
 
     @Override
diff --git 
a/regression-test/suites/external_table_p0/paimon/test_paimon_minmax_cast.groovy
 
b/regression-test/suites/external_table_p0/paimon/test_paimon_minmax_cast.groovy
new file mode 100644
index 00000000000..6e6e7ac7d36
--- /dev/null
+++ 
b/regression-test/suites/external_table_p0/paimon/test_paimon_minmax_cast.groovy
@@ -0,0 +1,105 @@
+// 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.
+
+suite("test_paimon_minmax_cast", 
"p0,external,paimon,external_docker,external_docker_paimon") {
+    String enabled = context.config.otherConfigs.get("enablePaimonTest")
+    if (enabled == null || !enabled.equalsIgnoreCase("true")) {
+        logger.info("Paimon test is disabled")
+        return
+    }
+
+    String catalogName = "test_paimon_minmax_cast"
+    String dbName = "paimon_minmax_cast_db"
+    String externalEnvIp = context.config.otherConfigs.get("externalEnvIp")
+    String minioPort = context.config.otherConfigs.get("iceberg_minio_port")
+    def originalSettings = ["enable_file_scanner_v2", "force_jni_scanner",
+            "enable_strict_cast", 
"enable_push_down_no_group_agg"].collectEntries { name ->
+        [(name): sql("show variables like '${name}'")[0][1]]
+    }
+
+    try {
+        // A single writer and bucket keep both overflowing endpoints and the 
valid interior value
+        // in one file. Separate files could let metadata aggregation retain 
the interior value.
+        spark_paimon_multi """
+            create database if not exists paimon.${dbName};
+            drop table if exists paimon.${dbName}.minmax_cast;
+            create table paimon.${dbName}.minmax_cast (value bigint, 
zero_value double, nan_value double)
+                using paimon tblproperties (
+                    'bucket'='1',
+                    'bucket-key'='value',
+                    'file.format'='parquet'
+                );
+            insert into paimon.${dbName}.minmax_cast
+                select /*+ coalesce(1) */ * from values
+                    (cast(-2147483649 as bigint), cast('0.0' as double), 
cast('0.0' as double)),
+                    (cast(0 as bigint), cast('-1e-320' as double), cast('NaN' 
as double)),
+                    (cast(2147483648 as bigint), cast(null as double), 
cast(null as double))
+                    as data(value, zero_value, nan_value) order by value;
+        """
+
+        sql """drop catalog if exists ${catalogName}"""
+        sql """create catalog ${catalogName} properties (
+            'type'='paimon',
+            'warehouse'='s3://warehouse/wh',
+            's3.endpoint'='http://${externalEnvIp}:${minioPort}',
+            's3.access_key'='admin',
+            's3.secret_key'='password',
+            's3.path.style.access'='true',
+            'meta.cache.paimon.table.ttl-second'='0'
+        )"""
+        sql """switch ${catalogName}"""
+        sql """use ${dbName}"""
+        sql "set enable_file_scanner_v2=true"
+        sql "set force_jni_scanner=false"
+        sql "set enable_strict_cast=false"
+
+        def queries = [
+            "select min(cast(value as int)) from minmax_cast",
+            "select max(cast(value as int)) from minmax_cast",
+            "select min(cast(value as int)), max(cast(value as int)) from 
minmax_cast",
+            // Numeric equality hides signed zero, and footer extrema can omit 
NaN entirely.
+            "select signbit(min(cast(zero_value as float))), " +
+                    "signbit(max(cast(zero_value as float))) from minmax_cast",
+            "select isnan(max(cast(nan_value as float))) from minmax_cast"
+        ]
+        queries.each { query ->
+            sql "set enable_push_down_no_group_agg=false"
+            def fullScanResult = sql(query)
+            sql "set enable_push_down_no_group_agg=true"
+            // Compare with row-by-row evaluation so the reference cannot 
share the metadata bug.
+            assertEquals(fullScanResult, sql(query))
+            explain {
+                sql(query)
+                contains "pushdown agg=NONE"
+                contains "paimonNativeReadSplits=1/1"
+            }
+        }
+
+        // Keep a positive control: disabling all file MIN/MAX pushdown must 
not satisfy this test.
+        explain {
+            sql "select min(value), max(value) from minmax_cast"
+            contains "pushdown agg=MINMAX"
+            contains "inputSplitNum=1"
+            // One connector split alone does not rule out a JNI fallback.
+            contains "paimonNativeReadSplits=1/1"
+        }
+    } finally {
+        originalSettings.each { name, value -> sql "set ${name}=${value}" }
+        sql "switch internal"
+        sql """drop catalog if exists ${catalogName}"""
+    }
+}
diff --git a/regression-test/suites/nereids_p0/test_minmax_cast_pushdown.groovy 
b/regression-test/suites/nereids_p0/test_minmax_cast_pushdown.groovy
new file mode 100644
index 00000000000..2483fa4a74f
--- /dev/null
+++ b/regression-test/suites/nereids_p0/test_minmax_cast_pushdown.groovy
@@ -0,0 +1,64 @@
+// 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.
+
+suite("test_minmax_cast_pushdown", "p0") {
+    def originalSettings = ["enable_strict_cast", 
"enable_push_down_no_group_agg"].collectEntries { name ->
+        [(name): sql("show variables like '${name}'")[0][1]]
+    }
+    try {
+        sql "drop table if exists test_minmax_cast_pushdown"
+        sql """
+            create table test_minmax_cast_pushdown (
+                id int,
+                value bigint
+            ) duplicate key(id)
+            distributed by hash(id) buckets 1
+            properties("replication_num"="1")
+        """
+        // One insert into one tablet keeps the valid interior value between 
overflowing zone-map endpoints.
+        sql "insert into test_minmax_cast_pushdown values (1, -2147483649), 
(2, 0), (3, 2147483648)"
+        sql "set enable_strict_cast=false"
+
+        def queries = [
+            "select min(cast(value as int)) from test_minmax_cast_pushdown",
+            "select max(cast(value as int)) from test_minmax_cast_pushdown",
+            "select min(cast(value as int)), max(cast(value as int)) from 
test_minmax_cast_pushdown",
+            "select min(cast_value), max(cast_value) from " +
+                    "(select cast(value as int) as cast_value from 
test_minmax_cast_pushdown) projected"
+        ]
+        queries.each { query ->
+            sql "set enable_push_down_no_group_agg=false"
+            def fullScanResult = sql(query)
+            sql "set enable_push_down_no_group_agg=true"
+            assertEquals(fullScanResult, sql(query))
+            explain {
+                sql(query)
+                contains "pushAggOp=NONE"
+            }
+        }
+        explain {
+            sql "select min(value), max(value) from test_minmax_cast_pushdown"
+            contains "pushAggOp=MINMAX"
+        }
+        explain {
+            sql "select min(cast(id as bigint)), max(cast(id as bigint)) from 
test_minmax_cast_pushdown"
+            contains "pushAggOp=MINMAX"
+        }
+    } finally {
+        originalSettings.each { name, value -> sql "set ${name}=${value}" }
+    }
+}


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to