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:
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
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.
Summary
pyiceberg/io/pyarrow.pyis 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:
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 originalpyarrow.pybecomes a thin re-export shim so all existing imports continue to work.PyArrowFile,PyArrowFileIO) - #3738schema_to_pyarrow,pyarrow_to_schema, visitors)expression_to_pyarrow,_ConvertToArrowExpression)StatsAggregator,PyArrowStatisticsCollector,ParquetFormatWriter)write_file,_dataframe_to_data_files, partitioning, bin packing)ArrowScan,_task_to_record_batches, delete resolution)Each PR:
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
ComputeEngineprotocol for the operations that benefit from bounded-memory execution: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
DataFusionComputeEngineimplementation 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
Prior art / references
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.