This is an automated email from the ASF dual-hosted git repository.
yiguolei pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/branch-4.1 by this push:
new 2c84c7f56f1 [fix](iceberg) Fix historical scans after schema evolution
(#67479)
2c84c7f56f1 is described below
commit 2c84c7f56f18eb75bd94843b7a4400bd1267f976
Author: Gabriel <[email protected]>
AuthorDate: Wed Sep 9 09:32:16 2026 +0800
[fix](iceberg) Fix historical scans after schema evolution (#67479)
### What problem does this PR solve?
Issue Number: DORIS-28397
Related PR: https://github.com/apache/doris-shade/pull/63
Problem Summary:
Iceberg time-travel planning can fail when a predicate references a
column that was renamed or dropped after the selected snapshot. Iceberg
1.11.0 includes the upstream fix that resolves historical partition
specs with the correct schema. This PR upgrades the direct and shaded
Iceberg dependencies, adapts the DLF table operations constructor, and
synchronizes Doris's public `DeleteFileIndex` fork with Iceberg 1.11.0
while retaining Java 8 source compatibility.
The shaded catalog dependency temporarily uses `3.1.3-ICEBERG-SNAPSHOT`
until the corresponding Doris Shade change is released.
### Release note
Fix Iceberg time-travel queries after column rename or drop.
### Check List (For Author)
- Test
- [x] Regression test
- [x] Unit Test
- [ ] Manual test (add detailed scripts or steps below)
- [ ] No need to test or manual test. Explain why:
- [ ] This is a refactor/code format and no logic has been changed.
- [ ] Previous test can cover this change.
- [ ] No code files have been changed.
- [ ] Other reason
Validation performed:
- Full FE package build with tests skipped
- FE Checkstyle with zero violations
- `IcebergScanNodeTest`: 93 tests passed, including new rename and drop
cases
- DLF Iceberg unit tests: 7 tests passed
- Behavior changed:
- [ ] No.
- [x] Yes. Historical Iceberg predicates continue to plan after later
schema evolution.
- Does this need documentation?
- [x] No.
- [ ] Yes.
### Check List (For Reviewer who merge this PR)
- [ ] Confirm the release note
- [ ] Confirm test cases
- [ ] Confirm document
- [ ] Add branch pick label
---------
Co-authored-by: Calvin Kirs <[email protected]>
---
.../datasource/iceberg/dlf/DLFTableOperations.java | 3 +-
.../java/org/apache/iceberg/DeleteFileIndex.java | 63 +++++++++++++++++-----
.../iceberg/IcebergExternalMetaCacheTest.java | 10 ++--
.../iceberg/source/IcebergScanNodeTest.java | 57 ++++++++++++++++++++
fe/pom.xml | 4 +-
.../iceberg_schema_change_ddl_with_branch.out | 5 ++
.../iceberg/test_iceberg_sys_table.out | 2 +
.../iceberg_schema_change_ddl_with_branch.groovy | 7 +--
8 files changed, 128 insertions(+), 23 deletions(-)
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/dlf/DLFTableOperations.java
b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/dlf/DLFTableOperations.java
index 2aab8e754ca..33bf0deccfe 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/dlf/DLFTableOperations.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/dlf/DLFTableOperations.java
@@ -32,6 +32,7 @@ public class DLFTableOperations extends HiveTableOperations {
String catalogName,
String database,
String table) {
- super(conf, metaClients, fileIO, catalogName, database, table);
+ // DLF does not configure an Iceberg KMS client; null preserves the
existing unencrypted behavior.
+ super(conf, metaClients, fileIO, null, catalogName, database, table);
}
}
diff --git a/fe/fe-core/src/main/java/org/apache/iceberg/DeleteFileIndex.java
b/fe/fe-core/src/main/java/org/apache/iceberg/DeleteFileIndex.java
index 36cf36b556d..630f19c81a4 100644
--- a/fe/fe-core/src/main/java/org/apache/iceberg/DeleteFileIndex.java
+++ b/fe/fe-core/src/main/java/org/apache/iceberg/DeleteFileIndex.java
@@ -32,6 +32,9 @@ import java.util.Queue;
import java.util.Set;
import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.ExecutorService;
+import java.util.function.Function;
+import java.util.function.Supplier;
+import java.util.stream.Collectors;
import org.apache.iceberg.exceptions.RuntimeIOException;
import org.apache.iceberg.exceptions.ValidationException;
import org.apache.iceberg.expressions.Expression;
@@ -64,7 +67,7 @@ import org.apache.iceberg.util.Tasks;
* DataFile)} or {@link #forEntry(ManifestEntry)} to get the delete files to
apply to a given data
* file.
*
- * Copied from
https://github.com/apache/iceberg/blob/apache-iceberg-1.9.1/core/src/main/java/org/apache/iceberg/DeleteFileIndex.java
+ * Copied from
https://github.com/apache/iceberg/blob/apache-iceberg-1.11.0/core/src/main/java/org/apache/iceberg/DeleteFileIndex.java
* Change DeleteFileIndex and some methods to public.
*/
public class DeleteFileIndex {
@@ -371,6 +374,7 @@ public class DeleteFileIndex {
private final Iterable<DeleteFile> deleteFiles;
private long minSequenceNumber = 0L;
private Map<Integer, PartitionSpec> specsById = null;
+ private Map<Integer, Schema> schemasById = null;
private Expression dataFilter = Expressions.alwaysTrue();
private Expression partitionFilter = Expressions.alwaysTrue();
private PartitionSet partitionSet = null;
@@ -396,6 +400,11 @@ public class DeleteFileIndex {
return this;
}
+ Builder schemasById(Map<Integer, Schema> newSchemasById) {
+ this.schemasById = newSchemasById;
+ return this;
+ }
+
public Builder specsById(Map<Integer, PartitionSpec> newSpecsById) {
this.specsById = newSpecsById;
return this;
@@ -459,8 +468,14 @@ public class DeleteFileIndex {
try (CloseableIterable<ManifestEntry<DeleteFile>> reader =
deleteFile) {
for (ManifestEntry<DeleteFile> entry : reader) {
if (entry.dataSequenceNumber() > minSequenceNumber) {
+ DeleteFile file = entry.file();
+ // keep minimum stats to avoid memory pressure
+ Set<Integer> columns =
+ file.content() == FileContent.POSITION_DELETES
+ ?
Collections.singleton(MetadataColumns.DELETE_FILE_PATH.fieldId())
+ : Sets.newHashSet(file.equalityFieldIds());
// copy with stats for better filtering against data
file stats
- files.add(entry.file().copy());
+ files.add(ContentFileUtil.copy(file, true, columns));
}
}
} catch (IOException e) {
@@ -470,10 +485,21 @@ public class DeleteFileIndex {
return files;
}
+ private Collection<Schema> schemas() {
+ if (schemasById != null) {
+ return schemasById.values();
+ } else {
+ return
specsById.values().stream().map(PartitionSpec::schema).collect(Collectors.toList());
+ }
+ }
+
public DeleteFileIndex build() {
+ // Equality deletes may reference fields from historical schemas, so
index every known field ID.
+ Map<Integer, Types.NestedField> fieldsById =
Schema.indexFields(schemas());
+ Function<Integer, Types.NestedField> fieldLookup = fieldsById::get;
Iterable<DeleteFile> files = deleteFiles != null ? filterDeleteFiles() :
loadDeleteFiles();
- EqualityDeletes globalDeletes = new EqualityDeletes();
+ EqualityDeletes globalDeletes = new EqualityDeletes(fieldLookup);
PartitionMap<EqualityDeletes> eqDeletesByPartition =
PartitionMap.create(specsById);
PartitionMap<PositionDeletes> posDeletesByPartition =
PartitionMap.create(specsById);
Map<String, PositionDeletes> posDeletesByPath = Maps.newHashMap();
@@ -489,7 +515,7 @@ public class DeleteFileIndex {
}
break;
case EQUALITY_DELETES:
- add(globalDeletes, eqDeletesByPartition, file);
+ add(globalDeletes, eqDeletesByPartition, file, fieldLookup);
break;
default:
throw new UnsupportedOperationException("Unsupported content: " +
file.content());
@@ -536,7 +562,8 @@ public class DeleteFileIndex {
private void add(
EqualityDeletes globalDeletes,
PartitionMap<EqualityDeletes> deletesByPartition,
- DeleteFile file) {
+ DeleteFile file,
+ Function<Integer, Types.NestedField> fieldLookup) {
PartitionSpec spec = specsById.get(file.specId());
EqualityDeletes deletes;
@@ -545,10 +572,11 @@ public class DeleteFileIndex {
} else {
int specId = spec.specId();
StructLike partition = file.partition();
- deletes = deletesByPartition.computeIfAbsent(specId, partition,
EqualityDeletes::new);
+ Supplier<EqualityDeletes> initEqDeletes = () -> new
EqualityDeletes(fieldLookup);
+ deletes = deletesByPartition.computeIfAbsent(specId, partition,
initEqDeletes);
}
- deletes.add(spec, file);
+ deletes.add(file);
}
private Iterable<CloseableIterable<ManifestEntry<DeleteFile>>>
deleteManifestReaders() {
@@ -725,6 +753,8 @@ public class DeleteFileIndex {
Comparator.comparingLong(EqualityDeleteFile::applySequenceNumber);
private static final EqualityDeleteFile[] EMPTY_EQUALITY_DELETES = new
EqualityDeleteFile[0];
+ private final Function<Integer, Types.NestedField> fieldLookup;
+
// indexed state
private long[] seqs = null;
private EqualityDeleteFile[] files = null;
@@ -732,9 +762,13 @@ public class DeleteFileIndex {
// a buffer that is used to hold files before indexing
private volatile List<EqualityDeleteFile> buffer = Lists.newArrayList();
- public void add(PartitionSpec spec, DeleteFile file) {
+ EqualityDeletes(Function<Integer, Types.NestedField> fieldLookup) {
+ this.fieldLookup = fieldLookup;
+ }
+
+ public void add(DeleteFile file) {
Preconditions.checkState(buffer != null, "Can't add files upon
indexing");
- buffer.add(new EqualityDeleteFile(spec, file));
+ buffer.add(new EqualityDeleteFile(fieldLookup, file));
}
public DeleteFile[] filter(long seq, DataFile dataFile) {
@@ -800,15 +834,15 @@ public class DeleteFileIndex {
// an equality delete file wrapper that caches the converted boundaries for
faster boundary checks
// this class is not meant to be exposed beyond the delete file index
private static class EqualityDeleteFile {
- private final PartitionSpec spec;
+ private final Function<Integer, Types.NestedField> fieldLookup;
private final DeleteFile wrapped;
private final long applySequenceNumber;
private volatile List<Types.NestedField> equalityFields = null;
private volatile Map<Integer, Object> convertedLowerBounds = null;
private volatile Map<Integer, Object> convertedUpperBounds = null;
- EqualityDeleteFile(PartitionSpec spec, DeleteFile file) {
- this.spec = spec;
+ EqualityDeleteFile(Function<Integer, Types.NestedField> fieldLookup,
DeleteFile file) {
+ this.fieldLookup = fieldLookup;
this.wrapped = file;
this.applySequenceNumber = wrapped.dataSequenceNumber() - 1;
}
@@ -827,7 +861,8 @@ public class DeleteFileIndex {
if (equalityFields == null) {
List<Types.NestedField> fields = Lists.newArrayList();
for (int id : wrapped.equalityFieldIds()) {
- Types.NestedField field = spec.schema().findField(id);
+ Types.NestedField field = fieldLookup.apply(id);
+ Preconditions.checkArgument(field != null, "Cannot find field
for ID %s", id);
fields.add(field);
}
this.equalityFields = fields;
@@ -890,7 +925,7 @@ public class DeleteFileIndex {
if (bounds != null) {
for (Types.NestedField field : equalityFields()) {
int id = field.fieldId();
- Type type = spec.schema().findField(id).type();
+ Type type = field.type();
if (type.isPrimitiveType()) {
ByteBuffer bound = bounds.get(id);
if (bound != null) {
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/IcebergExternalMetaCacheTest.java
b/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/IcebergExternalMetaCacheTest.java
index 7d559a23b07..1b51e437825 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/IcebergExternalMetaCacheTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/IcebergExternalMetaCacheTest.java
@@ -1955,14 +1955,18 @@ public class IcebergExternalMetaCacheTest {
public void testWeightedV2ManifestListMaterializesOnlyInQueryView() throws
Exception {
Schema schema = new Schema(Types.NestedField.required(1, "id",
Types.IntegerType.get()));
String tableLocation =
temporaryFolder.newFolder("v2-table").toURI().toString();
- Table liveTable = new HadoopTables(new Configuration()).create(
+ HadoopTables tables = new HadoopTables(new Configuration());
+ Table writerTable = tables.create(
schema, PartitionSpec.unpartitioned(), tableLocation);
- liveTable.newAppend().appendFile(
- DataFiles.builder(liveTable.spec())
+ writerTable.newAppend().appendFile(
+ DataFiles.builder(writerTable.spec())
.withPath(tableLocation + "/data/a.parquet")
.withFileSizeInBytes(10L)
.withRecordCount(1L)
.build()).commit();
+ // Iceberg 1.11 may retain eagerly loaded manifests on the writer-side
snapshot after commit,
+ // while the cache invariant applies to the lazy snapshots
reconstructed by catalog loads.
+ Table liveTable = tables.load(tableLocation);
Assert.assertNotNull(liveTable.currentSnapshot().manifestListLocation());
IcebergTableCacheValue value = new IcebergTableCacheValue(liveTable);
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 941f30c4dd7..8bb684ce1b8 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
@@ -2938,6 +2938,63 @@ public class IcebergScanNodeTest {
Mockito.verify(scan).filter(Mockito.argThat(expression ->
expression.toString().contains("old_name")));
}
+ @Test
+ public void testHistoricalPredicatePlansAfterColumnRename() throws
Exception {
+ assertHistoricalPredicatePlansAfterSchemaEvolution(false);
+ }
+
+ @Test
+ public void testHistoricalPredicatePlansAfterColumnDrop() throws Exception
{
+ assertHistoricalPredicatePlansAfterSchemaEvolution(true);
+ }
+
+ private void assertHistoricalPredicatePlansAfterSchemaEvolution(boolean
dropColumn) 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();
+ Table table = tables.create(
+ historicalSchema, PartitionSpec.unpartitioned(),
SortOrder.unsorted(),
+ ImmutableMap.of(TableProperties.FORMAT_VERSION, "2"),
tableLocation);
+ DataFile historicalDataFile = DataFiles.builder(table.spec())
+ .withPath(tableLocation + "/data/historical.parquet")
+ .withFormat(FileFormat.PARQUET)
+ .withFileSizeInBytes(10)
+ .withRecordCount(2)
+ .build();
+ table.newFastAppend().appendFile(historicalDataFile).commit();
+ long historicalSnapshotId = table.currentSnapshot().snapshotId();
+ int historicalSchemaId = table.currentSnapshot().schemaId();
+
+ if (dropColumn) {
+ table.updateSchema().deleteColumn("x").commit();
+ } 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());
+ }
+
@Test
public void testPinnedBranchUsesFrozenSnapshotWithCurrentSchema() throws
Exception {
Schema snapshotSchema = new Schema(11, ImmutableList.of(
diff --git a/fe/pom.xml b/fe/pom.xml
index d1a7ece5d31..8349b523c96 100644
--- a/fe/pom.xml
+++ b/fe/pom.xml
@@ -231,7 +231,7 @@ under the License.
<module>fe-authentication</module>
</modules>
<properties>
-
<doris.hive.catalog.shade.version>3.1.2</doris.hive.catalog.shade.version>
+
<doris.hive.catalog.shade.version>3.1.3</doris.hive.catalog.shade.version>
<!-- iceberg 1.9.1 depends avro on 1.12 -->
<avro.version>1.12.1</avro.version>
<parquet.version>1.17.0</parquet.version>
@@ -348,7 +348,7 @@ under the License.
<!-- ATTN: avro version must be consistent with Iceberg version -->
<!-- Please modify iceberg.version and avro.version together,
you can find avro version info in iceberg mvn repository -->
- <iceberg.version>1.10.1</iceberg.version>
+ <iceberg.version>1.11.0</iceberg.version>
<lance.version>9.1.0-beta.3</lance.version>
<substrait.version>0.40.0</substrait.version>
<!-- 0.56.1 has bug that "SplitMode" in query response may not be
set-->
diff --git
a/regression-test/data/external_table_p0/iceberg/iceberg_schema_change_ddl_with_branch.out
b/regression-test/data/external_table_p0/iceberg/iceberg_schema_change_ddl_with_branch.out
index cf3ccac5bc8..c6af1ae9181 100644
---
a/regression-test/data/external_table_p0/iceberg/iceberg_schema_change_ddl_with_branch.out
+++
b/regression-test/data/external_table_p0/iceberg/iceberg_schema_change_ddl_with_branch.out
@@ -250,6 +250,11 @@ phone text Yes true \N
4 [email protected]
-- !all_branches_have_phone --
+1 \N
+2 \N
+3 \N
+4 \N
+5 \N
-- !summary_main --
1 Alice 95.5 \N \N
diff --git
a/regression-test/data/external_table_p0/iceberg/test_iceberg_sys_table.out
b/regression-test/data/external_table_p0/iceberg/test_iceberg_sys_table.out
index 6a5298b356f..3692b06a99e 100644
--- a/regression-test/data/external_table_p0/iceberg/test_iceberg_sys_table.out
+++ b/regression-test/data/external_table_p0/iceberg/test_iceberg_sys_table.out
@@ -478,6 +478,7 @@ deleted_data_files_count int Yes true \N
NONE
deleted_delete_files_count int Yes true \N NONE
existing_data_files_count int Yes true \N NONE
existing_delete_files_count int Yes true \N NONE
+key_metadata text Yes true \N NONE
length bigint Yes true \N NONE
partition_spec_id int Yes true \N NONE
partition_summaries array<struct<contains_null:boolean not
null,contains_nan:boolean not null,lower_bound:text,upper_bound:text>> Yes
true \N NONE
@@ -986,6 +987,7 @@ deleted_data_files_count int Yes true \N
NONE
deleted_delete_files_count int Yes true \N NONE
existing_data_files_count int Yes true \N NONE
existing_delete_files_count int Yes true \N NONE
+key_metadata text Yes true \N NONE
length bigint Yes true \N NONE
partition_spec_id int Yes true \N NONE
partition_summaries array<struct<contains_null:boolean not
null,contains_nan:boolean not null,lower_bound:text,upper_bound:text>> Yes
true \N NONE
diff --git
a/regression-test/suites/external_table_p0/iceberg/iceberg_schema_change_ddl_with_branch.groovy
b/regression-test/suites/external_table_p0/iceberg/iceberg_schema_change_ddl_with_branch.groovy
index 047b75cf33e..b4f8124e1fe 100644
---
a/regression-test/suites/external_table_p0/iceberg/iceberg_schema_change_ddl_with_branch.groovy
+++
b/regression-test/suites/external_table_p0/iceberg/iceberg_schema_change_ddl_with_branch.groovy
@@ -198,10 +198,11 @@ suite("iceberg_schema_change_ddl_with_branch",
"p0,external,doris,external_docke
// All branches expose the current table columns: id, name, grade, email,
phone.
- // Verify all branches have the latest columns
- qt_all_branches_have_grade """ SELECT id, grade FROM
${branch_table_name}@branch(branch1) WHERE grade > 0 ORDER BY id """
+ // Iceberg validates filters against the referenced snapshot schema, so
columns renamed or
+ // added later are verified through projection instead of predicates on
historical branches.
+ qt_all_branches_have_grade """ SELECT id, grade FROM
${branch_table_name}@branch(branch1) ORDER BY id """
qt_all_branches_have_email """ SELECT id, email FROM
${branch_table_name}@branch(branch2) WHERE email IS NOT NULL ORDER BY id """
- qt_all_branches_have_phone """ SELECT id, phone FROM
${branch_table_name}@branch(branch3) WHERE phone IS NOT NULL ORDER BY id """
+ qt_all_branches_have_phone """ SELECT id, phone FROM
${branch_table_name}@branch(branch3) ORDER BY id """
// All branches should NOT have old columns that were dropped/renamed
test {
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]