Skip to content

fix(upsert): make the match filter flat for composite keys - #4027

Open
azwanzuharimi wants to merge 3 commits into
apache:mainfrom
azwanzuharimi:upsert-flat-match-filter
Open

azwanzuharimi wants to merge 3 commits into
apache:mainfrom
azwanzuharimi:upsert-flat-match-filter

Conversation

@azwanzuharimi

@azwanzuharimi azwanzuharimi commented Sep 29, 2026 •

Copy link
Copy Markdown

Closes #3508

Rationale for this change

Table.upsert with two or more join columns builds one Or(And(EqualTo, ...)) disjunct per key. PyArrow flattens that chain and recurses over it. From about 1,000 keys the C++ stack overflows and the process dies without a Python traceback. All three uses of the filter crash: the scan for matching rows, the insert filter and the overwrite filter. A balanced Or tree does not help, because PyArrow flattens it again.

This change keeps every filter flat:

  • create_match_filter returns one In per join column. For a composite key this filter can match more rows than the keys in the source.
  • A new upsert_util.exclude_keys does the exact key match in Arrow with an anti join. The insert path uses it instead of an Arrow expression filter.
  • The overwrite filter is the flat filter on the updated keys. Before the overwrite, the transaction scans the rows that this filter removes. After the overwrite it appends back the rows that were not updated.

Related work: #3509 groups keys by prefix and stays exact, but it still emits one disjunct per distinct prefix. #3420 bounds the scan filter, but the insert and overwrite paths keep the exact filter.

Results, measured locally on macOS with pyarrow 25.0.1. Target table 100,000 rows, upsert 20,000 rows (10,000 updates and 10,000 inserts). The check reads the table back and compares every row. rc 138 is a stack overflow crash. The script is available on request.

case main #3509 #3420 this PR
a: 2 columns, one with ~10 values rc 138 pass 0.3 s rc 138 pass 0.3 s
b: 2 columns, all values unique rc 138 rc 138 rc 138 pass 0.5 s
c: 3 columns, all values unique rc 138 rc 138 rc 138 pass 0.7 s
d: null in a key column TypeError TypeError TypeError TypeError
e: keys match nothing pass pass pass pass
f: partitioned (day, bucket[16]) rc 138 rc 138 rc 138 pass 32 s
control: 1 column pass pass pass pass

Case d rejects null keys on every branch. This PR does not change that.

Cost for a composite key: the overwrite reads the affected files one extra time. Rows that share a key column value with an updated key, but are not updated, are rewritten unchanged. In the worst case the filter matches every combination of the updated key values, for example every date times every id. Single column keys are not changed.

Are these changes tested?

Yes. tests/table/test_upsert.py adds:

  • test_create_match_filter_composite_key_is_flat
  • test_upsert_composite_key_keeps_rows_outside_source_keys: every key column value of the source exists in the target, but only some key tuples do. This guards against deleting too much.
  • test_upsert_composite_key_large_batch: 5,000 key tuples. On main this crashes pytest with exit code 132.
  • test_upsert_composite_key_all_source_keys_matched_across_files: the first data file removes every source key from the insert set, a later file still matches the filter.
  • test_exclude_keys_rejects_reserved_column_name

All 4239 unit tests pass.

Are there any user-facing changes?

Yes. An upsert on a composite key with more than about 1,000 keys no longer crashes. An upsert on a composite key can add one extra APPEND snapshot when the overwrite filter removes rows that are not updated. No API change.

An upsert with two or more join columns built one Or disjunct per key.
PyArrow flattens that chain and overflows its stack from about 1,000 keys.
The scan, the insert filter and the overwrite filter all crashed.

Build one In per join column instead. Do the exact key match in Arrow
with an anti join. Append back the rows that the overwrite filter removes
but does not update.

Closes apache#3508
@azwanzuharimi
azwanzuharimi force-pushed the upsert-flat-match-filter branch from 0c60294 to fbcc6fd Compare September 29, 2026 08:34
An empty table gave the index column the null type and the anti join
rejected it. The insert path reaches this when one scan batch removes
every source key and a later batch still matches the filter.
Match the check that get_rows_to_update already does, so a join column
named __index gives a clear error instead of an Arrow join failure.
@azwanzuharimi azwanzuharimi changed the title Keep the upsert match filter flat for composite keys fix(upsert): keep the match filter flat for composite keys Sep 29, 2026
@azwanzuharimi azwanzuharimi changed the title fix(upsert): keep the match filter flat for composite keys fix(upsert): make the match filter flat for composite keys Sep 29, 2026

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Segfault on large multi-column Iceberg upserts

1 participant