Skip to content

[Data] DSv2 Reorg [6/6]: Generalize OnlineBinPacker to any file format with one FileChunk type for run metadata - #66577

Open
AarryaSaraf wants to merge 4 commits into
ray-project:masterfrom
AarryaSaraf:dsv2-chunk-run-metadata
Open

AarryaSaraf wants to merge 4 commits into
ray-project:masterfrom
AarryaSaraf:dsv2-chunk-run-metadata

Conversation

@AarryaSaraf

@AarryaSaraf AarryaSaraf commented Sep 29, 2026 •

Copy link
Copy Markdown
Contributor

Description

Last PR of the DSv2 Reorg series. Three changes, two types:

  1. OnlineBinPacker packs any file, not just Parquet. It read and wrote ParquetRowGroupChunkMetadata, so only the Parquet footer indexer could feed it. It now works on a format-neutral row: "this file's units 3-5, N rows, S bytes", where a unit is whatever a format's reader can scan on its own (a Parquet row group today, a CSV byte range or a Lance fragment tomorrow). The packer only sums size_bytes against its budget and never asks what the number measures; the indexer that writes the row decides. The packing arithmetic is unchanged.
  2. One dataclass for "a run of units" instead of four look-alikes. RowGroupInfo (footer reader and coalescer), ParquetRowGroupChunkMetadata (manifest row), the copied fields on the packer's BinItem, and the create_chunk_metadata helper that built the row all carried the same ids / rows / bytes / per-unit tuples under different names. They collapse into one frozen FileChunk in interfaces/file_manifest.py, with to_metadata() / from_metadata() for the manifest's __file_chunk_metadata column. (Named UnitRun in the first draft; renamed per review because it read too close to ReadUnit.)
  3. One shape for "a file and its size" from listing to packer. FileInfo(path, size) already existed in interfaces/file_indexer.py, but the listing helpers passed (path, size) tuples, OrderedFileResult carried the path twice, the footer actor took list[tuple[str, int]], its result FileChunks had its own path/size, and the packer's BinItem had its own path/size_bytes. All of them now carry a FileInfo. An indexer author emits one type; a reviewer checks no conversions.

Cost of the unified types. Measured with pympler.asizeof on Python 3.12, same 45-character path for old and new. A single-row-group run is 888 bytes as a FileChunk versus 920 as a RowGroupInfo; a three-group coalesced run is 1120 versus 1072. A BinItem is 1888 bytes versus 1416 for the old flat BinItem(path, size_bytes, run): the extra ~470 bytes is the FileInfo instance it now points at instead of copying its two fields. BinItems are transient packing state, bounded by the open bins, so this is a few hundred bytes per in-flight item. The manifest's Arrow column is unchanged at 36 bytes per row, because the struct has the same fields under new names.

