This is an automated email from the ASF dual-hosted git repository. yiguolei pushed a commit to branch branch-4.2 in repository https://gitbox.apache.org/repos/asf/doris.git
commit 8b3cccfef42439ba2a89d8cdbd336283ce5df13c Author: Gabriel <[email protected]> AuthorDate: Tue Sep 29 09:42:00 2026 +0800 [fix](iceberg) Bind historical scan specs to the selected schema (#68568) Historical Iceberg predicates can still fail after #67479: Doris batch-mode detection and manifest-cache planning use table specs bound to the current schema, and Iceberg 1.11 skips rebinding when a schema-only update leaves the data snapshot unchanged. Bind partition specs by field ID to the schema selected by the Doris relation for both SDK and Doris planning. Preserve the existing scan context and reuse specs already bound to that schema. When rebinding is required, restrict it to specs referenced by the selected snapshot's data and delete manifests; later unused specs may conflict with historical column names. Tests now exercise Doris scan creation and batch-mode detection across rename/drop, schema-only/append, partitioned/unpartitioned tables, and reused field names. A new SQL regression covers snapshot and tag reads with batch mode and manifest caching enabled and disabled. Validation: - The four targeted tests fail on the unmodified branch-4.1 baseline with `Cannot find field`. - FE compilation and 148 Iceberg scan/utility unit tests passed. - The new regression suite passed against a local FE/BE and Iceberg REST/MinIO environment: 288 historical-query result assertions and 24 positive cache-access/zero-failure checks, no skipped tests. - FE Checkstyle passed with zero violations. Additional coverage verifies metadata-only partition evolution and pruning of a nonmatching manifest within the selected historical snapshot, including batch-mode file estimation. --- fe/check/checkstyle/suppressions.xml | 2 + .../datasource/iceberg/source/IcebergScanNode.java | 10 +- .../org/apache/iceberg/DorisDataTableScan.java | 71 +++++++++++ .../iceberg/source/IcebergScanNodeTest.java | 140 +++++++++++++++++---- .../test_iceberg_historical_filter_planning.groovy | 116 +++++++++++++++++ 5 files changed, 313 insertions(+), 26 deletions(-) diff --git a/fe/check/checkstyle/suppressions.xml b/fe/check/checkstyle/suppressions.xml index 8e1c40e3162..3300755eba1 100644 --- a/fe/check/checkstyle/suppressions.xml +++ b/fe/check/checkstyle/suppressions.xml @@ -69,6 +69,8 @@ under the License. <!-- ignore hudi disk map copied from hudi/common/util/collection/DiskMap.java --> <suppress files="org[\\/]apache[\\/]hudi[\\/]common[\\/]util[\\/]collection[\\/]DiskMap\.java" checks="[a-zA-Z0-9]*"/> + <!-- The scan adapter needs Iceberg package access to preserve the original scan context. --> + <suppress files="org[\\/]apache[\\/]iceberg[\\/]DorisDataTableScan\.java" checks="ImportControl"/> <!-- ignore iceberg delete file index copied from iceberg/DeleteFileIndex.java --> <suppress files="org[\\/]apache[\\/]iceberg[\\/]DeleteFileIndex\.java" checks="[a-zA-Z0-9]*"/> diff --git a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/source/IcebergScanNode.java b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/source/IcebergScanNode.java index c6e6d17af0f..75489d13b37 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/source/IcebergScanNode.java +++ b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/source/IcebergScanNode.java @@ -102,6 +102,7 @@ import org.apache.iceberg.ContentScanTask; import org.apache.iceberg.DataFile; import org.apache.iceberg.DeleteFile; import org.apache.iceberg.DeleteFileIndex; +import org.apache.iceberg.DorisDataTableScan; import org.apache.iceberg.FileContent; import org.apache.iceberg.FileFormat; import org.apache.iceberg.FileScanTask; @@ -1886,6 +1887,11 @@ public class IcebergScanNode extends FileQueryScanNode { selectedSchema, "Schema %s for Iceberg scan is null", info.getSchemaId())); } } + if (!isSystemTable) { + // Iceberg 1.11 skips spec rebinding when only the schema changed, and its snapshot + // schema can differ from the current schema deliberately projected for a branch. + scan = DorisDataTableScan.wrap(scan); + } Schema scanSchema = scan.schema(); // set filter @@ -2387,7 +2393,7 @@ public class IcebergScanNode extends FileQueryScanNode { .reduce(Expressions.alwaysTrue(), Expressions::and); // Get all partition specs by their IDs for later use - Map<Integer, PartitionSpec> specsById = icebergTable.specs(); + Map<Integer, PartitionSpec> specsById = DorisDataTableScan.specsForScan(scan); boolean caseSensitive = true; // Create residual evaluators for each partition spec @@ -2987,7 +2993,7 @@ public class IcebergScanNode extends FileQueryScanNode { try (CloseableIterator<ManifestFile> matchingManifest = IcebergUtils.getMatchingManifest( createTableScan().snapshot().dataManifests(icebergTable.io()), - icebergTable.specs(), + DorisDataTableScan.specsForScan(createTableScan()), createTableScan().filter()).iterator()) { int cnt = 0; while (matchingManifest.hasNext()) { diff --git a/fe/fe-core/src/main/java/org/apache/iceberg/DorisDataTableScan.java b/fe/fe-core/src/main/java/org/apache/iceberg/DorisDataTableScan.java new file mode 100644 index 00000000000..08c11ad8cef --- /dev/null +++ b/fe/fe-core/src/main/java/org/apache/iceberg/DorisDataTableScan.java @@ -0,0 +1,71 @@ +// 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 snapshot specs 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. + 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(scan.table().io())) { + int specId = manifest.partitionSpecId(); + specs.computeIfAbsent(specId, id -> { + PartitionSpec spec = tableSpecs.get(id); + return spec.schema() == schema ? spec : spec.toUnbound().bind(schema, true); + }); + } + } + return specs; + } + + @Override + protected Map<Integer, PartitionSpec> specs() { + return specsForScan(this); + } + + @Override + protected TableScan newRefinedScan(Table table, Schema schema, TableScanContext context) { + return new DorisDataTableScan(table, schema, context); + } +} diff --git a/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/source/IcebergScanNodeTest.java b/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/source/IcebergScanNodeTest.java index cd708220e37..a104f5982ce 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/source/IcebergScanNodeTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/source/IcebergScanNodeTest.java @@ -85,6 +85,7 @@ import org.apache.iceberg.BatchScan; import org.apache.iceberg.DataFile; import org.apache.iceberg.DataFiles; import org.apache.iceberg.DeleteFile; +import org.apache.iceberg.DorisDataTableScan; import org.apache.iceberg.FileContent; import org.apache.iceberg.FileFormat; import org.apache.iceberg.FileMetadata; @@ -233,6 +234,10 @@ public class IcebergScanNodeTest { return super.createTableScan(); } + boolean isRealBatchMode() { + return super.isBatchMode(); + } + @Override public boolean isBatchMode() { return batchMode; @@ -3108,32 +3113,75 @@ public class IcebergScanNodeTest { @Test public void testHistoricalPredicatePlansAfterColumnRename() throws Exception { - assertHistoricalPredicatePlansAfterSchemaEvolution(false); + assertHistoricalPredicatePlansAfterSchemaEvolution(false, true); } @Test public void testHistoricalPredicatePlansAfterColumnDrop() throws Exception { - assertHistoricalPredicatePlansAfterSchemaEvolution(true); + assertHistoricalPredicatePlansAfterSchemaEvolution(true, true); + } + + @Test + public void testHistoricalPredicatePlansAfterSchemaOnlyRename() throws Exception { + assertHistoricalPredicatePlansAfterSchemaEvolution(false, false); + } + + @Test + public void testHistoricalPredicatePlansAfterSchemaOnlyDrop() throws Exception { + assertHistoricalPredicatePlansAfterSchemaEvolution(true, false); } - private void assertHistoricalPredicatePlansAfterSchemaEvolution(boolean dropColumn) throws Exception { + @Test + public void testHistoricalPredicatePlansAfterMetadataOnlyPartitionEvolution() throws Exception { + for (boolean partitioned : new boolean[] {false, true}) { + assertHistoricalPredicatePlansAfterSchemaEvolution(false, false, partitioned, false, true); + } + } + + private void assertHistoricalPredicatePlansAfterSchemaEvolution(boolean dropColumn, boolean append) + throws Exception { + for (boolean partitioned : new boolean[] {false, true}) { + for (boolean reuseName : new boolean[] {false, true}) { + assertHistoricalPredicatePlansAfterSchemaEvolution(dropColumn, append, partitioned, reuseName, false); + } + } + } + + private void assertHistoricalPredicatePlansAfterSchemaEvolution( + boolean dropColumn, boolean append, boolean partitioned, boolean reuseName, boolean evolveSpec) + throws Exception { 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())); HadoopTables tables = new HadoopTables(new Configuration()); String tableLocation = temporaryFolder.getRoot().toPath() - .resolve("historical_predicate_after_" + (dropColumn ? "drop" : "rename")).toUri().toString(); + .resolve("historical_predicate_after_" + (dropColumn ? "drop" : "rename") + + "_" + partitioned + "_" + reuseName).toUri().toString(); + PartitionSpec spec = partitioned ? PartitionSpec.builderFor(historicalSchema).identity("part").build() + : PartitionSpec.unpartitioned(); Table table = tables.create( - historicalSchema, PartitionSpec.unpartitioned(), SortOrder.unsorted(), + historicalSchema, spec, SortOrder.unsorted(), ImmutableMap.of(TableProperties.FORMAT_VERSION, "2"), tableLocation); DataFile historicalDataFile = DataFiles.builder(table.spec()) .withPath(tableLocation + "/data/historical.parquet") + .withPartitionPath(partitioned ? "part=2" : "") .withFormat(FileFormat.PARQUET) .withFileSizeInBytes(10) .withRecordCount(2) .build(); table.newFastAppend().appendFile(historicalDataFile).commit(); + if (partitioned) { + // A separate manifest in the same snapshot must be rejected by partition pruning. + DataFile nonmatchingDataFile = DataFiles.builder(table.spec()) + .withPath(tableLocation + "/data/nonmatching.parquet") + .withPartitionPath("part=3") + .withFormat(FileFormat.PARQUET) + .withFileSizeInBytes(10) + .withRecordCount(1) + .build(); + table.newFastAppend().appendFile(nonmatchingDataFile).commit(); + } long historicalSnapshotId = table.currentSnapshot().snapshotId(); int historicalSchemaId = table.currentSnapshot().schemaId(); @@ -3142,25 +3190,69 @@ public class IcebergScanNodeTest { } 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 (evolveSpec) { + // This later spec is unused by the snapshot and conflicts with its old column name. + table.updateSpec().addField("x", Expressions.ref("y")).commit(); + } + 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); + Assert.assertEquals(1, tasks.size()); + Assert.assertEquals(historicalDataFile.path().toString(), tasks.get(0).file().path().toString()); + Assert.assertTrue(node.isRealBatchMode()); + if (partitioned) { + Assert.assertEquals(2, scan.snapshot().dataManifests(table.io()).size()); + // Counting the nonmatching manifest would incorrectly enable batch mode at two files. + sessionVariable.numFilesInBatchMode = 2; + setPrivateField(node, "isBatchMode", null); + Assert.assertFalse(node.isRealBatchMode()); + } + if (evolveSpec) { + Assert.assertFalse(DorisDataTableScan.specsForScan(scan).containsKey(table.spec().specId())); + } + } finally { + ConnectContext.remove(); + executor.shutdownNow(); + } } @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]
