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 c6a4f9512b8 [fix](iceberg) Scope historical scan specs to selected 
snapshot (#68582)
c6a4f9512b8 is described below

commit c6a4f9512b8de3996c6484358847abf98833ed82
Author: Gabriel <[email protected]>
AuthorDate: Tue Sep 29 09:44:26 2026 +0800

    [fix](iceberg) Scope historical scan specs to selected snapshot (#68582)
    
    Ports #68568 to master, adapted to the Iceberg connector architecture.
    
    Master already uses `SchemaAwareDataTableScan` for historical schemas,
    but still rebinds every table partition spec. Renaming a column and then
    adding an unused partition field with that old name can make a
    historical query fail with `Cannot create identity partition sourced
    from different field in schema` even though the selected snapshot never
    uses the new spec.
    
    Bind only specs referenced by the selected snapshot's data and delete
    manifests to the selected scan schema. Use the same helper for SDK file
    planning, streaming file estimates, and manifest-cache planning. Retain
    the existing fast path when the scan uses the current schema.
    
    Adapt the historical scan tests to exercise partitioned/unpartitioned
    tables, rename/drop, schema-only/append, reused names,
    synchronous/streaming planning, and cache on/off. Assert that a
    nonmatching manifest in the selected snapshot is pruned and that cache
    planning succeeds without fallback. Also port the SQL snapshot/tag
    regression suite from #68568 unchanged.
    
    Validation:
    - The new metadata-only partition-evolution test reproduces the
    identity-partition error on unmodified master.
    - All 181 `IcebergScanPlanProviderTest` tests passed with the Maven
    build cache disabled.
    - Connector reactor compilation, Checkstyle, and the connector import
    gate passed.
    - The SQL regression suite is included; end-to-end SQL execution was not
    repeated on master in this port.
---
 .../connector/iceberg/IcebergScanPlanProvider.java |   4 +-
 .../apache/iceberg/SchemaAwareDataTableScan.java   |  21 +++-
 .../iceberg/IcebergScanPlanProviderTest.java       | 121 ++++++++++++++-------
 .../test_iceberg_historical_filter_planning.groovy | 116 ++++++++++++++++++++
 4 files changed, 214 insertions(+), 48 deletions(-)

diff --git 
a/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergScanPlanProvider.java
 
b/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergScanPlanProvider.java
index 8774b6a5d7a..09b77b83bb4 100644
--- 
a/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergScanPlanProvider.java
+++ 
b/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergScanPlanProvider.java
@@ -530,7 +530,7 @@ public class IcebergScanPlanProvider implements 
ConnectorScanPlanProvider {
         long fileCount = 0;
         try (CloseableIterable<ManifestFile> matching = getMatchingManifest(
                 snapshot.dataManifests(table.io()),
-                SchemaAwareDataTableScan.specsFor(table, scan.schema()), 
scan.filter())) {
+                SchemaAwareDataTableScan.specsFor(scan), scan.filter())) {
             for (ManifestFile manifest : matching) {
                 // Manifest metadata counts (cheap — no per-file read). Null 
guard for ancient manifests that
                 // omit the counts (legacy summed them unguarded; 0 is the 
safe under-count, never over-streams).
@@ -2811,7 +2811,7 @@ public class IcebergScanPlanProvider implements 
ConnectorScanPlanProvider {
         }
         Schema scanSchema = scan.schema();
         Expression filterExpr = combineFilter(filter, scanSchema, session);
-        Map<Integer, PartitionSpec> specsById = 
SchemaAwareDataTableScan.specsFor(table, scanSchema);
+        Map<Integer, PartitionSpec> specsById = 
SchemaAwareDataTableScan.specsFor(scan);
         boolean caseSensitive = true;
 
         Map<Integer, ResidualEvaluator> residualEvaluators = new HashMap<>();
diff --git 
a/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/iceberg/SchemaAwareDataTableScan.java
 
b/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/iceberg/SchemaAwareDataTableScan.java
index dbef7b6c44c..31ec1e3ea68 100644
--- 
a/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/iceberg/SchemaAwareDataTableScan.java
+++ 
b/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/iceberg/SchemaAwareDataTableScan.java
@@ -32,14 +32,25 @@ public final class SchemaAwareDataTableScan extends 
DataTableScan {
         return new SchemaAwareDataTableScan(table, table.schema(), 
TableScanContext.empty());
     }
 
-    /** Returns every table spec rebound to {@code schema}. */
-    public static Map<Integer, PartitionSpec> specsFor(Table table, Schema 
schema) {
+    /** Returns the selected snapshot's specs bound to the scan schema. */
+    public static Map<Integer, PartitionSpec> specsFor(TableScan scan) {
+        Table table = scan.table();
+        Schema schema = scan.schema();
+        Map<Integer, PartitionSpec> tableSpecs = table.specs();
         if (schema.sameSchema(table.schema())) {
-            return table.specs();
+            return tableSpecs;
         }
 
         Map<Integer, PartitionSpec> specs = new LinkedHashMap<>();
-        table.specs().forEach((id, spec) -> specs.put(id, 
spec.toUnbound().bind(schema, true)));
+        Snapshot snapshot = scan.snapshot();
+        if (snapshot != null) {
+            // Later, unused specs can have partition names that conflict with 
historical columns.
+            // Bind only specs referenced by this snapshot, including those 
needed for delete files.
+            for (ManifestFile manifest : snapshot.allManifests(table.io())) {
+                specs.computeIfAbsent(manifest.partitionSpecId(),
+                        id -> tableSpecs.get(id).toUnbound().bind(schema, 
true));
+            }
+        }
         return Collections.unmodifiableMap(specs);
     }
 
@@ -47,7 +58,7 @@ public final class SchemaAwareDataTableScan extends 
DataTableScan {
     protected Map<Integer, PartitionSpec> specs() {
         // A metadata-only schema commit preserves the current snapshot ID, so 
schema identity—not snapshot
         // identity—must decide whether historical partition specs need 
rebinding.
-        return specsFor(table(), tableSchema());
+        return specsFor(this);
     }
 
     @Override
diff --git 
a/fe/fe-connector/fe-connector-iceberg/src/test/java/org/apache/doris/connector/iceberg/IcebergScanPlanProviderTest.java
 
b/fe/fe-connector/fe-connector-iceberg/src/test/java/org/apache/doris/connector/iceberg/IcebergScanPlanProviderTest.java
index 42530caf9f1..c9592884791 100644
--- 
a/fe/fe-connector/fe-connector-iceberg/src/test/java/org/apache/doris/connector/iceberg/IcebergScanPlanProviderTest.java
+++ 
b/fe/fe-connector/fe-connector-iceberg/src/test/java/org/apache/doris/connector/iceberg/IcebergScanPlanProviderTest.java
@@ -22,6 +22,7 @@ import org.apache.doris.connector.spi.ConnectorStatementScope;
 import org.apache.doris.connector.spi.ConnectorType;
 import org.apache.doris.connector.spi.DorisConnectorException;
 import org.apache.doris.connector.spi.handle.ConnectorColumnHandle;
+import org.apache.doris.connector.spi.pushdown.ConnectorAnd;
 import org.apache.doris.connector.spi.pushdown.ConnectorColumnRef;
 import org.apache.doris.connector.spi.pushdown.ConnectorComparison;
 import org.apache.doris.connector.spi.pushdown.ConnectorExpression;
@@ -65,6 +66,7 @@ import org.apache.iceberg.TableProperties;
 import org.apache.iceberg.TableScan;
 import org.apache.iceberg.catalog.Namespace;
 import org.apache.iceberg.catalog.TableIdentifier;
+import org.apache.iceberg.expressions.Expressions;
 import org.apache.iceberg.expressions.Literal;
 import org.apache.iceberg.hadoop.HadoopTables;
 import org.apache.iceberg.inmemory.InMemoryCatalog;
@@ -2131,57 +2133,94 @@ public class IcebergScanPlanProviderTest {
     }
 
     @Test
-    public void planScanHistoricalPredicateSurvivesColumnRename() {
-        assertHistoricalPredicatePlansAfterSchemaEvolution(false);
+    public void planScanHistoricalPredicateSurvivesColumnRename() throws 
IOException {
+        assertHistoricalPredicatePlansAfterSchemaEvolution(false, false);
     }
 
     @Test
-    public void planScanHistoricalPredicateSurvivesColumnDrop() {
-        assertHistoricalPredicatePlansAfterSchemaEvolution(true);
+    public void planScanHistoricalPredicateSurvivesColumnDrop() throws 
IOException {
+        assertHistoricalPredicatePlansAfterSchemaEvolution(true, false);
     }
 
-    private void assertHistoricalPredicatePlansAfterSchemaEvolution(boolean 
dropColumn) {
-        Schema historicalSchema = new Schema(
-                Types.NestedField.optional(1, "x", Types.IntegerType.get()),
-                Types.NestedField.optional(2, "y", Types.IntegerType.get()),
-                Types.NestedField.optional(3, "part", 
Types.IntegerType.get()));
-        Table table = createTable(
-                "historical_predicate_after_" + (dropColumn ? "drop" : 
"rename"),
-                historicalSchema, PartitionSpec.unpartitioned(),
-                Collections.singletonMap(TableProperties.FORMAT_VERSION, "2"));
-        table.newFastAppend()
-                .appendFile(dataFile(table.spec(), 
"s3://b/db/historical.parquet", 1024, null, null))
-                .commit();
-        long historicalSnapshotId = table.currentSnapshot().snapshotId();
-        int historicalSchemaId = table.currentSnapshot().schemaId();
+    @Test
+    public void historicalPlanningIgnoresUnusedConflictingPartitionSpecs() 
throws IOException {
+        assertHistoricalPredicatePlansAfterSchemaEvolution(false, true);
+    }
 
-        if (dropColumn) {
-            table.updateSchema().deleteColumn("x").commit();
-        } else {
-            table.updateSchema().renameColumn("x", "renamed_x").commit();
+    private void assertHistoricalPredicatePlansAfterSchemaEvolution(boolean 
dropColumn, boolean evolveSpec)
+            throws IOException {
+        for (boolean partitioned : new boolean[] {false, true}) {
+            Schema historicalSchema = new Schema(
+                    Types.NestedField.optional(1, "x", 
Types.IntegerType.get()),
+                    Types.NestedField.optional(2, "y", 
Types.IntegerType.get()),
+                    Types.NestedField.optional(3, "part", 
Types.IntegerType.get()));
+            PartitionSpec spec = partitioned ? 
PartitionSpec.builderFor(historicalSchema).identity("part").build()
+                    : PartitionSpec.unpartitioned();
+            Table table = createTable("historical_predicate_" + dropColumn + 
"_" + partitioned,
+                    historicalSchema, spec, 
Collections.singletonMap(TableProperties.FORMAT_VERSION, "2"));
+            table.newFastAppend().appendFile(dataFile(table.spec(), 
"s3://b/db/historical.parquet",
+                    1024, null, partitioned ? "part=2" : null)).commit();
+            if (partitioned) {
+                // Both manifests belong to the selected snapshot; snapshot 
selection alone cannot prune this one.
+                table.newFastAppend().appendFile(dataFile(table.spec(), 
"s3://b/db/nonmatching.parquet",
+                        1024, null, "part=3")).commit();
+            }
+            long snapshotId = table.currentSnapshot().snapshotId();
+            int schemaId = table.currentSnapshot().schemaId();
+            if (dropColumn) {
+                table.updateSchema().deleteColumn("x").commit();
+            } else {
+                table.updateSchema().renameColumn("x", "renamed_x").commit();
+            }
+            if (evolveSpec) {
+                // This newer spec is unused by the snapshot and conflicts 
with its old column name.
+                table.updateSpec().addField("x", 
Expressions.ref("y")).commit();
+            }
+            Assertions.assertEquals(snapshotId, 
table.currentSnapshot().snapshotId());
+            assertHistoricalPredicatePlans(table, snapshotId, schemaId, 
partitioned);
+            if (evolveSpec) {
+                continue;
+            }
+            table.newFastAppend().appendFile(dataFile(table.spec(), 
"s3://b/db/current.parquet",
+                    1024, null, partitioned ? "part=3" : null)).commit();
+            assertHistoricalPredicatePlans(table, snapshotId, schemaId, 
partitioned);
+            // A reused name has a different field ID and must not change 
historical predicate binding.
+            table.updateSchema().addColumn("x", 
Types.IntegerType.get()).commit();
+            assertHistoricalPredicatePlans(table, snapshotId, schemaId, 
partitioned);
         }
-
-        assertHistoricalPredicatePlans(table, historicalSnapshotId, 
historicalSchemaId);
-
-        table.newFastAppend()
-                .appendFile(dataFile(table.spec(), 
"s3://b/db/current.parquet", 1024, null, null))
-                .commit();
-
-        assertHistoricalPredicatePlans(table, historicalSnapshotId, 
historicalSchemaId);
     }
 
     private static void assertHistoricalPredicatePlans(
-            Table table, long historicalSnapshotId, int historicalSchemaId) {
-        IcebergTableHandle historicalHandle = new IcebergTableHandle("db1", 
"t1")
-                .withSnapshot(historicalSnapshotId, null, historicalSchemaId);
-        List<ConnectorScanRange> ranges = providerOver(table).planScan(
-                emptySession(), ConnectorScanRequest.builder(historicalHandle, 
Collections.emptyList())
-                        .filter(Optional.of(eqInt("x", 1)))
-                        .build());
-
-        // Historical predicates must remain bound to the snapshot schema 
after later schema evolution.
-        Assertions.assertEquals(1, ranges.size());
-        
Assertions.assertTrue(ranges.get(0).getPath().get().endsWith("historical.parquet"));
+            Table table, long snapshotId, int schemaId, boolean partitioned) 
throws IOException {
+        IcebergTableHandle handle = new IcebergTableHandle("db1", 
"t1").withSnapshot(snapshotId, null, schemaId);
+        Optional<ConnectorExpression> filter = Optional.of(new ConnectorAnd(
+                Arrays.asList(eqInt("x", 1), eqInt("part", 2))));
+        for (boolean cacheEnabled : new boolean[] {false, true}) {
+            IcebergManifestCache cache = new IcebergManifestCache();
+            IcebergScanPlanProvider provider = cacheEnabled
+                    ? manifestProvider(manifestCacheProps(), table, cache) : 
providerOver(table);
+            ConnectorSession session = emptySession();
+            List<ConnectorScanRange> ranges = provider.planScan(session,
+                    ConnectorScanRequest.builder(handle, 
Collections.emptyList()).filter(filter).build());
+            Assertions.assertEquals(1, ranges.size());
+            
Assertions.assertTrue(ranges.get(0).getPath().get().endsWith("historical.parquet"));
+            if (cacheEnabled) {
+                // Correct rows alone do not prove the cache path succeeded 
instead of falling back to the SDK.
+                long[] stats = cache.takeStats(session.getQueryId());
+                Assertions.assertTrue(stats[0] + stats[1] > 0);
+                Assertions.assertEquals(0L, stats[2]);
+            }
+            Assertions.assertEquals(1L, 
provider.streamingSplitEstimate(batchSession(1, true), handle, filter, false));
+            if (partitioned) {
+                Assertions.assertEquals(2, 
table.snapshot(snapshotId).dataManifests(table.io()).size());
+                Assertions.assertEquals(-1L,
+                        provider.streamingSplitEstimate(batchSession(2, true), 
handle, filter, false),
+                        "the nonmatching manifest must not enable streaming at 
a two-file threshold");
+            }
+            List<ConnectorScanRange> streamed = drain(provider.streamSplits(
+                    batchSession(1, true), handle, Collections.emptyList(), 
filter, -1L));
+            Assertions.assertEquals(sortedPaths(ranges), 
sortedPaths(streamed));
+        }
     }
 
     @Test
diff --git 
a/regression-test/suites/external_table_p0/iceberg/test_iceberg_historical_filter_planning.groovy
 
b/regression-test/suites/external_table_p0/iceberg/test_iceberg_historical_filter_planning.groovy
new file mode 100644
index 00000000000..c569b2e3b8f
--- /dev/null
+++ 
b/regression-test/suites/external_table_p0/iceberg/test_iceberg_historical_filter_planning.groovy
@@ -0,0 +1,116 @@
+// 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}"
+                        [false, true].each { batchMode ->
+                            sql "set enable_external_table_batch_mode = 
${batchMode}"
+                            [" for version as of ${snapshotId}", 
"@tag(before_change)"].each { ref ->
+                                String historical = 
"${catalog}.${dbName}.${tableName}${ref}"
+                                assertEquals([[1, 1, 1], [2, 1, 2], [3, 2, 2]],
+                                        sql("select id, x, part from 
${historical} order by id"))
+                                assertEquals([[1, 1, 1], [2, 1, 2]],
+                                        sql("select id, x, part from 
${historical} where x = 1 order by id"))
+                                String filteredQuery = "select id, x, part 
from ${historical} where x = 1 and part = 2 order by id"
+                                assertEquals([[2, 1, 2]], sql(filteredQuery))
+                                if (catalog == cacheCatalog && !batchMode) {
+                                    // Row results alone also pass if cache 
planning silently falls back to the SDK.
+                                    String plan = sql("explain verbose 
${filteredQuery}").collect { it[0] }.join("\n")
+                                    def cacheStats = plan =~ /manifest cache: 
hits=(\d+), misses=(\d+), failures=(\d+)/
+                                    assertTrue(cacheStats.find(), plan)
+                                    assertTrue(cacheStats.group(1).toLong() + 
cacheStats.group(2).toLong() > 0, plan)
+                                    assertEquals(0L, 
cacheStats.group(3).toLong())
+                                }
+                            }
+                        }
+                    }
+                }
+
+                // Schema-only changes retain the data snapshot ID; SDK-only 
time-travel tests
+                // that append immediately after DDL miss this case and Doris 
batch-mode pruning.
+                assertEquals(snapshotId, sql("select snapshot_id from 
${table}\$snapshots")[0][0])
+                assertHistoricalFilters()
+
+                if (change == "rename") {
+                    sql "insert into ${table} values (4, 1, 3)"
+                } else {
+                    sql "insert into ${table} values (4, 3)"
+                }
+                assertHistoricalFilters()
+
+                // The same name with a new field ID must not change 
historical predicate binding.
+                sql "alter table ${table} add column x int"
+                assertHistoricalFilters()
+            }
+        }
+    } finally {
+        sql "set enable_external_table_batch_mode = ${oldBatchMode}"
+        sql "set num_files_in_batch_mode = ${oldBatchSize}"
+        sql "drop database if exists ${baseCatalog}.${dbName} force"
+        catalogs.reverseEach { catalog -> sql "drop catalog if exists 
${catalog}" }
+    }
+}


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

Reply via email to