qzyu999 opened a new issue, #3737:
URL: https://github.com/apache/iceberg-python/issues/3737
## 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 - requires external merge sort
- Data compaction (#1092) - requires sort + join + rewrite pipeline
- CoW deletes on large files - full file materialization causes OOM
An earlier attempt to address this (pluggable backend PR) was rejected as
too large. 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`) | ~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
- Community sync discussion (June 30, 2026): established that
read/write/compute should be separable
- #3554: Original DataFusion integration proposal
- #1210, #3270: Equality delete resolution
- #1092: Data compaction
---
I plan to start with PR A (FileIO extraction) as a proof of concept for the
approach. Feedback on the overall direction is welcome before I proceed further.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]