apache / apache/iceberg-python

Decompose io/pyarrow.py into focused modules to enable pluggable compute engines

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

説明

## Summary

`pyiceberg/io/pyarrow.py` is a 3,100+ line monolith that handles six unrelated concerns: filesystem I/O, schema conversion, expression translation, scan/read orchestration, write logic, and Parquet statistics. This makes it difficult to test individual components, extend behavior, or substitute alternative engines for specific operations.

This issue proposes an incremental decomposition - a series of small, independently-reviewable refactoring PRs that split the file by concern while maintaining full backward compatibility via re-exports. The end goal is clean seam points where bounded-memory compute engines (DataFusion, etc.) can be introduced for operations that currently OOM on large data.

## Motivation

Several open issues depend on bounded-memory compute that PyArrow's kernel library cannot provide:

- Equality delete resolution (#1210, #3270) - requires anti-join with spill
- Sort-on-write (#271) - requires external merge sort
- Data compaction (#1092) - requires sort + join + rewrite pipeline
- CoW deletes on large files - full file materialization causes OOM

A previous attempt to deliver all of this at once (#3715, PR #3716) was rejected for being too large to review. This issue takes the opposite approach: decompose first, add capabilities later.

## Approach

### Phase 1: Split the monolith (pure refactoring)

Extract each concern into its own module under `pyiceberg/io/`. The original `pyarrow.py` becomes a thin re-export shim so all existing imports continue to work.

| PR | Extraction | Approximate scope |
|----|-----------|-------------------|
| A | FileIO (`PyArrowFile`, `PyArrowFileIO`) - #3738 | ~700 lines |
| B | Schema conversion (`schema_to_pyarrow`, `pyarrow_to_schema`, visitors) | ~900 lines |
| C | Expression translation (`expression_to_pyarrow`, `_ConvertToArrowExpression`) | ~300 lines |
| D | Statistics (`StatsAggregator`, `PyArrowStatisticsCollector`, `ParquetFormatWriter`) | ~500 lines |
| E | Write path (`write_file`, `_dataframe_to_data_files`, partitioning, bin packing) | ~1200 lines |
| F | Scan/Read (`ArrowScan`, `_task_to_record_batches`, delete resolution) | ~300 lines |

Each PR:
- Moves code, does not change behavior
- Re-exports from the original module path
- All existing tests pass unchanged
- No new dependencies

These PRs are largely independent of each other (no strict ordering required).

### Phase 2: Introduce a compute protocol

Once concerns are separated, introduce a thin `ComputeEngine` protocol for the operations that benefit from bounded-memory execution:

```python
class ComputeEngine(Protocol):
def filter_batches(self, batches, expr, schema) -> Iterator[RecordBatch]: ...
def sort_batches(self, batches, sort_order, schema) -> Iterator[RecordBatch]: ...
def anti_join(self, left, right, keys) -> Iterator[RecordBatch]: ...
```

The default implementation delegates to the existing PyArrow code. No behavior change, just an indirection point.

### Phase 3: DataFusion as optional compute engine

With the protocol in place, a `DataFusionComputeEngine` implementation slots in as an optional extra. Each capability (equality delete resolution, sort-on-write, etc.) is its own PR wiring the protocol into the specific code path.

## What this is NOT

- Not a rewrite. Phase 1 is purely moving existing code into new files.
- Not adding DataFusion as a hard dependency. It remains an optional extra.
- Not changing the public API. All existing imports and behaviors are preserved.

## Prior art / references

- #3715 / PR #3716: Previous pluggable backend attempt (rejected as too large)
- Community sync discussion (June 30, 2026): established that read/write/compute should be separable
- #3554: Original DataFusion integration proposal
- #271: Sort-on-write (requires external merge sort)
- #1210, #3270: Equality delete resolution
- #1092: Data compaction

---

I plan to start with PR A (FileIO extraction, #3738) as a proof of concept for the approach. Feedback on the overall direction is welcome before I proceed further.

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

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

調査の方向性

pyiceberg/io/pyarrow.py と、PR #3738 で説明されている FileIO の抽出から始めます。既存の PyArrowFile および PyArrowFileIO のコードと関連するテストを確認します。FileIO に関する部分だけを移動し、元のモジュールパスからの再エクスポートを維持して、既存のテストスイートを実行し、動作とインポートに変更がないことを確認します。

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

評価

技術スタック
python
領域
data-engineering, databases
issue の種類
リファクタリング
難易度
5/5
見積もり時間
1週間以上
活発さ
活発
明瞭さ
おおむね明確
初心者へのやさしさ
35/100

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

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