Decompose io/pyarrow.py into focused modules to enable pluggable compute engines
还没有人认领这个 Issue。
评估
- 难度
- 5/5
- 预计耗时
- 一周以上
- 新手友好度
- 35/100
- Issue 类型
- 重构
- 描述清晰度
- 基本清楚
- 活跃度
- 活跃
- 技术栈
- python
调研方向
从 pyiceberg/io/pyarrow.py 以及 PR #3738 中描述的 FileIO 提取开始;检查现有的 PyArrowFile 和 PyArrowFileIO 代码及相关测试。仅移动 FileIO 相关职责,保留从原始模块路径进行的重新导出,并运行现有测试套件,以验证行为和导入保持不变。
由索引模型根据 Issue 内容生成。
描述
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:
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.
- 主要语言
- Python
- 星标
- 1.1k
- 派生
- 589
- 平均合并
- 2 天 2 小时
- 30 天内合并 PR
- 70
贡献指南
这个仓库没有索引到贡献指南
从这里开始
- 先读完整个 Issue,再读项目的贡献指南。
- 在 Issue 下留言说明你要接手 —— 这能避免两个人做同样的事。
- Fork 仓库,在一个分支上完成修改。
- 提交 Pull Request,并在描述里引用这个 Issue 编号。
apache/iceberg-python 的其他 Issue
-
kind:bug
难度 1/5 1 小时以内 新手友好度 92/100
apache/iceberg-python#4006 ·
-
难度 2/5 1-3 小时 新手友好度 78/100
apache/iceberg-python#3996 ·
-
bug
难度 2/5 1-3 小时 新手友好度 72/100
apache/iceberg-python#3979 ·
-
难度 2/5 1-3 小时 新手友好度 78/100
apache/iceberg-python#3885 ·
-
[Bug] PyArrowFileIO fails to propagate s3.ssl.ca-cert to pyarrow.fs.S3FileSystem tls_ca_file_path 未关闭
难度 2/5 1-3 小时 新手友好度 76/100
apache/iceberg-python#3866 · 1 条评论 ·
查看 apache/iceberg-python 的全部 Issue
相似的 Issue
-
essnmx good first issue
难度 1/5 1 小时以内 新手友好度 95/100
-
难度 2/5 1-3 小时 新手友好度 65/100
syfoud/Simulated_Scepter#174 ·
-
难度 2/5 1-3 小时 新手友好度 75/100
Giskard-AI/giskard-oss#2840 · 1 条评论 ·
-
A claim comment carrying the issue number is silently declined while the workflow reports success 未关闭area: repo bug perceived difficulty: 2
难度 2/5 1-3 小时 新手友好度 70/100
-
难度 2/5 1-3 小时 新手友好度 75/100
yeti-platform/yeti#1380 ·