Repository navigation
Support range-based reads for deletion vectors - #3478
KaiqiJinWow wants to merge 4 commits into
Conversation
859efdc to
118c561
Compare
amogh-jahagirdar
left a comment
There was a problem hiding this comment.
Thanks @KaiqiJinWow, main comment is that I think we should introduce a new deletion_vector module which exposes a read_deletion_vector API and hides all the I/O, deserialization, validation. Looks like currently that's all kinda spread out over different classes.
Also just for transparency on what's driving this change to others, currently Databricks Runtime produces deletion vectors that are Iceberg spec compliant DV blobs but they are not neccessarily written in literal Puffin files (they're written in .bin files as a single blob) . The current PyIceberg implementation has strict checks that the DVs must be in literal Puffin files but that's not strictly neccessary. As long as the blob is spec compliant I think there's a reasonable argument that we can consume them regardless of what kind of literal file the blob is stored in. For context, the Java implementation also just works off a similar principle of just reading a spec compliant blob from a range.
46bdf9f to
b7c2ef4
Compare
|
I was asked for a review on this PR. @KaiqiJinWow is this WIP or is it ready for review? If you can get the integration test passing, I'd love to take a look. |
b7c2ef4 to
fa81f14
Compare
rambleraptor
left a comment
There was a problem hiding this comment.
I've got some questions around APIs mostly.
d92432d to
4786782
Compare
|
Hi @rambleraptor @amogh-jahagirdar @ebyhr, thanks for your reviews! I updated this PR to address the review feedback. The latest revision keeps the content-range DV path strict, preserves whole-Puffin reads, and cleans up the deletion_vector API surface. Could you take another look when you get a chance? Thanks! |
| if has_deletion_vector_content_reference(data_file): | ||
| return [_read_deletion_vector(io, data_file)] | ||
|
|
||
| with io.new_input(data_file.file_path).open() as fi: | ||
| return deletion_vectors_from_puffin_file(PuffinFile(fi.read())) |
There was a problem hiding this comment.
I don't think we need both branches. In both cases we are reading a single DV in a given byte range. Whether it's in a puffin or not should be inconsequential.
| if content_offset is None: | ||
| raise ValueError(f"Invalid deletion vector, content offset is missing: {data_file.file_path}") | ||
| if content_size_in_bytes is None: | ||
| raise ValueError(f"Invalid deletion vector, content size is missing: {data_file.file_path}") | ||
| if content_offset < 0: | ||
| raise ValueError(f"Invalid deletion vector, content offset cannot be negative: {content_offset}") |
There was a problem hiding this comment.
I am fine with having a more defensive implementation (the spec requires writers to produce the offset/size/refereenced file for DVs anyways) but just mentioning i think we only need to do these checks once and in one place only rather than in multiple places.
| if cardinality != record_count: | ||
| raise ValueError(f"Invalid cardinality: {cardinality}, expected {record_count}") |
There was a problem hiding this comment.
I think this is fine, again as we expect these two values to be the same but just remember implementations can choose to be a bit more relaxed (or vice versa more strict) than the actual spec. Is it worth failing the read of the DV if there's a mismatch? On one hand it indicates something incorrect in the metadata, on the other hand, we could be blocking a read of the data unnecessarily (because it wouldn't affect correctness of the result anyways). So in this case I'd probably bias to the latter of not doing this check. But I'll leave it up to you cc @kevinjqliu @rambleraptor in case you folks have opinions here.
There was a problem hiding this comment.
Java does the check, so we should probably keep it just to match the implementations. This is the kind of thing that I imagine iceberg-verification will be checking at some point and we don't want to have to add the check back in to help keep that repository green.
That being said, I always bias towards removing checks on user data, since we can't always assume the writer did a valid job. It's such a waste to not read a valid DV because of a mismatch.
amogh-jahagirdar
left a comment
There was a problem hiding this comment.
iceberg-python/pyiceberg/manifest.py
Line 557 in 2c75523
@KaiqiJinWow I think we need to double check the equals implementation. The code prior to this change only uses path to dedupe even for delete files/DVs, which are collected into a set. This used to work for the case where multiple DVs exist in a Puffin prior to this change just based off luck because we would read the whole puffin file in the end anyways. But in a world where we just do the range based reads which are more generic we cannot rely on this because we're not reading the whole puffin
f4b6a20 to
7de9394
Compare
|
The @dataclass(frozen=True, slots=True)
class DeleteFileKey:
file_path: str
content_offset: int | None
content_size_in_bytes: int | None
@classmethod
def from_file(cls, delete_file: DataFile) -> "DeleteFileKey":
return cls(
file_path=delete_file.file_path,
content_offset=delete_file.content_offset,
content_size_in_bytes=delete_file.content_size_in_bytes,
)
self._files: dict[DeleteFileKey, DataFile]For example: def add(self, delete_file: DataFile) -> None:
key = DeleteFileKey.from_file(delete_file)
self._files.setdefault(key, delete_file)
def discard(self, delete_file: DataFile) -> None:
key = DeleteFileKey.from_file(delete_file)
self._files.pop(key, None)This makes the intended identity clearer:
Using a named, immutable key also avoids relying on tuple ordering and keeps this specialized identity separate from the existing path-based |
|
This pull request has been marked as stale due to 30 days of inactivity. It will be closed in 1 week if no further activity occurs. If you think that's incorrect or this pull request requires a review, please simply write any comment. If closed, you can revive the PR at any time and @mention a reviewer or discuss it on the dev@iceberg.apache.org list. Thank you for your contributions. |
7de9394 to
d71d94c
Compare
d71d94c to
0e99544
Compare
|
Hi @amogh-jahagirdar @rambleraptor @ebyhr @kevinjqliu, #3690 is now merged, and this PR is rebased with CI green. I’ve addressed the previous feedback, including the range aware Could you please take another look? Thanks! |
rambleraptor
left a comment
There was a problem hiding this comment.
Can you add an integration test for this with Spark? I'd really love to see that we can successfully read a deletion vector written by an outside source.
Even copying in a fixture made by somebody else would be great.
| ) -> None: | ||
| self.file = data_file | ||
| self.delete_files = delete_files or set() | ||
| self.delete_files = DeleteFileSet(delete_files if delete_files is not None else []) |
There was a problem hiding this comment.
Doesn't look like you need the default, since DeleteFileSet already sets a default.
| if cardinality != record_count: | ||
| raise ValueError(f"Invalid cardinality: {cardinality}, expected {record_count}") |
There was a problem hiding this comment.
Java does the check, so we should probably keep it just to match the implementations. This is the kind of thing that I imagine iceberg-verification will be checking at some point and we don't want to have to add the check back in to help keep that repository green.
That being said, I always bias towards removing checks on user data, since we can't always assume the writer did a valid job. It's such a waste to not read a valid DV because of a mismatch.
Thanks @rambleraptor! I added two externally generated The tests verify the expected deleted row positions using the manifest-provided content ranges. The packed fixture test also checks that all six DVs survive deduplication and are associated with the correct data files. Could you take another look when you have a chance? Thanks! |
Fokko
left a comment
There was a problem hiding this comment.
Left some comments, but this looks great to me 👍
| ) | ||
|
|
||
|
|
||
| class DeleteFileSet(MutableSet[DataFile]): |
There was a problem hiding this comment.
I'm wondering if we could make it extend Set, rather than MutableSet. We don't used discard and update does an add operation. I think having this as immutable, that it makes it easier to reason about the code and flow.
There was a problem hiding this comment.
+1
I don't love the idea of us creating our own set class if possible (that may have unexpected semantics under the hood)
There was a problem hiding this comment.
Could we make this a frozenset?
|
|
||
| def has_deletion_vector_content_reference(dv: "DataFile") -> bool: | ||
| """Return whether a deletion vector is described by manifest content-range metadata.""" | ||
| return dv.content_offset is not None or dv.content_size_in_bytes is not None or dv.referenced_data_file is not None |
There was a problem hiding this comment.
I think these should be and, rather than or, since we require them all downstream in _read_deletion_vector
There was a problem hiding this comment.
Hi Fokko, the or is intentional: if any content-reference field is present, we take the range-read path and validate that all required fields are set later. The whole-Puffin fallback is only used when all three fields are absent. Changing this to and would let partial metadata bypass validation. I’ll clarify the docstring and add a test for that case.
| other_keys: set[DeleteFileKey] = set() | ||
| other_count = 0 | ||
| for delete_file in other: | ||
| if not isinstance(delete_file, DataFile): | ||
| return False | ||
| other_keys.add(DeleteFileKey.from_file(delete_file)) | ||
| other_count += 1 | ||
|
|
||
| return len(other_keys) == other_count and set(self._files) == other_keys |
There was a problem hiding this comment.
This part feels odd to me, why do we want to compare this to an iterable?
| from pyiceberg.manifest import DataFile | ||
|
|
||
|
|
||
| @dataclass(frozen=True, slots=True) |
There was a problem hiding this comment.
Love the slots=True. We define slots by hand throughout the codebase, but that's not needed anymore with Python 3.10
rambleraptor
left a comment
There was a problem hiding this comment.
Broadly, this looks great. Most of my questions are around the DeleteFileSet API
|
|
||
| file: DataFile | ||
| delete_files: set[DataFile] | ||
| delete_files: DeleteFileSet |
There was a problem hiding this comment.
This is going to change the API surface, since DeleteFileSet is a MutableSet, not a regular set.
Sets have a whole mess of built-in methods and we don't want to replicate all of them.
| ) | ||
|
|
||
|
|
||
| class DeleteFileSet(MutableSet[DataFile]): |
There was a problem hiding this comment.
Could we make this a frozenset?
Restore native scan task sets by routing deletion vectors to their referenced data files. Preserve REST range references and cover local and REST scans of packed DV fixtures.
|
Thanks @rambleraptor @Fokko @amogh-jahagirdar @kevinjqliu for the feedback! I’ve pushed 5c05f03. In this commit, I removed I also fixed REST planning to preserve the DV metadata and expanded the packed Could you take another look at the matching and deduplication changes? Thanks! |
Summary
Builds on #3690, which enables reading V3 deletion-vector content-range fields from manifests.
content_offsetandcontent_size_in_bytes, without requiring the physical file to be a complete Puffin file.set[DataFile]API andDataFileequality unchanged. Match DVs to their referenced data files before constructing task sets, and deduplicate across tasks internally by(file_path, content_offset, content_size_in_bytes).Testing
Automated
.binfixtures containing a single DV and six DVs sharing one physical file.Manual Validation
.binDV file and its Iceberg V3 metadata.StaticTable.scan()using the packed fixture with locally assembled V3 metadata and Parquet files. Full scans, filtered scans, and point lookups returned the expected results with both PyArrowFileIO and FsspecFileIO, includingto_arrow(),count(), and batch reads. This additional check is not part of CI.