Hacktoberfest 2026: the issues maintainers tagged for October, open and beginner-friendly. Browse Hacktoberfest issues

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

Open
#3,737 2 comments 1 reaction 0 assignees View on GitHub

Nobody has claimed this yet.

Assessment

Difficulty
5/5
Estimated time
Over a week
Newbie friendliness
35/100
Issue type
Refactor
Clarity
Mostly clear
Activity status
Active
Tech stack
python

Research direction

Start with pyiceberg/io/pyarrow.py and the FileIO extraction described in PR #3738; review the existing PyArrowFile and PyArrowFileIO code and related tests. Move only the FileIO concern, preserve re-exports from the original module path, and run the existing test suite to verify unchanged behavior and imports.

Written by the indexing model from the issue text.

Description

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:

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.

Dominant language
Python
Stars
1.1k
Forks
589
Avg merge
2d 2h
Merged PRs (30d)
70

Contributor guide

No contributing guide indexed for this repository

First steps

  1. Read the whole issue, then the project's contributing guide.
  2. Comment on the issue to say you are picking it up — it saves two people doing the same work.
  3. Fork the repository and make your change on a branch.
  4. Open a pull request that references the issue number.

More from apache/iceberg-python

All issues in apache/iceberg-python

Similar issues

More Python issues

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.