Gabriel39 commented on code in PR #68568: URL: https://github.com/apache/doris/pull/68568#discussion_r4122197484
########## 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: Fixed in d58c0e1027e. When rebinding is needed, the spec map now includes only IDs referenced by the selected snapshot's data and delete manifests. Added a metadata-only rename followed by a conflicting identity partition-field addition, for both initially partitioned and unpartitioned tables. The new test fails on the previous commit with the reported identity-partition error and passes after this change. All 148 related unit tests and FE Checkstyle pass. ########## 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: Fixed in d58c0e1027e. Each cache-enabled, batch-disabled historical filtered query now checks EXPLAIN VERBOSE for hits + misses > 0 and failures = 0. The full SQL suite passed with 288 row-result assertions and 24 cache-path checks, with no skipped tests. ########## 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: Fixed in d58c0e1027e. The selected historical snapshot now contains a part=3 file in a separate manifest alongside the matching part=2 file. Tests assert that only the matching file is planned, that both manifests exist in the selected snapshot, and that batch mode is enabled at a one-file threshold but disabled at a two-file threshold. These assertions pass across the historical rename/drop scenarios. -- 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]