What changes

  1. interfaces/file_manifest.py. New FileChunk(unit_ids, num_rows, size_bytes, fully_matched=True, unit_sizes=(), unit_rows=()) with to_metadata() / from_metadata(); __post_init__ asserts a per-unit breakdown is empty or one entry per unit. create_chunk_metadata and its TypeVar are removed; the empty ChunkMetadata TypedDict becomes a Dict[str, Any] alias, since a row's chunk metadata is now always None or FileChunk.to_metadata().
  2. formats/parquet/. parquet_footer_types.py is deleted: RowGroupInfo and ParquetRowGroupChunkMetadata are gone, and the one survivor, FileChunks → ChunkedFile(file: FileInfo, row_groups: Tuple[FileChunk, ...]), moves next to the FooterReader that builds it. footer_reader.py builds one FileChunk per row group and _read_and_chunk / read_footers take FileInfo; parquet_row_group_coalescing.py merges with dataclasses.replace; footer_file_indexer.py emits run.to_metadata() via _chunked_file_to_manifest with size_bytes = the run's projection-scoped uncompressed size, as before. parquet_file_reader.py reads unit_ids and size_bytes, and drops the "uncompressed_size" in chunk key check, since a non-None row always has the keys now. Docstrings in parquet_file_chunking_utils.py follow.
  3. common/indexing_utils.py, common/non_sampling_file_indexer.py. PathContents.files and _get_file_infos produce FileInfo; OrderedFileResult loses its duplicate file_path.
  4. common/online_bin_packer.py. BinItem(file: FileInfo, run: Optional[FileChunk]), with size_bytes a property (the run's bytes, else the file size); run is None is a whole-file row. Splitting a run slices the FileChunk; a sealed bin merges a file's runs back into one FileChunk row. Bin.total_uncompressed_size → total_bytes. Two behaviour details: a whole-file row now comes out of the packer as a whole-file row again (it used to come out claiming row_group_ids=(0,)), and every packed row keeps the on-disk __file_size the listing gave it (a chunk row used to carry the run's bytes there; nothing reads __file_size after the packer).
  5. README. One paragraph under "Adding a format": an indexer that lists FileChunk rows can hand them to OnlineBinPacker as they are.
  6. Tests. test_online_bin_packer.py builds rows directly and imports nothing from formats/; its coalesce_row_groups cases move verbatim to a new test_parquet_row_group_coalescing.py, since they are Parquet's; two tests added (whole-file rows, on-disk size kept). test_chunk_metadata.py becomes a FileChunk round-trip and invariant test. FileInfo / renamed fields in test_footer_reader.py, test_read_units.py, test_synthesized_columns.py, test_parquet_datasource_v2.py, test_file_indexer.py and tests/datasource/test_read_parquet_v2.py.

Additional information

Not a duplicate: no open PR generalises the packer's row or unifies the listing types. #66121 and #65887 both touch these files under the pre-#66497 paths and will need a rebase either way; #66121 adds a second size field to the old Parquet row; on this PR it needs no new field, since size_bytes is whatever the indexer records and the decoded-or-fallback choice belongs in the footer indexer's _chunked_file_to_manifest. Its footer-indexer changes will see ChunkedFile.file.path and read_footers(list[FileInfo]) rather than tuples.

Checks run locally from this branch (macOS, ~/ray/.venv):

pytest -q python/ray/data/tests/unit/datasource_v2 \
          python/ray/data/tests/datasource_v2/test_footer_reader.py  # 193 passed
pytest -q python/ray/data/tests/datasource/test_read_parquet_v2.py   # 37 passed

ruff, black and pydoclint at the pre-commit pins are clean on every changed file; pyrefly 0.51.0 with pyarrow==24.0.0 reports no errors in the changed files. grep -r datasource_v2.formats interfaces/ common/ is empty.

Written with Claude Code; I reviewed every changed line and ran the tests above.

Signed-off-by: Aarrya aarrya.saraf@anyscale.com

🤖 Generated with Claude Code

…t with one UnitRun type for run metadata

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Signed-off-by: Aarrya <aarrya.saraf@anyscale.com>
@AarryaSaraf
AarryaSaraf force-pushed the dsv2-chunk-run-metadata branch from 79aa7e3 to b84d4b3 Compare September 29, 2026 17:14
@AarryaSaraf
AarryaSaraf marked this pull request as ready for review September 29, 2026 19:51
@AarryaSaraf
AarryaSaraf requested a review from a team as a code owner September 29, 2026 19:51
@AarryaSaraf AarryaSaraf added data Ray Data-related issues go add ONLY when ready to merge, run all tests labels Sep 29, 2026

@gemini-code-assist gemini-code-assist Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Code Review

This pull request refactors the DataSourceV2 chunk metadata representation by replacing Parquet-specific metadata classes with a generic UnitRun dataclass, allowing the OnlineBinPacker to partition files of any format based on generic read units and byte sizes. The review feedback highlights opportunities to improve robustness: specifically, adding validation in OnlineBinPacker to ensure unit_sizes matches unit_ids before splitting, using safe .get() fallbacks during UnitRun deserialization to prevent KeyError or TypeError, and verifying the presence of "num_rows" in the chunk metadata before estimating batch sizes.

Comment on lines +330 to +331
if self._split_coalesced and run is not None and len(run.unit_ids) > 1:
return list(run.unit_sizes)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

If run.unit_sizes is empty or its length does not match run.unit_ids (for example, after runs are merged in _bin_to_manifest where unit_sizes is not populated), _units will return an empty list or mismatched sizes. This can lead to incorrect behavior or crashes in _slice_bin_item when slicing. We should ensure that we only split the run if run.unit_sizes is populated and matches the length of run.unit_ids.

Suggested change
if self._split_coalesced and run is not None and len(run.unit_ids) > 1:
return list(run.unit_sizes)
if (
self._split_coalesced
and run is not None
and len(run.unit_ids) > 1
and run.unit_sizes
and len(run.unit_sizes) == len(run.unit_ids)
):
return list(run.unit_sizes)
Comment on lines +50 to +57
return cls(
unit_ids=tuple(metadata["unit_ids"]),
num_rows=int(metadata["num_rows"]),
size_bytes=int(metadata["size_bytes"]),
fully_matched=bool(metadata["fully_matched"]),
unit_sizes=tuple(metadata["unit_sizes"]),
unit_rows=tuple(metadata["unit_rows"]),
)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

Accessing keys directly from metadata can raise a KeyError or TypeError (e.g., if unit_sizes or unit_rows is None or missing due to schema evolution or null values in Arrow/Parquet manifests). Using .get() with safe fallbacks makes deserialization much more robust.

Suggested change
return cls(
unit_ids=tuple(metadata["unit_ids"]),
num_rows=int(metadata["num_rows"]),
size_bytes=int(metadata["size_bytes"]),
fully_matched=bool(metadata["fully_matched"]),
unit_sizes=tuple(metadata["unit_sizes"]),
unit_rows=tuple(metadata["unit_rows"]),
)
return cls(
unit_ids=tuple(metadata.get("unit_ids") or ()),
num_rows=int(metadata.get("num_rows", 0)),
size_bytes=int(metadata.get("size_bytes", 0)),
fully_matched=bool(metadata.get("fully_matched", True)),
unit_sizes=tuple(metadata.get("unit_sizes") or ()),
unit_rows=tuple(metadata.get("unit_rows") or ()),
)
(md for md in manifest.file_chunk_metadatas if md is not None), None
)
if chunk is not None and "uncompressed_size" in chunk:
if chunk is not None and "size_bytes" in chunk:

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

We should check if "num_rows" is also present in chunk before accessing it on line 333 to avoid a potential KeyError if a custom or future metadata format omits it.

Suggested change
if chunk is not None and "size_bytes" in chunk:
if chunk is not None and "size_bytes" in chunk and "num_rows" in chunk:
…onstruction and drop the stale chunk key check in the Parquet reader

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Signed-off-by: Aarrya <aarrya.saraf@anyscale.com>


@dataclass(frozen=True)
class UnitRun:

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Can we come up wiht a different term? ReadUnit and UnitRun are pretty similar and confusing

# ---------------------------------------------------------------------------


@pytest.mark.parametrize(

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why are these tests being removed?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Theyre just being moved, verbatim

@@ -1,113 +1,61 @@
from typing import Any
from typing import Any, Optional

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Ensure that in this file we're not removing any of the relevant bin packing tests, but modifying them to cater to the new abstractions

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Theyre just being moved, verbatim

rg_count=count,
rg_sizes=sizes if count > 1 else (),
rg_rows=rows if count > 1 else (),
size_bytes=sum(sizes),

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Do we not have the prefix sums to compute this cheaply?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why would that be cheaper?

…o the one path+size shape from listing to packer

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Signed-off-by: Aarrya <aarrya.saraf@anyscale.com>
@AarryaSaraf AarryaSaraf changed the title [Data] DSv2 Reorg [6/6]: Generalize OnlineBinPacker to any file format with one UnitRun type for run metadata Oct 1, 2026

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

data Ray Data-related issues go add ONLY when ready to merge, run all tests

2 participants