From 288f9f26859381a64fe58603fdeb15bdb935c504 Mon Sep 17 00:00:00 2001 From: Han You Date: Wed, 12 Aug 2026 11:33:45 -0500 Subject: [PATCH 1/2] Flink: Avoid per-record Set allocation in the dynamic sink forward check DynamicRecordProcessor.collect calls isForwardEligible for every record. That call resolves to DynamicSinkUtil.resolveEqualityFieldNames, which for a record with no user-supplied equality fields returns Schema.identifierFieldNames(), which uses Java Stream API to allocate and build a new HashSet. This HashSet is immediately discarded after isEmpty(). So each record on the hottest and commonest path (a table without identifier fields) allocates a set purely to ask whether it is empty. Ask the same question without materialising anything: test the user-supplied equality fields directly, then consult Schema.identifierFieldIds(), which Schema already memoizes as an ImmutableSet. Behaviour is unchanged and test still covers every branch. --- .../flink/sink/dynamic/DynamicRecordProcessor.java | 8 ++++++-- .../flink/sink/dynamic/DynamicRecordProcessor.java | 8 ++++++-- .../flink/sink/dynamic/DynamicRecordProcessor.java | 8 ++++++-- 3 files changed, 18 insertions(+), 6 deletions(-) diff --git a/flink/v1.20/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/DynamicRecordProcessor.java b/flink/v1.20/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/DynamicRecordProcessor.java index ba2d816ae9ea..1c378edb9876 100644 --- a/flink/v1.20/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/DynamicRecordProcessor.java +++ b/flink/v1.20/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/DynamicRecordProcessor.java @@ -19,6 +19,7 @@ package org.apache.iceberg.flink.sink.dynamic; import java.util.Map; +import java.util.Set; import org.apache.flink.annotation.Internal; import org.apache.flink.annotation.VisibleForTesting; import org.apache.flink.api.common.functions.OpenContext; @@ -233,8 +234,11 @@ public void close() { */ @VisibleForTesting static boolean isForwardEligible(DynamicRecord data) { + Set equalityFields = data.equalityFields(); + // Checks identifier field ids directly rather than resolving equality field names, which would + // allocate per record. return data.distributionMode() == null - && DynamicSinkUtil.resolveEqualityFieldNames(data.equalityFields(), data.schema()) - .isEmpty(); + && (equalityFields == null || equalityFields.isEmpty()) + && data.schema().identifierFieldIds().isEmpty(); } } diff --git a/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/DynamicRecordProcessor.java b/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/DynamicRecordProcessor.java index ba2d816ae9ea..1c378edb9876 100644 --- a/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/DynamicRecordProcessor.java +++ b/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/DynamicRecordProcessor.java @@ -19,6 +19,7 @@ package org.apache.iceberg.flink.sink.dynamic; import java.util.Map; +import java.util.Set; import org.apache.flink.annotation.Internal; import org.apache.flink.annotation.VisibleForTesting; import org.apache.flink.api.common.functions.OpenContext; @@ -233,8 +234,11 @@ public void close() { */ @VisibleForTesting static boolean isForwardEligible(DynamicRecord data) { + Set equalityFields = data.equalityFields(); + // Checks identifier field ids directly rather than resolving equality field names, which would + // allocate per record. return data.distributionMode() == null - && DynamicSinkUtil.resolveEqualityFieldNames(data.equalityFields(), data.schema()) - .isEmpty(); + && (equalityFields == null || equalityFields.isEmpty()) + && data.schema().identifierFieldIds().isEmpty(); } } diff --git a/flink/v2.1/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/DynamicRecordProcessor.java b/flink/v2.1/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/DynamicRecordProcessor.java index ba2d816ae9ea..1c378edb9876 100644 --- a/flink/v2.1/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/DynamicRecordProcessor.java +++ b/flink/v2.1/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/DynamicRecordProcessor.java @@ -19,6 +19,7 @@ package org.apache.iceberg.flink.sink.dynamic; import java.util.Map; +import java.util.Set; import org.apache.flink.annotation.Internal; import org.apache.flink.annotation.VisibleForTesting; import org.apache.flink.api.common.functions.OpenContext; @@ -233,8 +234,11 @@ public void close() { */ @VisibleForTesting static boolean isForwardEligible(DynamicRecord data) { + Set equalityFields = data.equalityFields(); + // Checks identifier field ids directly rather than resolving equality field names, which would + // allocate per record. return data.distributionMode() == null - && DynamicSinkUtil.resolveEqualityFieldNames(data.equalityFields(), data.schema()) - .isEmpty(); + && (equalityFields == null || equalityFields.isEmpty()) + && data.schema().identifierFieldIds().isEmpty(); } } From 2d61a1b6e20bf003f80cffabf2a4ed97e8571761 Mon Sep 17 00:00:00 2001 From: Han You Date: Thu, 13 Aug 2026 09:59:33 -0500 Subject: [PATCH 2/2] PR comments: move the fix to resolveEqualityFieldNames --- api/src/main/java/org/apache/iceberg/Schema.java | 15 ++++++++++++--- .../java/org/apache/iceberg/SchemaUpdate.java | 2 +- .../sink/dynamic/DynamicRecordProcessor.java | 8 ++------ .../sink/dynamic/DynamicRecordProcessor.java | 8 ++------ .../sink/dynamic/DynamicRecordProcessor.java | 8 ++------ 5 files changed, 19 insertions(+), 22 deletions(-) diff --git a/api/src/main/java/org/apache/iceberg/Schema.java b/api/src/main/java/org/apache/iceberg/Schema.java index 3e59998be476..ed4d1d60e901 100644 --- a/api/src/main/java/org/apache/iceberg/Schema.java +++ b/api/src/main/java/org/apache/iceberg/Schema.java @@ -81,6 +81,7 @@ public class Schema implements Serializable { private transient Map> idToAccessor = null; private transient Map idToName = null; private transient Set identifierFieldIdSet = null; + private transient Set identifierFieldNameSet = null; private final transient Map idsToReassigned; private final transient Map idsToOriginal; @@ -254,6 +255,16 @@ private Set lazyIdentifierFieldIdSet() { return identifierFieldIdSet; } + private Set lazyIdentifierFieldNameSet() { + if (identifierFieldNameSet == null) { + this.identifierFieldNameSet = + lazyIdentifierFieldIdSet().stream() + .map(id -> lazyIdToName().get(id)) + .collect(ImmutableSet.toImmutableSet()); + } + return identifierFieldNameSet; + } + /** * Returns the schema ID for this schema. * @@ -331,9 +342,7 @@ public Set identifierFieldIds() { /** Returns the set of identifier field names. */ public Set identifierFieldNames() { - return identifierFieldIds().stream() - .map(id -> lazyIdToName().get(id)) - .collect(Collectors.toSet()); + return lazyIdentifierFieldNameSet(); } /** diff --git a/core/src/main/java/org/apache/iceberg/SchemaUpdate.java b/core/src/main/java/org/apache/iceberg/SchemaUpdate.java index 1fa6ebbe8fef..ed26bdb82446 100644 --- a/core/src/main/java/org/apache/iceberg/SchemaUpdate.java +++ b/core/src/main/java/org/apache/iceberg/SchemaUpdate.java @@ -87,7 +87,7 @@ private SchemaUpdate(TableOperations ops, TableMetadata base, Schema schema, int this.schema = schema; this.lastColumnId = lastColumnId; this.idToParent = Maps.newHashMap(TypeUtil.indexParents(schema.asStruct())); - this.identifierFieldNames = schema.identifierFieldNames(); + this.identifierFieldNames = Sets.newHashSet(schema.identifierFieldNames()); } @Override diff --git a/flink/v1.20/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/DynamicRecordProcessor.java b/flink/v1.20/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/DynamicRecordProcessor.java index 1c378edb9876..ba2d816ae9ea 100644 --- a/flink/v1.20/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/DynamicRecordProcessor.java +++ b/flink/v1.20/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/DynamicRecordProcessor.java @@ -19,7 +19,6 @@ package org.apache.iceberg.flink.sink.dynamic; import java.util.Map; -import java.util.Set; import org.apache.flink.annotation.Internal; import org.apache.flink.annotation.VisibleForTesting; import org.apache.flink.api.common.functions.OpenContext; @@ -234,11 +233,8 @@ public void close() { */ @VisibleForTesting static boolean isForwardEligible(DynamicRecord data) { - Set equalityFields = data.equalityFields(); - // Checks identifier field ids directly rather than resolving equality field names, which would - // allocate per record. return data.distributionMode() == null - && (equalityFields == null || equalityFields.isEmpty()) - && data.schema().identifierFieldIds().isEmpty(); + && DynamicSinkUtil.resolveEqualityFieldNames(data.equalityFields(), data.schema()) + .isEmpty(); } } diff --git a/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/DynamicRecordProcessor.java b/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/DynamicRecordProcessor.java index 1c378edb9876..ba2d816ae9ea 100644 --- a/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/DynamicRecordProcessor.java +++ b/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/DynamicRecordProcessor.java @@ -19,7 +19,6 @@ package org.apache.iceberg.flink.sink.dynamic; import java.util.Map; -import java.util.Set; import org.apache.flink.annotation.Internal; import org.apache.flink.annotation.VisibleForTesting; import org.apache.flink.api.common.functions.OpenContext; @@ -234,11 +233,8 @@ public void close() { */ @VisibleForTesting static boolean isForwardEligible(DynamicRecord data) { - Set equalityFields = data.equalityFields(); - // Checks identifier field ids directly rather than resolving equality field names, which would - // allocate per record. return data.distributionMode() == null - && (equalityFields == null || equalityFields.isEmpty()) - && data.schema().identifierFieldIds().isEmpty(); + && DynamicSinkUtil.resolveEqualityFieldNames(data.equalityFields(), data.schema()) + .isEmpty(); } } diff --git a/flink/v2.1/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/DynamicRecordProcessor.java b/flink/v2.1/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/DynamicRecordProcessor.java index 1c378edb9876..ba2d816ae9ea 100644 --- a/flink/v2.1/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/DynamicRecordProcessor.java +++ b/flink/v2.1/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/DynamicRecordProcessor.java @@ -19,7 +19,6 @@ package org.apache.iceberg.flink.sink.dynamic; import java.util.Map; -import java.util.Set; import org.apache.flink.annotation.Internal; import org.apache.flink.annotation.VisibleForTesting; import org.apache.flink.api.common.functions.OpenContext; @@ -234,11 +233,8 @@ public void close() { */ @VisibleForTesting static boolean isForwardEligible(DynamicRecord data) { - Set equalityFields = data.equalityFields(); - // Checks identifier field ids directly rather than resolving equality field names, which would - // allocate per record. return data.distributionMode() == null - && (equalityFields == null || equalityFields.isEmpty()) - && data.schema().identifierFieldIds().isEmpty(); + && DynamicSinkUtil.resolveEqualityFieldNames(data.equalityFields(), data.schema()) + .isEmpty(); } }