924060929 commented on code in PR #65851:
URL: https://github.com/apache/doris/pull/65851#discussion_r3654263268


##########
fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/IcebergWriteSchemaContext.java:
##########
@@ -0,0 +1,452 @@
+// 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.doris.datasource.iceberg;
+
+import org.apache.doris.catalog.Column;
+import org.apache.doris.nereids.exceptions.AnalysisException;
+import org.apache.doris.nereids.trees.expressions.Expression;
+import org.apache.doris.nereids.trees.expressions.functions.scalar.Array;
+import org.apache.doris.nereids.trees.expressions.functions.scalar.CreateMap;
+import 
org.apache.doris.nereids.trees.expressions.functions.scalar.CreateNamedStruct;
+import org.apache.doris.nereids.trees.expressions.functions.scalar.Unhex;
+import org.apache.doris.nereids.trees.expressions.literal.ArrayLiteral;
+import org.apache.doris.nereids.trees.expressions.literal.BigIntLiteral;
+import org.apache.doris.nereids.trees.expressions.literal.BooleanLiteral;
+import org.apache.doris.nereids.trees.expressions.literal.DateTimeV2Literal;
+import org.apache.doris.nereids.trees.expressions.literal.DateV2Literal;
+import org.apache.doris.nereids.trees.expressions.literal.DecimalV3Literal;
+import org.apache.doris.nereids.trees.expressions.literal.DoubleLiteral;
+import org.apache.doris.nereids.trees.expressions.literal.FloatLiteral;
+import org.apache.doris.nereids.trees.expressions.literal.IntegerLiteral;
+import org.apache.doris.nereids.trees.expressions.literal.Literal;
+import org.apache.doris.nereids.trees.expressions.literal.MapLiteral;
+import org.apache.doris.nereids.trees.expressions.literal.NullLiteral;
+import org.apache.doris.nereids.trees.expressions.literal.StringLiteral;
+import org.apache.doris.nereids.trees.expressions.literal.StructLiteral;
+import org.apache.doris.nereids.trees.expressions.literal.TimestampTzLiteral;
+import org.apache.doris.nereids.trees.expressions.literal.VarBinaryLiteral;
+import org.apache.doris.nereids.types.DataType;
+import org.apache.doris.nereids.types.DateTimeV2Type;
+import org.apache.doris.nereids.types.DecimalV3Type;
+import org.apache.doris.nereids.types.StructType;
+import org.apache.doris.nereids.types.TimeStampTzType;
+import org.apache.doris.nereids.types.VarBinaryType;
+import org.apache.doris.nereids.util.TypeCoercionUtils;
+
+import com.google.common.annotations.VisibleForTesting;
+import com.google.common.base.Preconditions;
+import com.google.common.collect.ImmutableList;
+import com.google.common.collect.ImmutableMap;
+import com.google.common.io.BaseEncoding;
+import org.apache.iceberg.Schema;
+import org.apache.iceberg.SchemaParser;
+import org.apache.iceberg.SnapshotRef;
+import org.apache.iceberg.StructLike;
+import org.apache.iceberg.Table;
+import org.apache.iceberg.types.Type;
+import org.apache.iceberg.types.Types;
+import org.apache.iceberg.util.SnapshotUtil;
+
+import java.math.BigDecimal;
+import java.nio.ByteBuffer;
+import java.time.Instant;
+import java.time.LocalDate;
+import java.time.LocalDateTime;
+import java.time.ZoneOffset;
+import java.util.ArrayList;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.Optional;
+import java.util.UUID;
+
+/**
+ * Statement-scoped Iceberg write schema and write-default values.
+ *
+ * <p>The context pins one Iceberg schema before analysis. The analyzer, 
planner sink and
+ * transaction preflight must all use this same instance so a concurrent 
schema change cannot
+ * combine expressions from one schema with a writer schema from another one.
+ */
+public final class IcebergWriteSchemaContext {
+    private final long tableId;
+    private final String tableName;
+    private final Schema schema;
+    private final int formatVersion;
+    private final Optional<String> branchName;
+    private final String schemaJson;
+    private final String mergeSchemaJson;
+    private final List<Column> columns;
+    private final List<Column> mergeColumns;
+    private final Map<Integer, Types.NestedField> fieldsById;
+    private final Map<Integer, Expression> writeDefaultsById;
+
+    /** Pin the current main or branch schema under the catalog authentication 
boundary. */
+    public static IcebergWriteSchemaContext create(
+            IcebergExternalTable dorisTable, Optional<String> branchName) {
+        Objects.requireNonNull(dorisTable, "dorisTable should not be null");
+        Objects.requireNonNull(branchName, "branchName should not be null");
+        try {
+            return 
dorisTable.getCatalog().getExecutionAuthenticator().execute(() -> {
+                Table table = dorisTable.getIcebergTable();
+                table.refresh();
+                Schema schema = resolveSchema(table, branchName, 
dorisTable.getName());
+                int formatVersion = IcebergUtils.getFormatVersion(table);
+                return new IcebergWriteSchemaContext(
+                        dorisTable.getId(), dorisTable.getName(), schema, 
formatVersion, branchName,
+                        dorisTable.getCatalog().getEnableMappingVarbinary(),
+                        dorisTable.getCatalog().getEnableMappingTimestampTz());
+            });
+        } catch (Exception e) {
+            throw new AnalysisException("Failed to pin Iceberg write schema 
for table "
+                    + dorisTable.getName() + ": " + e.getMessage(), e);
+        }
+    }
+
+    @VisibleForTesting
+    public static IcebergWriteSchemaContext forSchema(Schema schema, int 
formatVersion,
+            boolean enableMappingVarbinary, boolean enableMappingTimestampTz) {
+        return new IcebergWriteSchemaContext(-1L, "test_table", schema, 
formatVersion,
+                Optional.empty(), enableMappingVarbinary, 
enableMappingTimestampTz);
+    }
+
+    private IcebergWriteSchemaContext(long tableId, String tableName, Schema 
schema,
+            int formatVersion, Optional<String> branchName,
+            boolean enableMappingVarbinary, boolean enableMappingTimestampTz) {
+        this.tableId = tableId;
+        this.tableName = Objects.requireNonNull(tableName, "tableName should 
not be null");
+        this.schema = Objects.requireNonNull(schema, "schema should not be 
null");
+        this.formatVersion = formatVersion;
+        this.branchName = Objects.requireNonNull(branchName, "branchName 
should not be null");
+        this.schemaJson = SchemaParser.toJson(schema);
+        Schema mergeSchema = formatVersion >= 
IcebergUtils.ICEBERG_ROW_LINEAGE_MIN_VERSION
+                ? IcebergUtils.appendRowLineageFieldsForV3(schema) : schema;
+        this.mergeSchemaJson = SchemaParser.toJson(mergeSchema);
+
+        List<Column> parsedColumns = IcebergUtils.parseSchema(
+                schema, enableMappingVarbinary, enableMappingTimestampTz);
+        this.columns = ImmutableList.copyOf(parsedColumns);
+        List<Column> writerColumns = new ArrayList<>(parsedColumns);
+        writerColumns.add(IcebergRowId.createHiddenColumn());
+        if (formatVersion >= IcebergUtils.ICEBERG_ROW_LINEAGE_MIN_VERSION) {
+            Column rowIdColumn = IcebergUtils.parseField(
+                    org.apache.iceberg.MetadataColumns.ROW_ID,
+                    enableMappingVarbinary, enableMappingTimestampTz);
+            rowIdColumn.setIsVisible(false);
+            writerColumns.add(rowIdColumn);
+            Column sequenceColumn = IcebergUtils.parseField(
+                    
org.apache.iceberg.MetadataColumns.LAST_UPDATED_SEQUENCE_NUMBER,
+                    enableMappingVarbinary, enableMappingTimestampTz);
+            sequenceColumn.setIsVisible(false);
+            writerColumns.add(sequenceColumn);
+        }
+        this.mergeColumns = ImmutableList.copyOf(writerColumns);
+
+        ImmutableMap.Builder<Integer, Types.NestedField> byId = 
ImmutableMap.builder();
+        ImmutableMap.Builder<Integer, Expression> defaults = 
ImmutableMap.builder();
+        for (Types.NestedField field : schema.columns()) {
+            byId.put(field.fieldId(), field);
+            if (field.writeDefault() != null) {
+                DataType targetType = 
DataType.fromCatalogType(IcebergUtils.icebergTypeToDorisType(
+                        field.type(), enableMappingVarbinary, 
enableMappingTimestampTz));
+                defaults.put(field.fieldId(), toDorisExpression(
+                        field.type(), field.writeDefault(), targetType,
+                        enableMappingVarbinary, enableMappingTimestampTz));
+            }
+        }
+        this.fieldsById = byId.build();
+        this.writeDefaultsById = defaults.build();
+    }
+
+    private static Schema resolveSchema(Table table, Optional<String> 
branchName, String tableName) {
+        if (!branchName.isPresent()) {
+            return table.schema();
+        }
+        SnapshotRef ref = table.refs().get(branchName.get());
+        if (ref == null) {
+            throw new AnalysisException(branchName.get() + " is not founded in 
" + tableName);
+        }
+        if (!ref.isBranch()) {
+            throw new AnalysisException(branchName.get()
+                    + " is a tag, not a branch. Tags cannot be targets for 
producing snapshots");
+        }
+        return SnapshotUtil.schemaFor(table, ref.snapshotId());

Review Comment:
   **⚠️ WITHDRAWN — I was wrong about the Iceberg/Spark semantics here. Please 
disregard the blocker; only the release-note point below still stands.**
   
   Apache Spark's Iceberg connector resolves the write schema for a branch 
exactly the way this PR does:
   
   ```java
   // spark/v4.0/.../source/SparkTable.java
   private Schema snapshotSchema() {
     if (icebergTable instanceof BaseMetadataTable) { return 
icebergTable.schema(); }
     else if (branch != null) { return 
addLineageIfRequired(SnapshotUtil.schemaFor(icebergTable, branch)); }
     else { return addLineageIfRequired(SnapshotUtil.schemaFor(icebergTable, 
snapshotId, null)); }
   }
   
   @Override
   public StructType schema() {   // what the DSv2 analyzer resolves INSERT 
columns against
     if (lazyTableSchema == null) { this.lazyTableSchema = 
SparkSchemaUtil.convert(snapshotSchema()); }
     return lazyTableSchema;
   }
   ```
   
   So `SnapshotUtil.schemaFor(table, ref.snapshotId())` is the correct choice, 
and my claimed data-loss scenario does not hold either: a value written for a 
field id that is no longer in the current schema is invisible because the 
column was dropped, which is what DROP COLUMN means. On top of that, Doris 
already reads a branch through the branch snapshot's schema (`getQuerySchema()` 
resolves via `getSpecifiedSnapshot().getSchemaId()`), so this change makes 
branch reads and branch writes agree, whereas before writes used the 
table-level schema. That is a fix, not a regression. Sorry for the noise.
   
   The one point I would still keep, downgraded to P2: this is a user-visible 
behavior change that is not in the PR description or the release note. `INSERT 
INTO t@branch(b) (id, value, new_col)` used to succeed after `ALTER TABLE t ADD 
COLUMN new_col` and now fails with `Unknown column 'new_col' in target table`, 
and two existing regression contracts were flipped to match 
(`iceberg_branch_tag_schema_change_extended` `.out` goes from `3\t30\ttest` to 
`3\t30\t\\N`, and the T09 comment in `test_iceberg_schema_ref_actions_matrix` 
is inverted from *uses main's latest schema* to *uses the branch snapshot's 
schema*). Since it is orthogonal to default values, please call it out 
explicitly in the release note so users hitting the new error know it is 
intentional and that fast-forwarding or committing to the branch is the way 
forward.



-- 
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]

Reply via email to