apache / apache/iceberg-python
Allow user to define a subset of columns for update detection in UPSERT
- 主要語言
- Python
- 星號
- 1.1k
- 分支
- 581
- 平均合併
- 1 天 13 小時
- 30 天內合併 PR
- 76
描述
### 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 處理,讓選取的非主鍵欄位決定更新偵測,並確認當要求的欄位與可用欄位完全沒有交集時所採用的行為。
由索引模型根據 Issue 內容生成。
評估
- 技術堆疊
- python
- 領域
- data-engineering, databases
- Issue 類型
- 功能
- 難度
- 3/5
- 預估耗時
- 1-2 天
- 活躍度
- 冷清
- 描述清晰度
- 基本清楚
- 新手友好度
- 55/100