skip to content

Why does an Apache Hudi upsert create a duplicate when a record's partition value changes?

level: seniorimportance: should knowfreq 32%

answer

  1. the lookup happens somewhere, not everywhere
  2. a changed partition value is a new address
  3. nothing asked about the old address
  4. one index flavour searches the whole table
  5. or stop partitioning on something that changes

basics

~20 s

A non-global Hudi index enforces key uniqueness only within a partition. When the partition value changes, the lookup in the new partition finds nothing, so the record is inserted there while the original row stays in its old partition.

solid answer

~50 s

Hudi's default index types are partition-scoped: the index maps a record key to a file group *within* a partition path. When a record's partition column changes value — a customer moves country, an order's status column is the partition key — the upsert computes the new partition path, looks the key up there, finds no match, and inserts. Nothing deletes the old row, so the same key now exists in two partitions. The fixes are to use a global index (`GLOBAL_BLOOM`, `GLOBAL_SIMPLE`, or the record-level index, which is global by construction), which searches the whole table for the key; and to configure whether that global index moves the record — the update-partition-path setting decides between deleting the old row and writing the new one, or keeping the record in its original partition. The tradeoff is cost: a global lookup spans every partition.

go deeper

for a junior

Know that Hudi looks up a record key to decide update versus insert, and that the lookup is normally scoped to one partition, so a changed partition value looks like a brand-new record.

for a middle

Explain that uniqueness is enforced per partition path plus record key by default, and name the global index types that widen the search to the whole table.

for a senior

Diagnose it from the metadata columns, weigh the global-index cost against redesigning the partitioning, and know that whether the record moves or stays put is an explicit setting rather than a safe default.

for a principal

Treat mutable partition columns as a design rule for the platform: forbid them where possible, and where global uniqueness is a real requirement decide which index carries it and who pays the tagging cost.

## The mechanism Hudi identifies a row by its record key, but the index that finds the row is, by default, scoped to a partition. Conceptually the index answers "within partition P, which file group holds key K?" — not "where in the table is key K?". Uniqueness is therefore enforced per (partition path, record key), and that is a deliberate design choice: restricting the search to one partition is what keeps tagging affordable on a large table. Now take a record whose partition column is mutable. The write computes the partition path from the incoming record's values, so the changed value produces a *different* partition path. The upsert looks key K up in the new partition, finds nothing, and classifies the record as an insert. The row in the old partition is untouched, because nothing ever asked about it. A snapshot query now returns two rows for one business entity, in different partitions, with different `_hoodie_partition_path` values. This is the single most common cause of "Hudi produced duplicates", and it is not a bug — it is the documented consequence of a partition-scoped index meeting a mutable partition column. ## Detecting it The symptom is a count that exceeds the distinct count of the record key. Because every base file carries `_hoodie_record_key` and `_hoodie_partition_path` as metadata columns, the diagnosis is a group-by over those two columns: any key appearing under more than one partition path is a moved record that was never removed. If duplicates instead share a partition path, the cause is different — usually a bulk_insert over live data, or a record key that is not actually unique. ## The fixes **Use a global index.** `GLOBAL_BLOOM` and `GLOBAL_SIMPLE` search the whole table for the key rather than a single partition, and the record-level index in the metadata table is global by construction because it maps a key to both its partition and its file group. With a global index, the moved record's existing location is found, and Hudi can act on it. **Decide what "found in another partition" should mean.** A global index exposes a setting — the update-partition-path option on the bloom and simple global indexes — that chooses between two behaviours: delete the record from its old partition and write it into the new one, so the row *moves*; or ignore the new partition value and update the record where it already lives, so the partition column effectively becomes immutable after first insert. Both are legitimate; the wrong one silently produces either duplicates or a partition column that never changes. Pick deliberately and write it down. **Or remove the mutability.** The cheapest fix is often to stop partitioning by a mutable attribute. Partition by an immutable event or creation date, keep the mutable attribute as an ordinary column, and the problem disappears along with the cost of a global index. This is the usual senior recommendation: the duplicate is a modelling symptom, and partitioning by a value that changes is a modelling mistake in any table format. ## The cost of going global A global index cannot prune by partition, so every key lookup has the whole table as its search space. On a bloom-based global index that means testing far more files; the effect on job runtime can be severe on wide tables. The record-level index softens this because its lookup is a mapping probe rather than a file scan, which is why it is the usual recommendation when global uniqueness is a genuine requirement on a large table. There is also a correctness benefit to going global that is worth naming: a global index guarantees the record key is unique across the entire table, which is what people usually assume a "primary key" means. With a non-global index, two rows with the same key in different partitions are perfectly legal. ## Cleaning up existing duplicates Switching the index does not retroactively merge rows that are already duplicated. Existing duplicates must be removed explicitly — typically by identifying, per record key, which partition holds the surviving version, issuing deletes for the stale copies, and only then switching the pipeline to a global index so the problem does not recur. Doing it in the other order lets the pipeline keep producing new duplicates while you clean. ## What interviewers look for The strong answer names the scope of the index as the cause rather than blaming the merge logic or the precombine field, offers both the global-index fix and the immutable-partition-column fix, states the cost of the global lookup, and mentions that the choice of whether the record moves or stays is an explicit configuration, not a default you can assume.

  • How would you detect this class of duplicate in an existing Hudi table?
    Group by the `_hoodie_record_key` metadata column and count distinct `_hoodie_partition_path` values. Any key spanning more than one partition path is a record that moved and was never removed from its original location. Duplicates within a single partition path point at a different cause, usually a bulk_insert run over live data or a record key that is not genuinely unique.
  • What does a global Hudi index cost compared with a partition-scoped one?
    It loses partition pruning: every key lookup has the whole table as its search space rather than one partition. With a global bloom or simple index that means far more candidate files per key and much longer tagging as the table grows. The record-level index mitigates it because its lookup is a mapping probe rather than a scan of data files.
  • Does switching to a global index remove duplicates that already exist?
    No. The index changes how future writes resolve a key; it does not revisit stored rows. Existing duplicates must be cleaned explicitly — pick the surviving copy per record key, issue deletes for the stale ones, and switch the pipeline to the global index first so the cleanup is not undone by the next batch.

saying these in an interview costs you the question

  • Blames the precombine field for the duplicate rows
  • Assumes the record key is unique across the whole table by default
  • Thinks switching to a global index retroactively merges existing duplicates
  • Says a global index costs nothing extra because it is still an index
  • Never questions partitioning by a column whose value changes

context