kumarpritam863 opened a new pull request, #18296:
URL: https://github.com/apache/iceberg/pull/18296

   ## Summary
   Fixes an incorrect `IllegalStateException` in the Flink **dynamic sink** 
when an upsert job uses an
   equality (primary) key that is a **nested** field and also partitions by 
that same nested field.
   
   ## Symptom
   Upsert with primary key `user.id` and partition `bucket(16, user.id)` fails 
at runtime:
   
   ```
   java.lang.IllegalStateException: <table>: In 'hash' distribution mode with 
equality fields set,
   partition field '...: user.id_bucket: bucket[16](...)' should be included in 
equality fields: '[]'
       at 
org.apache.iceberg.flink.sink.dynamic.HashKeyGenerator.getKeySelector(HashKeyGenerator.java:...)
       at 
org.apache.iceberg.flink.sink.dynamic.HashKeyGenerator.generateKey(...)
       at 
org.apache.iceberg.flink.sink.dynamic.DynamicRecordProcessor.collect(...)
   ```
   
   ## Root cause
   `HashKeyGenerator.getKeySelector` (HASH branch) verifies every partition 
**source** field is one of the
   equality fields, but matched **by name**:
   
   ```java
   NestedField sourceField = schema.findField(partitionField.sourceId());
   Preconditions.checkState(
       sourceField != null && equalityFields.contains(sourceField.name()), ...);
   ```
   
   For a nested source, `sourceField.name()` is the **simple** leaf name 
(`"id"`), while the equality
   field set holds **dotted** names (`"user.id"`) — so they never match and a 
valid nested
   partition/equality field is wrongly rejected. The `'[]'` in the message is 
itself misleading: it prints
   the subset of *top-level* columns matching the equality names (empty for a 
nested key), not the
   equality set.
   
   ## Fix
   Resolve the partition **source** to its **fully-qualified column name** — 
from the spec's own schema —
   and check membership in the equality-field names:
   
   ```java
   for (PartitionField partitionField : spec.fields()) {
     String sourceName = 
spec.schema().findColumnName(partitionField.sourceId());
     Preconditions.checkState(
         sourceName != null && equalityFields.contains(sourceName),
         "%s: In 'hash' distribution mode with equality fields set, partition 
field '%s' "
             + "should be included in equality fields: '%s'",
         tableName, partitionField, equalityFields);
   }
   ```
   
   **Why fully-qualified name and not field ID** : field IDs are only meaningful
   within the schema that assigned them, and the final table schema may not be 
known at this point — the
   schema available here can differ from `spec.schema()`, so an ID comparison 
across the two is not
   reliable. Column names are stable across schema evolution, and the 
equality-field names are already
   fully qualified (`user.id`), so resolving the source's fully-qualified name 
via `spec.schema()` and
   comparing names matches directly. This also fixes the original nested-field 
bug (simple leaf name
   `id` vs dotted `user.id`).
   
   The existing guard is **preserved**: a partition source that is genuinely 
*not* an equality field is
   still rejected. That guard is intentional — Iceberg equality deletes are 
partition-scoped, so
   partitioning by a non-key column under HASH + upsert would corrupt dedup (a 
partition-changing update
   leaves a stale row in the old partition).
   
   ## Tests
   Added to `TestHashKeyGenerator` (all Flink versions):
   - `testHashDistributionModeWithNestedEqualityAndPartitionField` — a 
partition sourced from a nested
     field that is also the equality field is accepted and distributes by 
partition value; the check
     matches by fully-qualified name (`user.id`), not the source's simple leaf 
name (`id`).
   - `testHashDistributionModePartitionFieldMustBeEqualityField` — a 
non-equality partition source is
     still rejected, with the corrected message.
   
   ## Versions
   Applied identically to Flink **v1.20, v2.1, v2.2, v2.3**.
   


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