[Data] DSv2 Reorg [6/6]: Generalize OnlineBinPacker to any file format with one FileChunk type for run metadata - #66577
AarryaSaraf wants to merge 4 commits into
Conversation
…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>
79aa7e3 to
b84d4b3
Compare
There was a problem hiding this comment.
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.
| if self._split_coalesced and run is not None and len(run.unit_ids) > 1: | ||
| return list(run.unit_sizes) |
There was a problem hiding this comment.
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.
| 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) |
| 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"]), | ||
| ) |
There was a problem hiding this comment.
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.
| 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: |
There was a problem hiding this comment.
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.
| 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: |
There was a problem hiding this comment.
Can we come up wiht a different term? ReadUnit and UnitRun are pretty similar and confusing
| # --------------------------------------------------------------------------- | ||
|
|
||
|
|
||
| @pytest.mark.parametrize( |
There was a problem hiding this comment.
Why are these tests being removed?
There was a problem hiding this comment.
Theyre just being moved, verbatim
| @@ -1,113 +1,61 @@ | |||
| from typing import Any | |||
| from typing import Any, Optional | |||
There was a problem hiding this comment.
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
There was a problem hiding this comment.
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), |
There was a problem hiding this comment.
Do we not have the prefix sums to compute this cheaply?
There was a problem hiding this comment.
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>
Description
Last PR of the DSv2 Reorg series. Three changes, two types:
OnlineBinPackerpacks any file, not just Parquet. It read and wroteParquetRowGroupChunkMetadata, 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 sumssize_bytesagainst its budget and never asks what the number measures; the indexer that writes the row decides. The packing arithmetic is unchanged.RowGroupInfo(footer reader and coalescer),ParquetRowGroupChunkMetadata(manifest row), the copied fields on the packer'sBinItem, and thecreate_chunk_metadatahelper that built the row all carried the same ids / rows / bytes / per-unit tuples under different names. They collapse into one frozenFileChunkininterfaces/file_manifest.py, withto_metadata()/from_metadata()for the manifest's__file_chunk_metadatacolumn. (NamedUnitRunin the first draft; renamed per review because it read too close toReadUnit.)FileInfo(path, size)already existed ininterfaces/file_indexer.py, but the listing helpers passed(path, size)tuples,OrderedFileResultcarried the path twice, the footer actor tooklist[tuple[str, int]], its resultFileChunkshad its ownpath/size, and the packer'sBinItemhad its ownpath/size_bytes. All of them now carry aFileInfo. An indexer author emits one type; a reviewer checks no conversions.Cost of the unified types. Measured with
pympler.asizeofon Python 3.12, same 45-character path for old and new. A single-row-group run is 888 bytes as aFileChunkversus 920 as aRowGroupInfo; a three-group coalesced run is 1120 versus 1072. ABinItemis 1888 bytes versus 1416 for the old flatBinItem(path, size_bytes, run): the extra ~470 bytes is theFileInfoinstance 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
interfaces/file_manifest.py. NewFileChunk(unit_ids, num_rows, size_bytes, fully_matched=True, unit_sizes=(), unit_rows=())withto_metadata()/from_metadata();__post_init__asserts a per-unit breakdown is empty or one entry per unit.create_chunk_metadataand itsTypeVarare removed; the emptyChunkMetadataTypedDictbecomes aDict[str, Any]alias, since a row's chunk metadata is now alwaysNoneorFileChunk.to_metadata().formats/parquet/.parquet_footer_types.pyis deleted:RowGroupInfoandParquetRowGroupChunkMetadataare gone, and the one survivor,FileChunks→ChunkedFile(file: FileInfo, row_groups: Tuple[FileChunk, ...]), moves next to theFooterReaderthat builds it.footer_reader.pybuilds oneFileChunkper row group and_read_and_chunk/read_footerstakeFileInfo;parquet_row_group_coalescing.pymerges withdataclasses.replace;footer_file_indexer.pyemitsrun.to_metadata()via_chunked_file_to_manifestwithsize_bytes= the run's projection-scoped uncompressed size, as before.parquet_file_reader.pyreadsunit_idsandsize_bytes, and drops the"uncompressed_size" in chunkkey check, since a non-Nonerow always has the keys now. Docstrings inparquet_file_chunking_utils.pyfollow.common/indexing_utils.py,common/non_sampling_file_indexer.py.PathContents.filesand_get_file_infosproduceFileInfo;OrderedFileResultloses its duplicatefile_path.common/online_bin_packer.py.BinItem(file: FileInfo, run: Optional[FileChunk]), withsize_bytesa property (the run's bytes, else the file size);run is Noneis a whole-file row. Splitting a run slices theFileChunk; a sealed bin merges a file's runs back into oneFileChunkrow.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 claimingrow_group_ids=(0,)), and every packed row keeps the on-disk__file_sizethe listing gave it (a chunk row used to carry the run's bytes there; nothing reads__file_sizeafter the packer).FileChunkrows can hand them toOnlineBinPackeras they are.test_online_bin_packer.pybuilds rows directly and imports nothing fromformats/; itscoalesce_row_groupscases move verbatim to a newtest_parquet_row_group_coalescing.py, since they are Parquet's; two tests added (whole-file rows, on-disk size kept).test_chunk_metadata.pybecomes aFileChunkround-trip and invariant test.FileInfo/ renamed fields intest_footer_reader.py,test_read_units.py,test_synthesized_columns.py,test_parquet_datasource_v2.py,test_file_indexer.pyandtests/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_bytesis 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 seeChunkedFile.file.pathandread_footers(list[FileInfo])rather than tuples.Checks run locally from this branch (macOS,
~/ray/.venv):ruff,blackandpydoclintat the pre-commit pins are clean on every changed file;pyrefly0.51.0 withpyarrow==24.0.0reports 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