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]
