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]

Reply via email to