apache / apache/iceberg-python

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

Abierto
#3,737 2 comentarios 1 reacción 0 asignados Ver en GitHub
Lenguaje dominante
Python
Estrellas
1.1k
Forks
581
Merge medio
1 d 17 h
PR fusionados (30 d)
78

Descripción

## 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.

Guía de contribución

No hay ninguna guía de contribución indexada para este repositorio

Línea de trabajo

Comienza con pyiceberg/io/pyarrow.py y la extracción de FileIO descrita en PR #3738; revisa el código existente de PyArrowFile y PyArrowFileIO, así como las pruebas relacionadas. Mueve únicamente la responsabilidad de FileIO, conserva las reexportaciones desde la ruta del módulo original y ejecuta la suite de pruebas existente para verificar que el comportamiento y los imports no hayan cambiado.

Escrito por el modelo de indexación a partir del texto del issue.

Evaluación

Stack tecnológico
python
Área
data-engineering, databases
Tipo de issue
Refactorización
Dificultad
5/5
Tiempo estimado
Más de una semana
Estado de actividad
Activo
Claridad
Bastante claro
Aptitud para principiantes
35/100

Recibe los nuevos issues en tu correo

Un resumen breve de issues de GitHub para principiantes.