DorisDataTableScan.java
// 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);
}
}