apache / apache/iceberg-python

Allow user to define a subset of columns for update detection in UPSERT

オープン
#3,598 コメント 1 件 リアクション 0 件 担当者 0 名 GitHub で見る
主要言語
Python
スター
1.1k
フォーク
581
平均マージ
1日 17時間
マージ済み PR(30日)
78

説明

### Feature Request / Improvement

Currently, when detecting which rows should be updated in the `upsert`, all non-primary key columns are iterated over, converted to a Python type, and compared. This can be time-consuming and memory-consuming (e.g., with complex columns containing JSON data in `struct`, `list`, ...).

If the user knows that a change has occurred to specific columns, it would be a huge performance improvement to just iterate over these columns. Or if the user has a column that implies that any change has occurred (e.g., a hash of the data)

For example:
Hash column: if the user has the possibility to create a column containing a hash for each row, then upsert has to look only at this column when detecting changes. This way, pyiceberg doesn't have to convert all columns to Python type and compare them.

Proposition:

Update the function `get_rows_to_update` to also accept an optional parameter `difference_cols`. Update the upsert methods with this parameter and pass it to the `get_rows_to_update`.

In case there is no intersection between non-primary key columns and `difference_cols`, `pyiceberg` can either raise and error or it can fall-back to the default behaviour (iterating over all non-PK columns).

Usage:

```python
from pyiceberg.schema import Schema
from pyiceberg.types import IntegerType, NestedField, StringType

import pyarrow as pa

schema = Schema(
NestedField(1, "city", StringType(), required=True),
NestedField(2, "inhabitants", IntegerType(), required=True),
# Mark City as the identifier field, also known as the primary-key
identifier_field_ids=[1]
)

tbl = catalog.create_table("default.cities", schema=schema)

arrow_schema = pa.schema(
[
pa.field("city", pa.string(), nullable=False),
pa.field("inhabitants", pa.int32(), nullable=False),
]
)

# Write some data
df = pa.Table.from_pylist(
[
{"city": "Amsterdam", "inhabitants": 921402},
{"city": "San Francisco", "inhabitants": 808988},
{"city": "Drachten", "inhabitants": 45019},
{"city": "Paris", "inhabitants": 2103000},
],
schema=arrow_schema
)
tbl.append(df)

df = pa.Table.from_pylist(
[
# Will be updated, the inhabitants has been updated
{"city": "Drachten", "inhabitants": 45505},

# New row, will be inserted
{"city": "Berlin", "inhabitants": 3432000},

# Ignored, already exists in the table
{"city": "Paris", "inhabitants": 2103000},
],
schema=arrow_schema
)
upd = tbl.upsert(df, difference_cols=["inhabitants"])
```

I have already prepared how it can look in my fork 59d18337a61b106566c2e6a432b7d4899ca7f334.
If the proposition is accepted, I can prepare the PR.

コントリビューションガイド

このリポジトリのコントリビューションガイドは索引されていません

調査の方向性

まず get_rows_to_update と upsert メソッドを調べ、次に提案されている動作をコミット 59d18337a61b106566c2e6a432b7d4899ca7f334 と比較します。difference_cols を任意で処理できるように追加し、選択された主キー以外のカラムによって更新の検出を行えるようにします。また、要求されたカラムが利用可能なカラムと1つも重ならない場合に選択する動作を確認します。

索引モデルが issue の本文から書いたものです。

評価

技術スタック
python
領域
data-engineering, databases
issue の種類
機能追加
難易度
3/5
見積もり時間
1〜2日
活発さ
静か
明瞭さ
おおむね明確
初心者へのやさしさ
55/100

新しい issue をメールで受け取る

初心者向けの GitHub issue を短くまとめたダイジェスト。