This is an automated email from the ASF dual-hosted git repository.
mxm pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/iceberg.git
The following commit(s) were added to refs/heads/main by this push:
new d490240b17 Flink: Backport fix for RANGE distribution ignoring
user-specified equality fields to 1.20 and 2.0 (#17459)
d490240b17 is described below
commit d490240b17be248df3fa8ed78a5d1867b3738dea
Author: Eunbin Son <[email protected]>
AuthorDate: Sat Aug 1 13:42:16 2026 +0900
Flink: Backport fix for RANGE distribution ignoring user-specified equality
fields to 1.20 and 2.0 (#17459)
---
.../flink/sink/dynamic/HashKeyGenerator.java | 2 +-
.../flink/sink/dynamic/TestHashKeyGenerator.java | 71 ++++++++++++++++++++++
.../flink/sink/dynamic/HashKeyGenerator.java | 2 +-
.../flink/sink/dynamic/TestHashKeyGenerator.java | 71 ++++++++++++++++++++++
4 files changed, 144 insertions(+), 2 deletions(-)
diff --git
a/flink/v1.20/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/HashKeyGenerator.java
b/flink/v1.20/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/HashKeyGenerator.java
index 03541ed596..0fd6bed9e1 100644
---
a/flink/v1.20/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/HashKeyGenerator.java
+++
b/flink/v1.20/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/HashKeyGenerator.java
@@ -173,7 +173,7 @@ class HashKeyGenerator {
}
case RANGE:
- if (schema.identifierFieldIds().isEmpty()) {
+ if (equalityFields.isEmpty()) {
LOG.warn(
"{}: Fallback to use 'none' distribution mode, because there are
no equality fields set "
+ "and {}='range' is not supported yet in flink",
diff --git
a/flink/v1.20/flink/src/test/java/org/apache/iceberg/flink/sink/dynamic/TestHashKeyGenerator.java
b/flink/v1.20/flink/src/test/java/org/apache/iceberg/flink/sink/dynamic/TestHashKeyGenerator.java
index 4181f00d0a..7a50ca4356 100644
---
a/flink/v1.20/flink/src/test/java/org/apache/iceberg/flink/sink/dynamic/TestHashKeyGenerator.java
+++
b/flink/v1.20/flink/src/test/java/org/apache/iceberg/flink/sink/dynamic/TestHashKeyGenerator.java
@@ -164,6 +164,77 @@ class TestHashKeyGenerator {
assertThat(getSubTaskId(writeKey3, writeParallelism,
maxWriteParallelism)).isEqualTo(0);
}
+ @Test
+ void testRangeModeWithEqualityFields() throws Exception {
+ int writeParallelism = 2;
+ int maxWriteParallelism = 8;
+ HashKeyGenerator generator = new HashKeyGenerator(16, maxWriteParallelism);
+ PartitionSpec unpartitioned = PartitionSpec.unpartitioned();
+
+ GenericRowData row1 = GenericRowData.of(1, StringData.fromString("foo"));
+ GenericRowData row2 = GenericRowData.of(1, StringData.fromString("bar"));
+ GenericRowData row3 = GenericRowData.of(2, StringData.fromString("baz"));
+ // SCHEMA has no identifier fields, so the equality fields only come from
the record
+ Set<String> equalityColumns = Collections.singleton("id");
+
+ int writeKey1 =
+ getWriteKey(
+ generator,
+ unpartitioned,
+ DistributionMode.RANGE,
+ writeParallelism,
+ equalityColumns,
+ row1);
+ int writeKey2 =
+ getWriteKey(
+ generator,
+ unpartitioned,
+ DistributionMode.RANGE,
+ writeParallelism,
+ equalityColumns,
+ row2);
+ int writeKey3 =
+ getWriteKey(
+ generator,
+ unpartitioned,
+ DistributionMode.RANGE,
+ writeParallelism,
+ equalityColumns,
+ row3);
+
+ assertThat(writeKey1).isEqualTo(writeKey2);
+ assertThat(writeKey2).isNotEqualTo(writeKey3);
+ }
+
+ @Test
+ void testRangeModeFallsBackToDistributionModeNone() throws Exception {
+ int writeParallelism = 2;
+ int maxWriteParallelism = 8;
+ HashKeyGenerator generator = new HashKeyGenerator(16, maxWriteParallelism);
+ Schema noIdSchema = new Schema(Types.NestedField.required(1, "x",
Types.StringType.get()));
+ PartitionSpec unpartitioned = PartitionSpec.unpartitioned();
+
+ DynamicRecord record =
+ new DynamicRecord(
+ TABLE_IDENTIFIER,
+ BRANCH,
+ noIdSchema,
+ GenericRowData.of(StringData.fromString("v")),
+ unpartitioned,
+ DistributionMode.RANGE,
+ writeParallelism);
+
+ int writeKey1 = generator.generateKey(record);
+ int writeKey2 = generator.generateKey(record);
+ int writeKey3 = generator.generateKey(record);
+ assertThat(writeKey1).isNotEqualTo(writeKey2);
+ assertThat(writeKey3).isEqualTo(writeKey1);
+
+ assertThat(getSubTaskId(writeKey1, writeParallelism,
maxWriteParallelism)).isEqualTo(1);
+ assertThat(getSubTaskId(writeKey2, writeParallelism,
maxWriteParallelism)).isEqualTo(0);
+ assertThat(getSubTaskId(writeKey3, writeParallelism,
maxWriteParallelism)).isEqualTo(1);
+ }
+
@Test
void testHashModeWithPartitionFieldAndEqualityField() throws Exception {
int writeParallelism = 2;
diff --git
a/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/HashKeyGenerator.java
b/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/HashKeyGenerator.java
index 03541ed596..0fd6bed9e1 100644
---
a/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/HashKeyGenerator.java
+++
b/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/HashKeyGenerator.java
@@ -173,7 +173,7 @@ class HashKeyGenerator {
}
case RANGE:
- if (schema.identifierFieldIds().isEmpty()) {
+ if (equalityFields.isEmpty()) {
LOG.warn(
"{}: Fallback to use 'none' distribution mode, because there are
no equality fields set "
+ "and {}='range' is not supported yet in flink",
diff --git
a/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/sink/dynamic/TestHashKeyGenerator.java
b/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/sink/dynamic/TestHashKeyGenerator.java
index 4181f00d0a..7a50ca4356 100644
---
a/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/sink/dynamic/TestHashKeyGenerator.java
+++
b/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/sink/dynamic/TestHashKeyGenerator.java
@@ -164,6 +164,77 @@ class TestHashKeyGenerator {
assertThat(getSubTaskId(writeKey3, writeParallelism,
maxWriteParallelism)).isEqualTo(0);
}
+ @Test
+ void testRangeModeWithEqualityFields() throws Exception {
+ int writeParallelism = 2;
+ int maxWriteParallelism = 8;
+ HashKeyGenerator generator = new HashKeyGenerator(16, maxWriteParallelism);
+ PartitionSpec unpartitioned = PartitionSpec.unpartitioned();
+
+ GenericRowData row1 = GenericRowData.of(1, StringData.fromString("foo"));
+ GenericRowData row2 = GenericRowData.of(1, StringData.fromString("bar"));
+ GenericRowData row3 = GenericRowData.of(2, StringData.fromString("baz"));
+ // SCHEMA has no identifier fields, so the equality fields only come from
the record
+ Set<String> equalityColumns = Collections.singleton("id");
+
+ int writeKey1 =
+ getWriteKey(
+ generator,
+ unpartitioned,
+ DistributionMode.RANGE,
+ writeParallelism,
+ equalityColumns,
+ row1);
+ int writeKey2 =
+ getWriteKey(
+ generator,
+ unpartitioned,
+ DistributionMode.RANGE,
+ writeParallelism,
+ equalityColumns,
+ row2);
+ int writeKey3 =
+ getWriteKey(
+ generator,
+ unpartitioned,
+ DistributionMode.RANGE,
+ writeParallelism,
+ equalityColumns,
+ row3);
+
+ assertThat(writeKey1).isEqualTo(writeKey2);
+ assertThat(writeKey2).isNotEqualTo(writeKey3);
+ }
+
+ @Test
+ void testRangeModeFallsBackToDistributionModeNone() throws Exception {
+ int writeParallelism = 2;
+ int maxWriteParallelism = 8;
+ HashKeyGenerator generator = new HashKeyGenerator(16, maxWriteParallelism);
+ Schema noIdSchema = new Schema(Types.NestedField.required(1, "x",
Types.StringType.get()));
+ PartitionSpec unpartitioned = PartitionSpec.unpartitioned();
+
+ DynamicRecord record =
+ new DynamicRecord(
+ TABLE_IDENTIFIER,
+ BRANCH,
+ noIdSchema,
+ GenericRowData.of(StringData.fromString("v")),
+ unpartitioned,
+ DistributionMode.RANGE,
+ writeParallelism);
+
+ int writeKey1 = generator.generateKey(record);
+ int writeKey2 = generator.generateKey(record);
+ int writeKey3 = generator.generateKey(record);
+ assertThat(writeKey1).isNotEqualTo(writeKey2);
+ assertThat(writeKey3).isEqualTo(writeKey1);
+
+ assertThat(getSubTaskId(writeKey1, writeParallelism,
maxWriteParallelism)).isEqualTo(1);
+ assertThat(getSubTaskId(writeKey2, writeParallelism,
maxWriteParallelism)).isEqualTo(0);
+ assertThat(getSubTaskId(writeKey3, writeParallelism,
maxWriteParallelism)).isEqualTo(1);
+ }
+
@Test
void testHashModeWithPartitionFieldAndEqualityField() throws Exception {
int writeParallelism = 2;