apache / apache/iceberg-python
Pluggable Backend Interface with DataFusion for Bounded-Memory Compute
- 主要语言
- Python
- 星标
- 1.1k
- 派生
- 581
- 平均合并
- 1 天 17 小时
- 30 天内合并 PR
- 78
描述
## Summary
PyIceberg uses PyArrow as its sole execution engine. PyArrow is a kernel library with no memory management, no spill-to-disk, and no join operators. Operations that process more data than available memory (CoW deletes, equality delete resolution, scan planning for heavily-deleted tables, sorted writes) crash with OOM errors.
This issue tracks introducing a pluggable backend interface (`ReadBackend`, `WriteBackend`, `ComputeBackend` protocols) and integrating Apache DataFusion as the first bounded-memory compute backend.
## Problem
| Operation | Current Status | OOM Pattern |
|-----------|---------------|-------------|
| Equality delete reads | Hard `ValueError` | Anti-join requires all delete keys in memory |
| CoW delete (large files) | OOMs | Materializes entire Parquet file into RAM |
| Scan planning (>100K deletes) | OOMs | All delete entries in Python dict |
| Sort-on-write | Not implemented | Full sort before write |
| Positional deletes (millions) | OOMs | Python set of positions |
Tables written by Flink (which uses equality deletes) are completely unreadable by PyIceberg today.
## Solution
1. **Pluggable interface**: `ReadBackend`, `WriteBackend`, `ComputeBackend` protocols that decouple PyIceberg from PyArrow
2. **DataFusion integration**: Bounded-memory sort, join, and filter with spill-to-disk via `datafusion-python`
3. **Migration**: All existing data operations route through the interface with zero API changes
## Deliverables
- [ ] Equality delete resolution (NEW): tables with equality deletes can now be read
- [ ] CoW delete/overwrite streaming (FIX): statistics short-circuit + two-pass streaming
- [ ] Positional delete resolution (IMPROVED): bounded-memory for large delete sets
- [ ] Sort-on-write (NEW): external merge sort when DataFusion installed
- [ ] Bounded-memory scan planning (NEW): for tables with >100K delete files
## Related Issues
- #1210 - Support reading equality delete files
- #3270 - Equality Delete support
- #3554 - Integrate DataFusion as execution engine
## Acceptance Criteria
- All existing tests pass without `datafusion` installed (no regression)
- Tables with equality deletes return correct results
- CoW delete on 2GB+ files completes without OOM (with DataFusion)
- Sort-on-write produces sorted files when table has sort order and DataFusion installed
- Property-based tests verify PyArrow and DataFusion backends produce identical output
贡献指南
这个仓库没有索引到贡献指南
调研方向
未指定实现文件或测试。先阅读相关 issue #1210、#3270 和 #3554,然后确定所提议的 ReadBackend、WriteBackend 和 ComputeBackend 协议以及 DataFusion 集成的范围。完成标准是实现列出的删除、流式处理、排序和扫描规划行为,在没有 DataFusion 时不发生回归,并且 PyArrow 和 DataFusion 的结果等价。
由索引模型根据 Issue 内容生成。
评估
- 技术栈
- python
- 领域
- backend, data-engineering, databases
- Issue 类型
- 功能
- 难度
- 5/5
- 预计耗时
- 一周以上
- 活跃度
- 冷清
- 描述清晰度
- 需要澄清
- 新手友好度
- 25/100