This is an automated email from the ASF dual-hosted git repository.
Gabriel39 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/master by this push:
new 09b4b730c9f [fix](iceberg) Do not compare write defaults after the
write is bound (#68463)
09b4b730c9f is described below
commit 09b4b730c9f21b159e302433b4f336e2933213f5
Author: daidai <[email protected]>
AuthorDate: Tue Sep 29 18:11:42 2026 +0800
[fix](iceberg) Do not compare write defaults after the write is bound
(#68463)
### What problem does this PR solve?
Issue Number: None
Related PR: #66345
Problem Summary:
`rewrite_data_files` fails with `Iceberg table schema changed after the
write was bound` on any table whose column has a write default. The
rewrite binds the cached read schema, which carries no write default, so
the default comparison always mismatches.
For INSERT/OVERWRITE/UPDATE/MERGE the comparison is unreachable: a
default change always commits a new schema id, which the
schema-generation fences (#65851, #66345) reject first. This PR removes
the comparison; name, type, field-id and column-order checks are kept.
### Release note
Fix `rewrite_data_files` failing on Iceberg tables with column write
defaults.
### 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
- Behavior changed:
- [x] No.
- [ ] Yes.
- 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
---
.../iceberg/IcebergWritePlanProvider.java | 15 ++-----
.../iceberg/IcebergWritePlanProviderTest.java | 52 ++++++++++++++++------
.../iceberg/test_iceberg_write_default.groovy | 10 +++++
3 files changed, 53 insertions(+), 24 deletions(-)
diff --git
a/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergWritePlanProvider.java
b/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergWritePlanProvider.java
index 83483dc83b8..2611f56b074 100644
---
a/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergWritePlanProvider.java
+++
b/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergWritePlanProvider.java
@@ -75,7 +75,6 @@ import java.util.HashSet;
import java.util.List;
import java.util.Locale;
import java.util.Map;
-import java.util.Objects;
import java.util.Optional;
import java.util.Set;
import java.util.function.Function;
@@ -311,21 +310,15 @@ public class IcebergWritePlanProvider implements
ConnectorWritePlanProvider {
ConnectorColumn bound = boundColumns.get(i);
ConnectorType currentType = IcebergTypeMapping.fromIcebergType(
current.type(), enableVarbinary, enableTimestampTz);
- String boundDefaultSql = bound.getDefaultValueSql();
- if ("NULL".equalsIgnoreCase(boundDefaultSql)) {
- boundDefaultSql = null;
- }
- String currentDefaultSql = current.writeDefault() == null ? null
- : IcebergWriteSchemaContext.toDorisSql(current.type(),
current.writeDefault(),
- enableVarbinary, enableTimestampTz);
// Do not compare top-level nullability: Doris widens Iceberg
required columns in its read schema
// so evolution default-fill may yield NULL. Nested requiredness
remains authoritative in
// sameBoundType, while current schema JSON enforces writes at the
root.
+ // Do not compare write defaults either. A default change always
commits a new schema id, which the
+ // schema-generation fences reject for every write that pins a
schema context. REWRITE materializes
+ // no default and binds the cached read schema, which carries
none, so a comparison only rejects it.
if (!current.name().equalsIgnoreCase(bound.getName())
|| !sameBoundType(currentType, bound.getType())
- || (bound.getUniqueId() >= 0 && current.fieldId() !=
bound.getUniqueId())
- // Omitted columns and DEFAULT expressions were already
materialized from this value at bind.
- || !Objects.equals(boundDefaultSql, currentDefaultSql)) {
+ || (bound.getUniqueId() >= 0 && current.fieldId() !=
bound.getUniqueId())) {
// BE maps write expressions to schema-json by ordinal, so
accepting a reordered live
// schema here could silently place values under the wrong
Iceberg field names.
throw new DorisConnectorException("Iceberg table schema
changed after the write was bound; retry "
diff --git
a/fe/fe-connector/fe-connector-iceberg/src/test/java/org/apache/doris/connector/iceberg/IcebergWritePlanProviderTest.java
b/fe/fe-connector/fe-connector-iceberg/src/test/java/org/apache/doris/connector/iceberg/IcebergWritePlanProviderTest.java
index 4d5040992e4..fd2e4b27d78 100644
---
a/fe/fe-connector/fe-connector-iceberg/src/test/java/org/apache/doris/connector/iceberg/IcebergWritePlanProviderTest.java
+++
b/fe/fe-connector/fe-connector-iceberg/src/test/java/org/apache/doris/connector/iceberg/IcebergWritePlanProviderTest.java
@@ -1194,25 +1194,28 @@ public class IcebergWritePlanProviderTest {
}
@Test
- public void planWriteRejectsWriteDefaultEvolution() {
+ public void planWriteRejectsWriteDefaultEvolutionBySchemaGeneration() {
InMemoryCatalog catalog = freshCatalog();
- Table table = unpartitionedUnsortedTable(catalog);
+ Table table = formatVersionThreeTable(catalog);
table.updateSchema().updateColumnDefault("id",
Literal.of(42)).commit();
- List<ConnectorColumn> boundColumns = Arrays.asList(
- new ConnectorColumn("id", ConnectorType.of("INT"), "", false,
null)
- .withDefaultValueSql("42")
-
.withUniqueId(table.schema().findField("id").fieldId()),
- new ConnectorColumn("name", ConnectorType.of("STRING"), "",
true, null)
-
.withUniqueId(table.schema().findField("name").fieldId()));
+ RecordingConnectorContext context = contextWithStorage();
+ IcebergWritePlanProvider provider = providerFor(table, context);
+ WriteSession session = sessionFor(table, context);
+ IcebergTableHandle tableHandle = new IcebergTableHandle("db1", "tv3");
+ // Bind columns and generation in one statement scope, as a real
INSERT does.
+ List<ConnectorColumn> boundColumns = provider.getWriteColumns(
+ session, tableHandle,
Optional.empty()).orElseThrow(AssertionError::new);
+ String boundIdentity = provider.getWriteMetadataIdentity(session,
tableHandle);
table.updateSchema().updateColumnDefault("id", Literal.of(7)).commit();
+ // A default change always commits a new schema id, so the
schema-generation fences reject the stale
+ // write before any column comparison; write defaults need no
comparison of their own.
DorisConnectorException ex =
Assertions.assertThrows(DorisConnectorException.class,
- () -> planSink(table, contextWithStorage(),
- new WriteHandle(new IcebergTableHandle("db1", "t2"))
- .boundTargetColumns(boundColumns)));
- Assertions.assertTrue(ex.getMessage().contains("schema changed"),
- "a statement must retry instead of writing a value
materialized from the stale default");
+ () -> provider.planWrite(session, new WriteHandle(tableHandle)
+ .boundTargetColumns(boundColumns)
+ .boundWriteMetadataIdentity(boundIdentity)));
+ Assertions.assertTrue(ex.getMessage().contains("changed"),
ex.getMessage());
}
@Test
@@ -1234,6 +1237,29 @@ public class IcebergWritePlanProviderTest {
Assertions.assertTrue(sink.getSchemaJson().contains("\"write-default\":42"));
}
+ @Test
+ public void planRewriteAcceptsStableWriteDefault() {
+ InMemoryCatalog catalog = freshCatalog();
+ Table table = formatVersionThreeTable(catalog);
+ table.updateSchema().addColumn("bonus", Types.IntegerType.get(),
Literal.of(7)).commit();
+ table.updateSchema().updateColumnDefault("bonus",
Literal.of(9)).commit();
+ // rewrite_data_files binds the cached read schema, which deliberately
carries no write default.
+ List<ConnectorColumn> boundColumns = new
ArrayList<>(boundDataColumns(table));
+ boundColumns.add(new ConnectorColumn("bonus", ConnectorType.of("INT"),
"", true, null)
+ .withUniqueId(table.schema().findField("bonus").fieldId()));
+ boundColumns.add(new ConnectorColumn("_row_id",
ConnectorType.of("BIGINT"), "", true, null)
+ .invisible().reservedPassthrough());
+ boundColumns.add(new ConnectorColumn(
+ "_last_updated_sequence_number", ConnectorType.of("BIGINT"),
"", true, null)
+ .invisible().reservedPassthrough());
+
+ TIcebergTableSink sink = Assertions.assertDoesNotThrow(() ->
planSink(table, contextWithStorage(),
+ new WriteHandle(new IcebergTableHandle("db1", "tv3"))
+ .boundTargetColumns(boundColumns)
+ .writeOperation(WriteOperation.REWRITE)));
+ Assertions.assertEquals(TIcebergWriteType.REWRITE,
sink.getWriteType());
+ }
+
@Test
public void planRewriteRejectsV2BoundSchemaAtV3Planning() {
InMemoryCatalog catalog = freshCatalog();
diff --git
a/regression-test/suites/external_table_p0/iceberg/test_iceberg_write_default.groovy
b/regression-test/suites/external_table_p0/iceberg/test_iceberg_write_default.groovy
index 1212826cdb7..9759a32461e 100644
---
a/regression-test/suites/external_table_p0/iceberg/test_iceberg_write_default.groovy
+++
b/regression-test/suites/external_table_p0/iceberg/test_iceberg_write_default.groovy
@@ -77,6 +77,16 @@ suite("test_iceberg_write_default", "p0,external") {
assertEquals(1, rows.size())
assertEquals("42", rows[0][1].toString(),
"an INSERT omitting column c must apply the iceberg write
default 42, got: " + rows[0][1])
+
+ // 3) rewrite_data_files binds the cached schema, which has no write
default. The stable default must
+ // not be mistaken for a concurrent schema change, and rewritten rows
keep their stored values.
+ sql """ INSERT INTO ${tbl} (id, c) VALUES (2, NULL) """
+ def rewriteResult = sql """ ALTER TABLE ${tbl} EXECUTE
rewrite_data_files("rewrite-all" = "true") """
+ assertEquals(2, rewriteResult[0][0] as int, "both data files must be
rewritten, got: " + rewriteResult)
+ rows = sql """ SELECT id, c FROM ${tbl} ORDER BY id """
+ assertEquals(2, rows.size())
+ assertEquals("42", rows[0][1].toString())
+ assertNull(rows[1][1], "an explicit NULL must survive the rewrite,
got: " + rows[1][1])
} finally {
sql """drop table if exists ${catalog_name}.${db}.${tbl}"""
sql """drop database if exists ${catalog_name}.${db} force"""
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]