hubgeter commented on code in PR #65851: URL: https://github.com/apache/doris/pull/65851#discussion_r3666090032
########## fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/IcebergWriteSchemaContext.java: ########## @@ -0,0 +1,655 @@ +// 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.datasource.mvcc.MvccSnapshot; +import org.apache.doris.datasource.mvcc.MvccUtil; +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.FileFormat; +import org.apache.iceberg.MetricsConfig; +import org.apache.iceberg.PartitionField; +import org.apache.iceberg.PartitionSpec; +import org.apache.iceberg.PartitionSpecParser; +import org.apache.iceberg.Schema; +import org.apache.iceberg.SchemaParser; +import org.apache.iceberg.SnapshotRef; +import org.apache.iceberg.SortField; +import org.apache.iceberg.SortOrder; +import org.apache.iceberg.SortOrderParser; +import org.apache.iceberg.StructLike; +import org.apache.iceberg.Table; +import org.apache.iceberg.TableProperties; +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 Schema mergeSchema; + private final String mergeSchemaJson; + private final PartitionSpec partitionSpec; + private final String partitionSpecJson; + private final SortOrder sortOrder; + private final String sortOrderJson; + private final FileFormat fileFormat; + private final MetricsConfig metricsConfig; + private final String fileCompression; + private final String dataLocation; + private final Map<String, String> writerProperties; + 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 statement snapshot's current table 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(); + Schema schema = branchName.isPresent() + ? resolveBranchSchema(table, branchName.get(), dorisTable.getName()) + : resolveStatementSchema(table, dorisTable); + int formatVersion = IcebergUtils.getFormatVersion(table); + Map<String, String> properties = ImmutableMap.copyOf(table.properties()); + return new IcebergWriteSchemaContext( + dorisTable.getId(), dorisTable.getName(), schema, formatVersion, branchName, + bindPartitionSpec(table.spec(), schema, dorisTable.getName()), + bindSortOrder(table.sortOrder(), schema, dorisTable.getName()), + IcebergUtils.getFileFormat(table), MetricsConfig.forTable(table), + IcebergUtils.getFileCompress(table), IcebergUtils.dataLocation(table), properties, + 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(), PartitionSpec.unpartitioned(), SortOrder.unsorted(), + FileFormat.PARQUET, MetricsConfig.getDefault(), + TableProperties.PARQUET_COMPRESSION_DEFAULT_SINCE_1_4_0, + "file:///tmp/test_table/data", ImmutableMap.of(), + enableMappingVarbinary, enableMappingTimestampTz); + } + + @VisibleForTesting + public static IcebergWriteSchemaContext forSchema(Schema schema, int formatVersion, + PartitionSpec partitionSpec, SortOrder sortOrder, FileFormat fileFormat, + MetricsConfig metricsConfig, String fileCompression, String dataLocation, + Map<String, String> writerProperties, + boolean enableMappingVarbinary, boolean enableMappingTimestampTz) { + return new IcebergWriteSchemaContext(-1L, "test_table", schema, formatVersion, + Optional.empty(), partitionSpec, sortOrder, fileFormat, metricsConfig, + fileCompression, dataLocation, writerProperties, + enableMappingVarbinary, enableMappingTimestampTz); + } + + private IcebergWriteSchemaContext(long tableId, String tableName, Schema schema, + int formatVersion, Optional<String> branchName, + PartitionSpec partitionSpec, SortOrder sortOrder, FileFormat fileFormat, + MetricsConfig metricsConfig, String fileCompression, String dataLocation, + Map<String, String> writerProperties, + 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); + this.mergeSchema = formatVersion >= IcebergUtils.ICEBERG_ROW_LINEAGE_MIN_VERSION + ? IcebergUtils.appendRowLineageFieldsForV3(schema) : schema; + this.mergeSchemaJson = SchemaParser.toJson(mergeSchema); + this.partitionSpec = Objects.requireNonNull(partitionSpec, "partitionSpec should not be null"); + this.partitionSpecJson = PartitionSpecParser.toJson(partitionSpec); + this.sortOrder = Objects.requireNonNull(sortOrder, "sortOrder should not be null"); + this.sortOrderJson = SortOrderParser.toJson(sortOrder); + this.fileFormat = Objects.requireNonNull(fileFormat, "fileFormat should not be null"); + this.metricsConfig = Objects.requireNonNull(metricsConfig, "metricsConfig should not be null"); + this.fileCompression = Objects.requireNonNull( + fileCompression, "fileCompression should not be null"); + this.dataLocation = Objects.requireNonNull(dataLocation, "dataLocation should not be null"); + this.writerProperties = ImmutableMap.copyOf( + Objects.requireNonNull(writerProperties, "writerProperties should not be null")); + validateWriterMetadataSources(schema, partitionSpec, sortOrder, tableName); + + 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 PartitionSpec bindPartitionSpec( + PartitionSpec partitionSpec, Schema schema, String tableName) { + if (!partitionSpec.isPartitioned()) { + return PartitionSpec.builderFor(schema) + .withSpecId(partitionSpec.specId()) + .build(); + } + try { + return PartitionSpecParser.fromJson(schema, PartitionSpecParser.toJson(partitionSpec)); + } catch (RuntimeException e) { + throw new AnalysisException("Iceberg partition spec " + partitionSpec.specId() + + " is incompatible with pinned schema " + schema.schemaId() + + " for table " + tableName + ": " + e.getMessage(), e); + } + } + + private static SortOrder bindSortOrder(SortOrder sortOrder, Schema schema, String tableName) { + if (!sortOrder.isSorted()) { + return SortOrder.unsorted(); + } + try { + return SortOrderParser.fromJson(schema, SortOrderParser.toJson(sortOrder)); + } catch (RuntimeException e) { + throw new AnalysisException("Iceberg sort order " + sortOrder.orderId() + + " is incompatible with pinned schema " + schema.schemaId() + + " for table " + tableName + ": " + e.getMessage(), e); + } + } + + private static void validateWriterMetadataSources( + Schema schema, PartitionSpec partitionSpec, SortOrder sortOrder, String tableName) { + Map<Integer, Types.NestedField> topLevelFields = schema.columns().stream() + .collect(ImmutableMap.toImmutableMap(Types.NestedField::fieldId, field -> field)); + for (PartitionField field : partitionSpec.fields()) { + if (!topLevelFields.containsKey(field.sourceId())) { + throw new AnalysisException("Iceberg partition field " + field.fieldId() + + " references source field " + field.sourceId() + + " outside pinned top-level schema " + schema.schemaId() + + " for table " + tableName); + } + } + for (SortField field : sortOrder.fields()) { + if (schema.findField(field.sourceId()) == null) { + throw new AnalysisException("Iceberg sort field references source field " + + field.sourceId() + " outside pinned schema " + schema.schemaId() + + " for table " + tableName); + } + } + } + + private static Schema resolveBranchSchema(Table table, String branchName, String tableName) { + SnapshotRef ref = table.refs().get(branchName); + if (ref == null) { + throw new AnalysisException(branchName + " is not founded in " + tableName); + } + if (!ref.isBranch()) { + throw new AnalysisException(branchName + + " is a tag, not a branch. Tags cannot be targets for producing snapshots"); + } + return SnapshotUtil.schemaFor(table, ref.snapshotId()); + } + + private static Schema resolveStatementSchema(Table table, IcebergExternalTable dorisTable) { + Optional<MvccSnapshot> snapshot = MvccUtil.getSnapshotFromContext(dorisTable); + if (!snapshot.isPresent()) { + return table.schema(); + } + Preconditions.checkState(snapshot.get() instanceof IcebergMvccSnapshot, + "Expected an Iceberg MVCC snapshot for table %s", dorisTable.getName()); + long schemaId = ((IcebergMvccSnapshot) snapshot.get()) + .getSnapshotCacheValue().getSnapshot().getSchemaId(); + Schema schema = table.schemas().get(Math.toIntExact(schemaId)); + return Preconditions.checkNotNull(schema, + "Iceberg schema %s is not available in the statement table metadata for %s", + schemaId, dorisTable.getName()); + } + + /** Resolve a write default by the pinned target field name. */ + public Expression resolveWriteDefault(String columnName) { + Column column = columns.stream() + .filter(targetColumn -> targetColumn.getName().equalsIgnoreCase(columnName)) + .findFirst() + .orElseThrow(() -> new AnalysisException( + "Cannot find column information for DEFAULT(" + columnName + ")")); + return resolveWriteDefault(column); + } + + /** Resolve the value used for an omitted column or an explicit DEFAULT. */ + public Expression resolveWriteDefault(Column column) { + Types.NestedField field = fieldsById.get(column.getUniqueId()); + if (field == null) { + throw new AnalysisException("Column " + column.getName() + + " is not present in pinned Iceberg schema " + getSchemaId()); + } + Expression writeDefault = writeDefaultsById.get(field.fieldId()); + if (writeDefault != null) { + return writeDefault; + } + DataType targetType = DataType.fromCatalogType(column.getType()); + if (field.isOptional()) { + return new NullLiteral(targetType); + } + throw new AnalysisException("Column has no write default and is required, column=" + field.name()); + } + + /** Validate that the fresh table can commit files described by the pinned writer metadata. */ + public void validateCurrentSchema(Table table) { + validateCurrentSchema(table, false); + } + + /** + * Validate that the fresh table can commit files described by the pinned writer metadata. + * + * <p>Static partition overwrite additionally requires the pinned spec to remain current because + * its replacement filter was planned from that spec. Appends can safely write an older retained + * spec, so they only require the pinned definition to remain available. + */ + public void validateCurrentSchema(Table table, boolean requireCurrentPartitionSpec) { + Schema currentSchema = branchName.isPresent() Review Comment: Fixed in 990c536eeff. Branch planning now compares the pinned branch-head writer schema with the table-current schema that Iceberg stamps on the new snapshot, rejecting required/no-initial-default fields that the branch writer cannot satisfy. The same check runs again against the fresh table before the transaction starts. Added tests for both an already-diverged first branch append and a schema change between planning and fresh-table validation; the expanded Iceberg FE suite passes 96/96 and the full FE build passes. -- 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]
