github-actions[bot] commented on code in PR #68568:
URL: https://github.com/apache/doris/pull/68568#discussion_r4121046740


##########
fe/fe-core/src/main/java/org/apache/iceberg/DorisDataTableScan.java:
##########
@@ -0,0 +1,61 @@
+// 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.iceberg;
+
+import java.util.HashMap;
+import java.util.Map;
+
+/** Keeps partition pruning aligned with the schema selected by the Doris 
relation. */
+public class DorisDataTableScan extends DataTableScan {
+    private DorisDataTableScan(Table table, Schema schema, TableScanContext 
context) {
+        super(table, schema, context);
+    }
+
+    public static TableScan wrap(TableScan scan) {
+        if (!(scan instanceof DataTableScan)) {
+            return scan;
+        }
+        DataTableScan dataScan = (DataTableScan) scan;
+        return new DorisDataTableScan(dataScan.table(), 
dataScan.tableSchema(), dataScan.context());
+    }
+
+    /** Rebinds each spec by field ID to the full schema selected for the 
scan. */
+    public static Map<Integer, PartitionSpec> specsForScan(TableScan scan) {
+        Map<Integer, PartitionSpec> tableSpecs = scan.table().specs();
+        Schema schema = scan.schema();
+        if (tableSpecs.values().stream().allMatch(spec -> spec.schema() == 
schema)) {
+            return tableSpecs;
+        }
+        Map<Integer, PartitionSpec> specs = new HashMap<>();
+        // A schema-only update keeps the snapshot ID unchanged. Snapshot IDs 
therefore cannot
+        // determine whether table specs still use the schema selected by a 
historical query.
+        tableSpecs.forEach((id, spec) ->

Review Comment:
   A valid metadata-only sequence can make this loop fail every historical read 
of the still-current snapshot: create S1 with `a` (id 1) and `b` (id 2), rename 
`a`, then add an identity partition field named `a` sourced from `b` without 
appending. `FOR VERSION AS OF S1` selects the old schema, but `table().specs()` 
also contains the newer spec. Iceberg's `PartitionSpec.Builder.add` rejects 
rebinding that spec because old `a` is id 1 while its source is id 2 (`Cannot 
create identity partition sourced from different field in schema: a`); 
`bind(..., true)` does not skip that check. The selected snapshot has no 
manifest using the newer spec. Please scope rebinding to specs the snapshot can 
use, or retain unrelated specs without rebinding, and cover this metadata-only 
case.



##########
regression-test/suites/external_table_p0/iceberg/test_iceberg_historical_filter_planning.groovy:
##########
@@ -0,0 +1,108 @@
+// 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_iceberg_historical_filter_planning", 
"p0,external,doris,external_docker,external_docker_doris") {
+    if 
(!"true".equalsIgnoreCase(context.config.otherConfigs.get("enableIcebergTest")))
 {
+        return
+    }
+
+    String restPort = context.config.otherConfigs.get("iceberg_rest_uri_port")
+    String minioPort = context.config.otherConfigs.get("iceberg_minio_port")
+    String externalEnvIp = context.config.otherConfigs.get("externalEnvIp")
+    String baseCatalog = "iceberg_historical_filter_planning"
+    String cacheCatalog = "iceberg_historical_filter_planning_cache"
+    String dbName = "historical_filter_planning_db"
+    def catalogs = [baseCatalog, cacheCatalog]
+    def oldBatchMode = sql("show variables like 
'enable_external_table_batch_mode'")[0][1]
+    def oldBatchSize = sql("show variables like 
'num_files_in_batch_mode'")[0][1]
+
+    try {
+        catalogs.each { catalog ->
+            sql "drop catalog if exists ${catalog}"
+            sql """create catalog ${catalog} properties (
+                'type' = 'iceberg',
+                'iceberg.catalog.type' = 'rest',
+                'uri' = 'http://${externalEnvIp}:${restPort}',
+                's3.access_key' = 'admin',
+                's3.secret_key' = 'password',
+                's3.endpoint' = 'http://${externalEnvIp}:${minioPort}',
+                's3.region' = 'us-east-1',
+                'meta.cache.iceberg.manifest.enable' = '${catalog == 
cacheCatalog}'
+            )"""
+        }
+        sql "create database if not exists ${baseCatalog}.${dbName}"
+        sql "set num_files_in_batch_mode = 1"
+
+        [false, true].each { partitioned ->
+            ["rename", "drop"].each { change ->
+                String tableName = "${change}_${partitioned ? 'partitioned' : 
'unpartitioned'}"
+                String table = "${baseCatalog}.${dbName}.${tableName}"
+                sql "drop table if exists ${table}"
+                sql """create table ${table} (id int, x int, part int)
+                    ${partitioned ? 'partition by list (part) ()' : ''}
+                    properties ('format-version' = '2')"""
+                sql "insert into ${table} values (1, 1, 1), (2, 1, 2), (3, 2, 
2)"
+                def snapshotId = sql("select snapshot_id from 
${table}\$snapshots")[0][0]
+                sql "alter table ${table} create tag before_change"
+                if (change == "rename") {
+                    sql "alter table ${table} rename column x renamed_x"
+                } else {
+                    sql "alter table ${table} drop column x"
+                }
+
+                def assertHistoricalFilters = {
+                    catalogs.each { catalog ->
+                        sql "refresh table ${catalog}.${dbName}.${tableName}"

Review Comment:
   Please assert that the cache-enabled, batch-disabled case actually uses 
manifest-cache planning. `planFileScanTaskWithoutReuse` catches any failure 
from `loadFileScanTasksWithManifestCache` (including the new spec rebinding) 
and falls back to the SDK scan, so every row assertion here can pass while this 
PR's cache fix is broken. As in `test_iceberg_manifest_cache.groovy`, check 
`EXPLAIN VERBOSE` for a nonzero cache miss or hit **and** `failures=0`; mere 
presence of the cache line can also mean all counters are zero.



##########
fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/source/IcebergScanNodeTest.java:
##########
@@ -3136,25 +3164,55 @@ private void 
assertHistoricalPredicatePlansAfterSchemaEvolution(boolean dropColu
         } else {
             table.updateSchema().renameColumn("x", "renamed_x").commit();
         }
-        DataFile currentDataFile = DataFiles.builder(table.spec())
-                .withPath(tableLocation + "/data/current.parquet")
-                .withFormat(FileFormat.PARQUET)
-                .withFileSizeInBytes(10)
-                .withRecordCount(1)
-                .build();
-        table.newFastAppend().appendFile(currentDataFile).commit();
-
-        // Historical filters must be resolved with the snapshot schema after 
later schema evolution.
-        TableScan scan = table.newScan()
-                .useSnapshot(historicalSnapshotId)
-                .project(table.schemas().get(historicalSchemaId));
-        BinaryPredicate conjunct = new 
BinaryPredicate(BinaryPredicate.Operator.EQ,
-                new SlotRef(new TableName(), "x"), new IntLiteral(1, 
Type.INT));
-        org.apache.iceberg.expressions.Expression predicate =
-                IcebergUtils.convertToIcebergExpr(conjunct, scan.schema());
-        Assert.assertNotNull(predicate);
-        scan = scan.filter(predicate);
-        Assert.assertEquals(1, materializeTasks(scan).size());
+        if (reuseName) {
+            table.updateSchema().addColumn("x", 
Types.IntegerType.get()).commit();
+        }
+        if (append) {
+            DataFile currentDataFile = DataFiles.builder(table.spec())
+                    .withPath(tableLocation + "/data/current.parquet")
+                    .withPartitionPath(partitioned ? "part=3" : "")
+                    .withFormat(FileFormat.PARQUET)
+                    .withFileSizeInBytes(10)
+                    .withRecordCount(1)
+                    .build();
+            table.newFastAppend().appendFile(currentDataFile).commit();
+        }
+        table = tables.load(tableLocation);
+        Assert.assertEquals(!append, historicalSnapshotId == 
table.currentSnapshot().snapshotId());
+
+        SessionVariable sessionVariable = new SessionVariable();
+        sessionVariable.enableExternalTableBatchMode = true;
+        sessionVariable.numFilesInBatchMode = 1;
+        TestIcebergScanNode node = Mockito.spy(new 
TestIcebergScanNode(sessionVariable));
+        IcebergSource source = Mockito.mock(IcebergSource.class);
+        IcebergExternalCatalog catalog = 
Mockito.mock(IcebergExternalCatalog.class);
+        ThreadPoolExecutor executor = (ThreadPoolExecutor) 
Executors.newFixedThreadPool(1);
+        Mockito.when(source.getCatalog()).thenReturn(catalog);
+        Mockito.when(catalog.getThreadPoolWithPreAuth()).thenReturn(executor);
+        setIcebergSource(node, source);
+        setIcebergTable(node, table);
+        setPreExecutionAuthenticator(node, new ExecutionAuthenticator() {});
+        Mockito.doReturn(new IcebergTableQueryInfo(historicalSnapshotId, null, 
historicalSchemaId))
+                .when(node).getSpecifiedSnapshot();
+        node.addConjunct(new BinaryPredicate(BinaryPredicate.Operator.EQ,
+                new SlotRef(new TableName(), "x"), new IntLiteral(1, 
Type.INT)));
+        node.addConjunct(new BinaryPredicate(BinaryPredicate.Operator.EQ,
+                new SlotRef(new TableName(), "part"), new IntLiteral(2, 
Type.INT)));
+        ConnectContext context = new ConnectContext();
+        context.setStatementContext(new StatementContext());
+        context.setThreadLocalInfo();
+        try {
+            // Exercise both Doris planning paths; an SDK-only scan misses 
batch-mode manifest filtering.
+            TableScan scan = node.createRealTableScan();
+            node.setTableScan(scan);
+            List<FileScanTask> tasks = materializeTasks(scan);

Review Comment:
   This task assertion cannot show that the historical partition predicate 
prunes files: S1 contains only `historicalDataFile` at `part=2`; the `part=3` 
file is appended after S1 and is excluded by snapshot selection. 
`scan.planFiles()` still returns exactly this task if partition pruning is 
disabled, and the batch assertion only checks a one-file threshold. Please put 
a nonmatching partition file in S1 and assert that only the `part=2` file is 
planned (and, for the batch estimate, make matching versus nonmatching 
manifests distinguishable).



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


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

Reply via email to