github-actions[bot] commented on code in PR #65851:
URL: https://github.com/apache/doris/pull/65851#discussion_r3671112037
##########
be/src/format_v2/table/iceberg_reader.cpp:
##########
@@ -511,35 +1061,80 @@ Status
IcebergTableReader::_append_row_position_output_column(format::FileScanRe
return Status::OK();
}
-const format::ColumnDefinition*
IcebergTableReader::_find_equality_delete_data_field(
- const EqualityDeleteFilter& filter, size_t key_idx) const {
+Status IcebergTableReader::_find_equality_delete_data_field(
+ const EqualityDeleteFilter& filter, size_t key_idx,
+ EqualityDeleteColumnPath* data_path) const {
DORIS_CHECK(key_idx < filter.field_ids.size());
DORIS_CHECK(key_idx < filter.field_names.size());
+ DORIS_CHECK(data_path != nullptr);
+ data_path->clear();
if (mapping_mode() != format::TableColumnMappingMode::BY_NAME) {
const int field_id = filter.field_ids[key_idx];
- const auto field_it = std::ranges::find_if(
- _data_reader.file_schema, [field_id](const
format::ColumnDefinition& field) {
- return field.has_identifier_field_id() &&
- field.get_identifier_field_id() == field_id;
- });
- return field_it == _data_reader.file_schema.end() ? nullptr :
&*field_it;
+ static_cast<void>(
+ find_equality_delete_column_path(_data_reader.file_schema,
field_id, data_path));
+ return Status::OK();
}
// Equality keys are hidden scan dependencies and need not appear in the
query projection.
- // Resolve their current name and aliases from the full table schema
supplied by FE, falling
- // back to the delete-file name when history metadata is unavailable.
Reuse ColumnMapper's
- // exact BY_NAME rules so case, string identifiers, and aliases on either
side stay consistent.
- auto table_field = _find_equality_delete_table_field(filter, key_idx);
- return format::find_column_by_name(*table_field, _data_reader.file_schema);
+ // Reuse ColumnMapper's exact BY_NAME rules at every ancestor so a nested
key keeps its
+ // physical path, including historical aliases for ID-less files.
+ auto schema_path =
+
_find_table_column_identity_path_by_field_id(filter.field_ids[key_idx], true);
+ std::vector<const format::ColumnDefinition*> table_path;
+ if (schema_path.has_value()) {
+ for (const auto& field : *schema_path) {
+ table_path.push_back(&field);
+ }
+ } else {
+ static_cast<void>(find_equality_delete_column_path(_projected_columns,
+
filter.field_ids[key_idx], &table_path));
+ }
+ std::optional<format::ColumnDefinition> legacy_table_field;
+ if (table_path.empty() &&
!supports_iceberg_scan_semantics_v2(_scan_params)) {
+ legacy_table_field.emplace();
+ legacy_table_field->name = filter.field_names[key_idx];
+ legacy_table_field->type = filter.key_types[key_idx];
+ table_path.push_back(&*legacy_table_field);
+ }
+ if (table_path.empty()) {
+ return Status::InvalidArgument(
+ "Iceberg equality delete field id {} is absent from current
and historical table "
+ "schema metadata",
+ filter.field_ids[key_idx]);
+ }
+ const std::vector<format::ColumnDefinition>* candidates =
&_data_reader.file_schema;
+ for (size_t index = 0; index < table_path.size(); ++index) {
+ const auto* table_field = table_path[index];
+ DORIS_CHECK(table_field != nullptr);
+ const auto* data_field = format::find_column_by_name(*table_field,
*candidates);
+ if (data_field == nullptr && index + 1 == table_path.size() &&
+ !table_field->has_name_mapping) {
+ // Schema-history fallback can carry a post-snapshot leaf rename
when the target
+ // snapshot's parent has expired. Retry the delete file's original
leaf name, but
+ // never bypass an explicit authoritative Iceberg mapping.
+ format::ColumnDefinition delete_file_field;
+ delete_file_field.name = filter.field_names[key_idx];
+ data_field = format::find_column_by_name(delete_file_field,
*candidates);
+ }
+ if (data_field == nullptr) {
+ data_path->clear();
Review Comment:
[P1] Keep nullable ancestors when defaulting a missing nested key
If a physical parent matches but a descendant is absent, this clears the
whole path; the caller then substitutes one leaf initial-default literal for
every row. Consider an old file with physical `optional struct payload` that
predates child `payload.k` with initial default 7. For a row where `payload` is
NULL, `payload.k` must remain NULL; only a present parent gets 7. An equality
delete for `k = 7` currently sees 7 for both rows and wrongly removes the
parent-NULL row.
V1 has the same batch-wide substitution in `iceberg_reader.cpp` at the
Parquet/ORC missing-key branches and `iceberg_reader_mixin.h`'s
synthesized-column handler. This is distinct from the earlier nested-key fix,
whose tests physically write the leaf. Please retain/read the deepest physical
ancestor and materialize the missing descendant through the recursive default
path so ancestor null maps dominate, with forced V1/V2 Parquet/ORC tests for
NULL and non-NULL parents.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/source/IcebergScanNode.java:
##########
@@ -530,30 +601,837 @@ void enableCurrentIcebergScanSemantics() {
params.setIcebergScanSemanticsVersion(ICEBERG_SCAN_SEMANTICS_VERSION);
}
+ /**
+ * Build the schema metadata carrier used by both scanners and
equality-delete readers.
+ *
+ * <p>Batch-mode delete files are planned asynchronously after scan
parameters are sent to BE.
+ * The authenticated manifest preflight therefore supplies the live
equality field IDs before
+ * the schema carrier is serialized. Only historical fields referenced by
those delete files are
+ * added, so an unrelated dropped type cannot make an otherwise supported
scan fail.
+ */
+ @VisibleForTesting
+ List<NestedField> getSchemaFieldsForScan(
+ Schema scanSchema, Set<Integer> equalityDeleteFieldIds) throws
UserException {
+ List<NestedField> fields = new ArrayList<>(scanSchema.columns());
+ if (isSystemTable || equalityDeleteFieldIds.isEmpty()) {
+ return fields;
+ }
+
+ Set<Integer> missingFieldIds = new HashSet<>(equalityDeleteFieldIds);
+
missingFieldIds.removeAll(TypeUtil.indexById(scanSchema.asStruct()).keySet());
+ if (missingFieldIds.isEmpty()) {
+ return fields;
+ }
+
+ List<Schema> schemaHistory = getMetadataSchemaHistory();
+ // Schema IDs may be reused when evolution returns to an earlier
schema, while the metadata
+ // list may also contain schemas committed after a time-travel or
branch target. Follow the
+ // actual scan snapshot's parent chain first so the field definition
active on that lineage
+ // wins. Then use the complete metadata list as a fallback for
schema-only changes and
+ // expired ancestors. A fallback definition may come from a later
rename, so BE resolves an
+ // ID-less equality key through the target mapping first and the
delete file's original key
+ // name second. Initial-default and field identity remain bound to the
stable field ID.
+ Snapshot snapshot = createTableScan().snapshot();
+ while (snapshot != null) {
+ Integer schemaId = snapshot.schemaId();
+ if (schemaId != null) {
+ Schema historicalSchema = icebergTable.schemas().get(schemaId);
+ Preconditions.checkState(historicalSchema != null,
+ "Iceberg snapshot schema %s is absent from table
metadata", schemaId);
+ addHistoricalEqualityFields(fields, missingFieldIds,
historicalSchema);
+ }
+ Long parentId = snapshot.parentId();
+ snapshot = parentId == null ? null :
icebergTable.snapshot(parentId);
+ }
+ for (int index = schemaHistory.size() - 1; index >= 0; index--) {
+ addHistoricalEqualityFields(fields, missingFieldIds,
schemaHistory.get(index));
+ }
+ Preconditions.checkState(missingFieldIds.isEmpty(),
+ "Iceberg equality-delete fields are absent from schema
history: %s",
+ missingFieldIds);
+ return fields;
+ }
+
+ private List<Schema> getMetadataSchemaHistory() {
+ Preconditions.checkState(icebergTable instanceof HasTableOperations,
+ "Iceberg table does not expose metadata schema history: %s",
icebergTable.name());
+ return ((HasTableOperations)
icebergTable).operations().current().schemas();
+ }
+
+ /**
+ * Return only schemas that can describe files visible from the selected
target.
+ *
+ * <p>The query schema is included explicitly because a schema-only commit
does not create a
+ * snapshot. Other schemas are taken from the selected snapshot's parent
lineage and from
+ * cherry-picked source snapshots (including their ancestry), excluding
later main-branch and
+ * unrelated branch schemas from the rolling-upgrade fence. An empty
optional means snapshot
+ * expiration truncated any required lineage, so callers must
conservatively require current
+ * scan semantics.
+ */
+ @VisibleForTesting
+ Optional<List<Schema>> getRequiredFieldSchemaHistory(Schema scanSchema)
throws UserException {
+ List<Schema> schemas = new ArrayList<>();
+ Set<Integer> schemaIds = new HashSet<>();
+ schemas.add(scanSchema);
+ schemaIds.add(scanSchema.schemaId());
+
+ Snapshot selectedSnapshot = createTableScan().snapshot();
+ Deque<Snapshot> snapshots = new ArrayDeque<>();
+ if (selectedSnapshot != null) {
+ snapshots.add(selectedSnapshot);
+ }
+ Set<Long> visitedSnapshotIds = new HashSet<>();
+ while (!snapshots.isEmpty()) {
+ Snapshot snapshot = snapshots.removeFirst();
+ if (!visitedSnapshotIds.add(snapshot.snapshotId())) {
+ continue;
+ }
+ Integer schemaId = snapshot.schemaId();
+ if (schemaId != null && schemaIds.add(schemaId)) {
+ Schema lineageSchema = icebergTable.schemas().get(schemaId);
+ Preconditions.checkState(lineageSchema != null,
+ "Iceberg snapshot schema %s is absent from table
metadata", schemaId);
+ schemas.add(lineageSchema);
+ }
+ Long parentId = snapshot.parentId();
+ if (parentId != null) {
+ Snapshot parent = icebergTable.snapshot(parentId);
+ if (parent == null) {
+ return Optional.empty();
+ }
+ snapshots.addLast(parent);
+ }
+ String sourceSnapshotId =
+
snapshot.summary().get(SnapshotSummary.SOURCE_SNAPSHOT_ID_PROP);
+ if (sourceSnapshotId != null) {
+ Snapshot sourceSnapshot =
+
icebergTable.snapshot(Long.parseLong(sourceSnapshotId));
+ if (sourceSnapshot == null) {
+ return Optional.empty();
+ }
+ snapshots.addLast(sourceSnapshot);
+ }
+ }
+ return Optional.of(schemas);
+ }
+
+ private static void addHistoricalEqualityFields(List<NestedField> fields,
+ Set<Integer> missingFieldIds, Schema historicalSchema) {
+ Map<Integer, NestedField> historicalFields =
+ TypeUtil.indexById(historicalSchema.asStruct());
+ Set<Integer> selectedFieldIds = new HashSet<>();
+ for (Integer fieldId : missingFieldIds) {
+ NestedField field = historicalFields.get(fieldId);
+ if (field != null) {
+ Preconditions.checkState(field.type().isPrimitiveType(),
+ "Iceberg equality-delete field %s must be primitive",
fieldId);
+ selectedFieldIds.add(fieldId);
+ }
+ }
+ if (selectedFieldIds.isEmpty()) {
+ return;
+ }
+
+ Schema selectedSchema = TypeUtil.select(historicalSchema,
selectedFieldIds);
+ mergeHistoricalEqualityFields(fields, selectedSchema.columns());
+ missingFieldIds.removeAll(selectedFieldIds);
+ }
+
+ private static void mergeHistoricalEqualityFields(
+ List<NestedField> fields, List<NestedField> historicalFields) {
+ for (NestedField historicalField : historicalFields) {
+ int currentIndex = -1;
+ for (int index = 0; index < fields.size(); index++) {
+ if (fields.get(index).fieldId() == historicalField.fieldId()) {
+ currentIndex = index;
+ break;
+ }
+ }
+ if (currentIndex < 0) {
+ fields.add(historicalField);
+ continue;
+ }
+
+ NestedField currentField = fields.get(currentIndex);
+ Type mergedType = mergeHistoricalEqualityType(
+ currentField.type(), historicalField.type());
+ if (mergedType != currentField.type()) {
+ fields.set(currentIndex, Types.NestedField.from(currentField)
+ .ofType(mergedType)
+ .build());
+ }
+ }
+ }
+
+ private static Type mergeHistoricalEqualityType(Type currentType, Type
historicalType) {
+ Preconditions.checkState(currentType.typeId() ==
historicalType.typeId(),
+ "Iceberg equality-delete ancestor type changed from %s to %s",
+ historicalType, currentType);
+ switch (currentType.typeId()) {
+ case STRUCT:
+ List<NestedField> mergedFields =
+ new ArrayList<>(currentType.asStructType().fields());
+ mergeHistoricalEqualityFields(
+ mergedFields, historicalType.asStructType().fields());
+ if (mergedFields.equals(currentType.asStructType().fields())) {
+ return currentType;
+ }
+ return Types.StructType.of(mergedFields);
+ case LIST:
+ Types.ListType currentList = currentType.asListType();
+ Types.ListType historicalList = historicalType.asListType();
+ Preconditions.checkState(currentList.elementId() ==
historicalList.elementId(),
+ "Iceberg equality-delete list element id changed from
%s to %s",
+ historicalList.elementId(), currentList.elementId());
+ Type mergedElement = mergeHistoricalEqualityType(
+ currentList.elementType(),
historicalList.elementType());
+ if (mergedElement == currentList.elementType()) {
+ return currentType;
+ }
+ return currentList.isElementOptional()
+ ? Types.ListType.ofOptional(currentList.elementId(),
mergedElement)
+ : Types.ListType.ofRequired(currentList.elementId(),
mergedElement);
+ case MAP:
+ Types.MapType currentMap = currentType.asMapType();
+ Types.MapType historicalMap = historicalType.asMapType();
+ Preconditions.checkState(currentMap.keyId() ==
historicalMap.keyId()
+ && currentMap.valueId() ==
historicalMap.valueId(),
+ "Iceberg equality-delete map field ids changed from
(%s, %s) to (%s, %s)",
+ historicalMap.keyId(), historicalMap.valueId(),
+ currentMap.keyId(), currentMap.valueId());
+ Type mergedKey = mergeHistoricalEqualityType(
+ currentMap.keyType(), historicalMap.keyType());
+ Type mergedValue = mergeHistoricalEqualityType(
+ currentMap.valueType(), historicalMap.valueType());
+ if (mergedKey == currentMap.keyType()
+ && mergedValue == currentMap.valueType()) {
+ return currentType;
+ }
+ return currentMap.isValueOptional()
+ ? Types.MapType.ofOptional(
+ currentMap.keyId(), currentMap.valueId(),
+ mergedKey, mergedValue)
+ : Types.MapType.ofRequired(
+ currentMap.keyId(), currentMap.valueId(),
+ mergedKey, mergedValue);
+ default:
+ Preconditions.checkState(currentType.equals(historicalType),
+ "Iceberg equality-delete field type changed from %s to
%s",
+ historicalType, currentType);
+ return currentType;
+ }
+ }
+
+ @VisibleForTesting
+ static boolean requiresRecursiveInitialDefaultMaterialization(
+ Schema scanSchema, List<SlotDescriptor> projectedSlots) {
+ return requiresProjectedIcebergField(scanSchema, projectedSlots,
+ (field, isTopLevel) -> field.initialDefault() != null
+ && (!isTopLevel || field.type().isNestedType()));
+ }
+
+ @VisibleForTesting
+ static boolean requiresMissingRequiredFieldRejection(
+ Schema scanSchema, List<SlotDescriptor> projectedSlots,
+ Optional<List<Schema>> historicalSchemas) {
+ return !historicalSchemas.isPresent()
+ || requiresMissingRequiredFieldRejection(
+ scanSchema, projectedSlots, historicalSchemas.get());
+ }
+
+ @VisibleForTesting
+ static boolean requiresMissingRequiredFieldRejection(
+ Schema scanSchema, List<SlotDescriptor> projectedSlots,
+ List<Schema> historicalSchemas) {
+ Map<Integer, NestedField> fieldById =
TypeUtil.indexById(scanSchema.asStruct());
+ Map<Integer, Integer> parentById =
TypeUtil.indexParents(scanSchema.asStruct());
+ Set<Integer> collectionWrapperFieldIds = new HashSet<>();
+ collectCollectionWrapperFieldIds(scanSchema.asStruct(),
collectionWrapperFieldIds);
+ Set<Integer> potentiallyMissingRequiredFieldIds = new HashSet<>();
+ for (Schema historicalSchema : historicalSchemas) {
+ Map<Integer, NestedField> historicalFieldById =
+ TypeUtil.indexById(historicalSchema.asStruct());
+ for (NestedField field : fieldById.values()) {
+ NestedField historicalField =
historicalFieldById.get(field.fieldId());
+ if (historicalField != null) {
+ if (!collectionWrapperFieldIds.contains(field.fieldId())
+ && field.isRequired() && field.initialDefault() ==
null
+ && historicalField.isOptional()) {
+
potentiallyMissingRequiredFieldIds.add(field.fieldId());
+ }
+ continue;
+ }
+ NestedField highestMissingField = field;
+ Integer parentId = parentById.get(field.fieldId());
+ while (parentId != null &&
!historicalFieldById.containsKey(parentId)) {
+ highestMissingField =
Preconditions.checkNotNull(fieldById.get(parentId),
+ "Iceberg parent field %s is absent from scan
schema", parentId);
+ parentId = parentById.get(parentId);
+ }
+ // If the highest missing ancestor is optional, the old
physical subtree is NULL
+ // and no required descendant is materialized. A non-null
initial default is
+ // already covered by
requiresRecursiveInitialDefaultMaterialization().
+ if
(!collectionWrapperFieldIds.contains(highestMissingField.fieldId())
+ && highestMissingField.isRequired()
+ && highestMissingField.initialDefault() == null) {
+
potentiallyMissingRequiredFieldIds.add(highestMissingField.fieldId());
+ }
+ }
+ }
+ return requiresProjectedIcebergField(scanSchema, projectedSlots,
+ (field, isTopLevel) ->
potentiallyMissingRequiredFieldIds.contains(
+ field.fieldId()));
+ }
+
+ private static void collectCollectionWrapperFieldIds(
+ Type type, Set<Integer> collectionWrapperFieldIds) {
+ switch (type.typeId()) {
+ case STRUCT:
+ for (NestedField field : type.asStructType().fields()) {
+ collectCollectionWrapperFieldIds(field.type(),
collectionWrapperFieldIds);
+ }
+ break;
+ case LIST:
+ Types.ListType listType = (Types.ListType) type;
+ collectionWrapperFieldIds.add(listType.elementId());
+ collectCollectionWrapperFieldIds(
+ listType.elementType(), collectionWrapperFieldIds);
+ break;
+ case MAP:
+ Types.MapType mapType = (Types.MapType) type;
+ collectionWrapperFieldIds.add(mapType.keyId());
+ collectionWrapperFieldIds.add(mapType.valueId());
+ collectCollectionWrapperFieldIds(mapType.keyType(),
collectionWrapperFieldIds);
+ collectCollectionWrapperFieldIds(mapType.valueType(),
collectionWrapperFieldIds);
+ break;
+ default:
+ break;
+ }
+ }
+
+ private static boolean requiresProjectedIcebergField(
+ Schema scanSchema, List<SlotDescriptor> projectedSlots,
+ ProjectedFieldRequirement requirement) {
+ Map<Integer, NestedField> fieldById =
TypeUtil.indexById(scanSchema.asStruct());
+ Set<Integer> topLevelFieldIds = new HashSet<>();
+ for (NestedField field : scanSchema.columns()) {
+ topLevelFieldIds.add(field.fieldId());
+ }
+ for (SlotDescriptor slot : projectedSlots) {
+ Column column = slot.getColumn();
+ List<ColumnAccessPath> accessPaths = slot.getAllAccessPaths();
+ if (accessPaths != null && !accessPaths.isEmpty()) {
+ for (ColumnAccessPath accessPath : accessPaths) {
+ List<String> path = accessPath.getPath();
+ Preconditions.checkState(!path.isEmpty(),
+ "Iceberg column access path must not be empty");
+
Preconditions.checkState(matchesAccessPathComponent(column, path.get(0)),
+ "Iceberg access path root %s does not match column
%s", path.get(0),
+ column.getName());
+ if (requiresProjectedIcebergField(
+ column, path, 1, fieldById,
+ topLevelFieldIds.contains(column.getUniqueId()),
requirement)) {
+ return true;
+ }
+ }
+ } else if (requiresProjectedIcebergField(
+ column, slot.getType(), fieldById,
+ topLevelFieldIds.contains(column.getUniqueId()),
requirement)) {
+ return true;
+ }
+ }
+ return false;
+ }
+
+ private static boolean requiresProjectedIcebergField(
+ Column column, org.apache.doris.catalog.Type projectedType,
+ Map<Integer, NestedField> fieldById, boolean isTopLevel,
+ ProjectedFieldRequirement requirement) {
+ if (requiresIcebergField(column, fieldById, isTopLevel, requirement)) {
+ return true;
+ }
+ if (column.getChildren() == null) {
+ return false;
+ }
+ if (projectedType.isStructType()) {
+ for (StructField projectedField : ((StructType)
projectedType).getFields()) {
+ Column child = findChildByName(column,
projectedField.getName());
+ Preconditions.checkState(child != null,
+ "Projected Iceberg child %s is absent from column %s",
+ projectedField.getName(), column.getName());
+ if (requiresProjectedIcebergField(
+ child, projectedField.getType(), fieldById, false,
requirement)) {
+ return true;
+ }
+ }
+ } else if (projectedType.isArrayType()) {
+ Preconditions.checkState(column.getChildren().size() == 1,
+ "Iceberg array column %s must have one child",
column.getName());
+ if (requiresProjectedIcebergField(
+ column.getChildren().get(0), ((ArrayType)
projectedType).getItemType(),
+ fieldById, false, requirement)) {
+ return true;
+ }
+ } else if (projectedType.isMapType()) {
+ Preconditions.checkState(column.getChildren().size() == 2,
+ "Iceberg map column %s must have two children",
column.getName());
+ MapType mapType = (MapType) projectedType;
+ if (requiresProjectedIcebergField(
+ column.getChildren().get(0), mapType.getKeyType(),
fieldById, false,
+ requirement)
+ || requiresProjectedIcebergField(
+ column.getChildren().get(1),
mapType.getValueType(), fieldById, false,
+ requirement)) {
+ return true;
+ }
+ }
+ return false;
+ }
+
+ private static boolean requiresProjectedIcebergField(
+ Column column, List<String> path, int pathIndex,
+ Map<Integer, NestedField> fieldById, boolean isTopLevel,
+ ProjectedFieldRequirement requirement) {
+ if (requiresIcebergField(column, fieldById, isTopLevel, requirement)) {
+ return true;
+ }
+ if (pathIndex == path.size()) {
+ return requiresProjectedIcebergField(column, fieldById,
requirement);
+ }
+
+ String component = path.get(pathIndex);
+ if (AccessPathInfo.ACCESS_NULL.equals(component)
+ || AccessPathInfo.ACCESS_OFFSET.equals(component)) {
+ return false;
+ }
+ Preconditions.checkState(column.getChildren() != null,
+ "Iceberg access path continues below primitive column %s",
column.getName());
+
+ if (AccessPathInfo.ACCESS_ALL.equals(component)) {
+ if (column.getType().isArrayType()) {
+ Preconditions.checkState(column.getChildren().size() == 1,
+ "Iceberg array column %s must have one child",
column.getName());
+ return requiresProjectedIcebergField(
+ column.getChildren().get(0), path, pathIndex + 1,
fieldById, false,
+ requirement);
+ }
+ Preconditions.checkState(column.getType().isMapType(),
+ "Unexpected Iceberg access-all path below column %s",
column.getName());
+ Preconditions.checkState(column.getChildren().size() == 2,
+ "Iceberg map column %s must have two children",
column.getName());
+ Column key = column.getChildren().get(0);
+ // element_at(map, key) reads the complete key subtree, while any
path after '*'
+ // describes only the selected value subtree.
+ if (requiresIcebergField(key, fieldById, false, requirement)
+ || requiresProjectedIcebergField(key, fieldById,
requirement)) {
+ return true;
+ }
+ return requiresProjectedIcebergField(
+ column.getChildren().get(1), path, pathIndex + 1,
fieldById, false,
+ requirement);
+ }
+ if (column.getType().isMapType()) {
+ Preconditions.checkState(column.getChildren().size() == 2,
+ "Iceberg map column %s must have two children",
column.getName());
+ int childIndex;
+ if (AccessPathInfo.ACCESS_MAP_KEYS.equals(component)) {
+ childIndex = 0;
+ } else {
+
Preconditions.checkState(AccessPathInfo.ACCESS_MAP_VALUES.equals(component),
+ "Unexpected Iceberg map access path component %s",
component);
+ childIndex = 1;
+ }
+ return requiresProjectedIcebergField(
+ column.getChildren().get(childIndex), path, pathIndex + 1,
fieldById, false,
+ requirement);
+ }
+
+ Column child = findAccessPathChild(column, component);
+ Preconditions.checkState(child != null,
+ "Iceberg access path child %s is absent from column %s",
component,
+ column.getName());
+ return requiresProjectedIcebergField(
+ child, path, pathIndex + 1, fieldById, false, requirement);
+ }
+
+ private static boolean requiresProjectedIcebergField(
+ Column column, Map<Integer, NestedField> fieldById,
+ ProjectedFieldRequirement requirement) {
+ if (column.getChildren() == null) {
+ return false;
+ }
+ for (Column child : column.getChildren()) {
+ if (requiresIcebergField(child, fieldById, false, requirement)
+ || requiresProjectedIcebergField(child, fieldById,
requirement)) {
+ return true;
+ }
+ }
+ return false;
+ }
+
+ private static boolean requiresIcebergField(
+ Column column, Map<Integer, NestedField> fieldById, boolean
isTopLevel,
+ ProjectedFieldRequirement requirement) {
+ NestedField field = fieldById.get(column.getUniqueId());
+ return field != null && requirement.requires(field, isTopLevel);
+ }
+
+ private interface ProjectedFieldRequirement {
+ boolean requires(NestedField field, boolean isTopLevel);
+ }
+
+ /**
+ * Detect a reused name that current BEs resolve before an older sibling's
historical alias.
+ *
+ * <p>A smooth-upgrade source BE recognizes only the original semantics
marker and performs one
+ * ordered name/alias pass. If a sibling retains another sibling's current
name as an alias, the
+ * two BE generations can bind the same projected path to different field
IDs and types.
+ */
+ @VisibleForTesting
+ static boolean hasCurrentNameAliasCollision(
+ Schema schema, Optional<Map<Integer, List<String>>> nameMapping) {
+ return !getCurrentNameAliasCollisionFieldIds(schema,
nameMapping).isEmpty();
+ }
+
+ @VisibleForTesting
+ static void checkNameMappingBackendCompatibility(
+ Schema schema,
+ List<SlotDescriptor> projectedSlots,
+ Set<Integer> equalityDeleteFieldIds,
+ Optional<Map<Integer, List<String>>> nameMapping,
+ Iterable<Backend> backends) throws UserException {
+ Set<Integer> collisionFieldIds =
+ getCurrentNameAliasCollisionFieldIds(schema, nameMapping);
+ if (collisionFieldIds.isEmpty()) {
+ return;
+ }
+ boolean projectedCollision = requiresProjectedIcebergField(
+ schema, projectedSlots,
+ (field, isTopLevel) ->
collisionFieldIds.contains(field.fieldId()));
+ if (!projectedCollision && !equalityDeleteFieldIds.isEmpty()) {
+ Map<Integer, Integer> parentById =
TypeUtil.indexParents(schema.asStruct());
+ for (Integer equalityDeleteFieldId : equalityDeleteFieldIds) {
+ Integer fieldId = equalityDeleteFieldId;
+ while (fieldId != null) {
+ if (collisionFieldIds.contains(fieldId)) {
+ projectedCollision = true;
+ break;
+ }
+ fieldId = parentById.get(fieldId);
+ }
+ if (projectedCollision) {
+ break;
+ }
+ }
+ }
+ if (projectedCollision) {
+ checkCurrentIcebergScanSemanticsBackendCompatibility(backends);
+ }
+ }
+
+ private static Set<Integer> getCurrentNameAliasCollisionFieldIds(
+ Schema schema, Optional<Map<Integer, List<String>>> nameMapping) {
+ Set<Integer> collisionFieldIds = new HashSet<>();
+ if (nameMapping.isPresent()) {
+ collectCurrentNameAliasCollisionFieldIds(
+ schema.asStruct(), nameMapping.get(), collisionFieldIds);
+ }
+ return collisionFieldIds;
+ }
+
+ private static void collectCurrentNameAliasCollisionFieldIds(
+ Type type, Map<Integer, List<String>> nameMapping,
+ Set<Integer> collisionFieldIds) {
+ switch (type.typeId()) {
+ case STRUCT:
+ List<NestedField> fields = type.asStructType().fields();
+ for (NestedField field : fields) {
+ List<String> aliases =
+ nameMapping.getOrDefault(field.fieldId(),
Collections.emptyList());
+ for (String alias : aliases) {
+ for (NestedField sibling : fields) {
+ if (sibling.fieldId() != field.fieldId()
+ && sibling.name().equalsIgnoreCase(alias))
{
+ collisionFieldIds.add(field.fieldId());
+ collisionFieldIds.add(sibling.fieldId());
+ }
+ }
+ }
+ collectCurrentNameAliasCollisionFieldIds(
+ field.type(), nameMapping, collisionFieldIds);
+ }
+ return;
+ case LIST:
+ collectCurrentNameAliasCollisionFieldIds(
+ type.asListType().elementType(), nameMapping,
collisionFieldIds);
+ return;
+ case MAP:
+ collectCurrentNameAliasCollisionFieldIds(
+ type.asMapType().keyType(), nameMapping,
collisionFieldIds);
+ collectCurrentNameAliasCollisionFieldIds(
+ type.asMapType().valueType(), nameMapping,
collisionFieldIds);
+ return;
+ default:
+ return;
+ }
+ }
+
+ private static boolean matchesAccessPathComponent(Column column, String
component) {
+ return Integer.toString(column.getUniqueId()).equals(component)
+ || column.getName().equalsIgnoreCase(component);
+ }
+
+ private static Column findAccessPathChild(Column column, String component)
{
+ for (Column child : column.getChildren()) {
+ if (matchesAccessPathComponent(child, component)) {
+ return child;
+ }
+ }
+ return null;
+ }
+
+ private static Column findChildByName(Column column, String childName) {
+ for (Column child : column.getChildren()) {
+ if (child.getName().equalsIgnoreCase(childName)) {
+ return child;
+ }
+ }
+ return null;
+ }
+
+ @VisibleForTesting
+ Set<Integer> getEqualityDeleteFieldIdsForScan() throws UserException {
+ TableScan scan = createTableScan();
+ if (scan.snapshot() == null) {
+ return Collections.emptySet();
+ }
+ try {
+ return preExecutionAuthenticator.execute(
+ () -> loadEqualityDeleteFieldIds(scan));
+ } catch (Exception e) {
+ Optional<NotSupportedException> opt =
checkNotSupportedException(e);
+ if (opt.isPresent()) {
+ throw opt.get();
+ }
+ throw new UserException(ExceptionUtils.getRootCauseMessage(e), e);
+ }
+ }
+
+ /**
+ * Skip exhaustive delete-file planning when the exact snapshot summary
already proves that
+ * metadata-only COUNT(*) is safe. A usable count requires the summary's
equality-delete total
+ * to be zero, so no equality field IDs can affect this scan.
+ */
+ @VisibleForTesting
+ Set<Integer> getEqualityDeleteFieldIdsForPlanning() throws UserException {
+ if (prepareTableLevelSnapshotCount()) {
+ return Collections.emptySet();
+ }
+ return getEqualityDeleteFieldIdsForScan();
+ }
+
+ @VisibleForTesting
+ Set<Integer> loadEqualityDeleteFieldIds(TableScan scan) {
+ ConnectContext context = ConnectContext.get();
+ Preconditions.checkNotNull(context);
+ Preconditions.checkNotNull(context.getStatementContext());
+ List<FileScanTask> rewriteTasks =
+ context.getStatementContext().getIcebergRewriteFileScanTasks();
+ if (rewriteTasks != null) {
+ return collectEqualityDeleteFieldIdsFromTasks(rewriteTasks);
+ }
+ if (isBatchMode()) {
+ return loadEqualityDeleteFieldIdsFromDeleteManifests(scan);
+ }
+
+ List<FileScanTask> tasks = new ArrayList<>();
+ try (CloseableIterable<FileScanTask> plannedTasks =
planFileScanTaskWithoutReuse(scan)) {
+ for (FileScanTask task : plannedTasks) {
+ tasks.add(task);
+ }
+ } catch (IOException e) {
+ throw new RuntimeException("Failed to close Iceberg file scan
tasks", e);
+ }
+ preplannedFileScanTasks = tasks;
+ return collectEqualityDeleteFieldIdsFromTasks(tasks);
+ }
+
+ /**
+ * Load only delete-manifest metadata for batch scans.
+ *
+ * <p>Batch split planning must remain asynchronous, so this preflight
cannot consume
+ * {@link TableScan#planFiles()}. The target snapshot's equality-delete
total avoids manifest
+ * reads for the common no-delete case. Otherwise, partition and row
filters prune the delete
+ * manifests conservatively before collecting the field IDs needed by the
schema carrier.
+ */
+ @VisibleForTesting
+ Set<Integer> loadEqualityDeleteFieldIdsFromDeleteManifests(TableScan scan)
{
+ Snapshot snapshot = Preconditions.checkNotNull(scan.snapshot());
+ String totalEqualityDeletes =
+ snapshot.summary().get(SnapshotSummary.TOTAL_EQ_DELETES_PROP);
+ if (totalEqualityDeletes != null &&
Long.parseLong(totalEqualityDeletes) == 0) {
+ return Collections.emptySet();
+ }
+
+ Expression dataFilter = scan.filter();
+ Set<Integer> equalityDeleteFieldIds = new HashSet<>();
+ for (ManifestFile manifest :
snapshot.deleteManifests(icebergTable.io())) {
+ if (manifest.content() != ManifestContent.DELETES
+ || (!manifest.hasAddedFiles() &&
!manifest.hasExistingFiles())) {
+ continue;
+ }
+ PartitionSpec spec = Preconditions.checkNotNull(
+ icebergTable.specs().get(manifest.partitionSpecId()),
+ "Iceberg partition spec %s is absent from table metadata",
+ manifest.partitionSpecId());
+ Expression partitionFilter =
+ Projections.inclusive(spec, true).project(dataFilter);
+ if (!ManifestEvaluator.forPartitionFilter(
+ partitionFilter, spec, true).eval(manifest)) {
+ continue;
+ }
+ try (ManifestReader<DeleteFile> reader =
ManifestFiles.readDeleteManifest(
+ manifest, icebergTable.io(), icebergTable.specs())) {
+ ManifestReader<DeleteFile> filteredReader = reader
+ .filterRows(dataFilter)
+ .filterPartitions(partitionFilter)
+ .caseSensitive(true);
+ equalityDeleteFieldIds.addAll(
Review Comment:
[P1] Preserve data-sequence applicability in the batch preflight
`filteredReader` prunes delete rows and partitions, but this path has no
selected data files and never applies the equality-delete data-sequence rule
that `DeleteFileIndex` enforces during dispatch. For example, if every selected
data file is at sequence 30 while a retained equality delete at sequence 20
uses dropped field ID 9, this line adds ID 9 even though dispatch attaches that
delete to no task. On a large scan (batch mode is enabled by default), that can
spuriously reject every smooth-upgrade source BE; if the historical key has an
unsupported type such as `TIMESTAMP_NANO`, schema-carrier construction can fail
the scan outright.
This reintroduces the sequence-applicability half of the earlier
applicable-task issue through the new batch-only optimization: cross-partition
deletes are pruned now, but delete/data sequence and spec pairing are still
missing. Please make the preflight use the same delete-to-data applicability
contract as final task planning while preserving asynchronous dispatch, and add
a batch regression with a live delete older than every selected data file.
--
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]