Conversation
sunchao
left a comment
There was a problem hiding this comment.
Reviewed b474ade26a33e58d5c46f2d31a1b1a6d7f177d0d against eabb5d4773091b983d8fce713f0e34b1cf93f877. No P1/P2 findings.
Required-accessor failures now propagate instead of substituting delete metadata, while valid null positional equality IDs remain supported. The successful serialization and exception-cause behavior were checked with real Iceberg objects in isolated component controls across the four pinned dependency versions.
All six new tests also passed in the full Linux Spark 3.4/JDK 11 and Spark 4.0/JDK 21 CI runs. Those runs checked out the PR merge commit, whose tree matches this head; the local Scala 2.12/2.13 component checks are separate evidence, not full Spark/JNI runs.
CI is not all green: the macOS scans job failed with a native SIGBUS during thread cleanup before any recorded execution of the new serde suite. The supplied logs and crash artifacts did not establish a causal changed-line defect. That crash remains unattributed; this approval does not classify it as pre-existing.
andygrove
left a comment
There was a problem hiding this comment.
Thanks for this. The core change is right, and the extraction into a testable serializeDeleteFile is a real improvement.
Two things I checked independently that back the change up:
content(),specId(),equalityFieldIds()andkeyMetadata()are all declared onorg.apache.iceberg.ContentFilein every pinned version. I ranjavapagainst iceberg-spark-runtime 1.5.2, 1.8.1, 1.10.0 and 1.11.0. So removing the fallbacks does not change behavior on any supported Iceberg. Incidentallylocation()only shows up from 1.8.1 onward, which is whyextractFileLocationstill needs itspath()branch.- The
nullhandling is load-bearing rather than defensive.BaseFile.equalityFieldIds()goes throughArrayUtil.toUnmodifiableIntList, which returnsnullfor a null backing array. So the oldcatchwas swallowing an NPE on every position-delete file, and without the new guard this PR would fail every position delete. Good catch, and the test covers it.
I have left comments on the equality-ids block, the content() match, and the new suite.
On CI: the macOS scans job is red because the Download native library step failed after 25 seconds, so the Java test steps were skipped and the new suite has never actually run on macOS. That looks like an artifact-download hiccup rather than anything in the change. Could you push or re-run to get a clean result? For what it is worth, I confirmed the suite runs and passes in the Linux scans jobs, it shows up in the MAVEN_SUITES list and reports CometIcebergDeleteFileSerdeSuite: in the log.
Separately, an earlier automated review here attributed the macOS scans failure to a native SIGBUS. That was a different job on an earlier head, so it is not what is failing now.
| val equalityIdsMethod = IcebergReflection.getMethod(deleteFileClass, "equalityFieldIds") | ||
| val equalityIds = equalityIdsMethod.invoke(deleteFile).asInstanceOf[java.util.List[Integer]] | ||
| // Iceberg's BaseFile stores equality field IDs in a nullable backing array, so | ||
| // equalityFieldIds() returns null for files without equality keys. A null return is a | ||
| // normal accessor result, unlike a reflective lookup or invocation failure, and does not | ||
| // make serialization fail. | ||
| if (equalityIds != null) { | ||
| equalityIds.forEach(id => deleteBuilder.addEqualityIds(id)) | ||
| } |
There was a problem hiding this comment.
Nice catch on the null return here. I checked BaseFile.equalityFieldIds() and it goes through ArrayUtil.toUnmodifiableIntList, which returns null for a null backing array, so the old catch really was swallowing an NPE on every position-delete file. Without this guard the PR would break every position delete.
One thing though. IcebergReflection.getEqualityFieldIds does almost exactly this, including the null-to-empty mapping, but keeps the blanket catch. CometIcebergNativeScan still calls it a bit further up to decide whether to union equality-delete field IDs into the task schema via schemaWithRequiredFields. That is the second instance #5256 calls out under finding 3:
Finding 3 has a second, independent instance:
IcebergReflection.getEqualityFieldIdsswallows the same failure into an empty list, and that result also drives the task-schema decision.
and its Expected behavior section asks for the helper to distinguish a genuine "no equality ids" from a reflective failure.
Would it work to drop the catch from getEqualityFieldIds, keep the null-to-empty mapping there, and have serializeDeleteFile call it? That fixes both call sites and avoids carrying two copies of this reflection with different failure semantics.
There is a wrinkle with the third caller. CometScanRule calls the same helper during planning, where falling back to Spark is still an option, so a blanket change would turn a fallback into a failure there. That is the same split-entry-point problem #5257 describes for buildFieldIdMapping, so it may be cleaner as a separate change.
If you would rather keep this PR tight, could we switch the description to "Part of #5256" and file a follow-up for the helper? I checked and it is not covered by #5257, which enumerates buildFieldIdMapping, pageIndexUnsupportedColumns and PartitionSpecParser.toJson but not this one, and #5258 is already closed. So as written the issue would get closed with half of finding 3 still open.
There was a problem hiding this comment.
Good point. I kept the planning helper unchanged since CometScanRule still needs to be able to fall back to Spark.
I added a strict serde-side helper instead and use it in both serializeDeleteFile and the task-schema union path. So lookup/invocation failures are fatal once native execution is committed, while a null return is still treated as empty for position deletes.
This should cover the second #5256 path without changing the planning-time fallback behavior.
| val contentType = contentMethod.invoke(deleteFile).toString match { | ||
| case IcebergReflection.ContentTypes.POSITION_DELETES => | ||
| IcebergReflection.ContentTypes.POSITION_DELETES | ||
| case IcebergReflection.ContentTypes.EQUALITY_DELETES => | ||
| IcebergReflection.ContentTypes.EQUALITY_DELETES | ||
| case other => other | ||
| } |
There was a problem hiding this comment.
All three branches of this match return their argument unchanged, so the whole thing is equivalent to contentMethod.invoke(deleteFile).toString.
The real validation is already on the native side. planner.rs:3910 rejects any content type that is not one of these two with Invalid delete content type '{}'. Since you are rewriting this block anyway, could we drop the match? As it stands it reads like the Scala side is filtering when it is not.
There was a problem hiding this comment.
Yep, agreed. Removed the match and now just serialize content().toString directly. The native planner remains responsible for validating the content type.
| test("position-delete file: null equalityFieldIds() serializes with no equality ids") { | ||
| val proto = serialize(new PositionDeleteFile) | ||
| assert(proto.getContentType == "POSITION_DELETES") | ||
| assert(proto.getPartitionSpecId == 7) | ||
| assert(proto.getEqualityIdsCount == 0) | ||
| assert(proto.getFilePath == "s3://bucket/pos-delete.parquet") | ||
| } | ||
|
|
||
| test("equality-delete file: declared equalityFieldIds() are serialized") { | ||
| val proto = serialize(new EqualityDeleteFile) | ||
| assert(proto.getContentType == "EQUALITY_DELETES") | ||
| assert(proto.getEqualityIdsCount == 2) | ||
| assert(proto.getEqualityIds(0) == 3) | ||
| assert(proto.getEqualityIds(1) == 5) | ||
| } | ||
|
|
||
| test("content() invocation failure propagates instead of defaulting to POSITION_DELETES") { | ||
| val ex = intercept[InvocationTargetException](serialize(new ThrowingContentDeleteFile)) | ||
| assert(ex.getCause.getMessage == "content boom") | ||
| } | ||
|
|
||
| test("specId() invocation failure propagates instead of defaulting to 0") { | ||
| val ex = intercept[InvocationTargetException](serialize(new ThrowingSpecIdDeleteFile)) | ||
| assert(ex.getCause.getMessage == "spec boom") | ||
| } | ||
|
|
||
| test("equalityFieldIds() invocation failure propagates instead of dropping equality ids") { | ||
| val ex = intercept[InvocationTargetException](serialize(new ThrowingEqualityIdsDeleteFile)) | ||
| assert(ex.getCause.getMessage == "ids boom") | ||
| } | ||
|
|
||
| test("missing content() accessor is fatal, not a default") { | ||
| // getMethod throws NoSuchMethodException directly (not wrapped) when the accessor is absent. | ||
| assertThrows[NoSuchMethodException](serialize(new NoContentAccessorDeleteFile)) | ||
| } |
There was a problem hiding this comment.
The synthetic stubs are a clean way to pin the control flow, thanks for that.
The argument for removing the fallbacks is that content(), specId() and equalityFieldIds() are always declared on the real Iceberg interfaces, and I do not think anything in the suite checks that. I verified it holds with javap across all four pinned versions, and they are all on ContentFile. But the consequence of that ceasing to be true changed with this PR. It used to be quietly wrong rows, now it is a hard query failure for anyone with delete files.
That is the regression that matters most now, and it is the one a test would let CI catch at the version bump rather than in the field. The pieces are already here, iceberg-spark-runtime is a test dependency in every Spark profile and IcebergReflectionSuite already resolves methods against real Iceberg classes. Would you add something like assert(IcebergReflection.findMethod(classOf[DeleteFile], "content").isDefined) for each of the three?
Also, the missing-path branch in serializeDeleteFile is the one fail-loud path without a test, and it is the one carrying the longest comment about why it has to be fatal. A stub declaring neither location() nor path() would round the suite out.
There was a problem hiding this comment.
Added both. The suite now checks content, specId, and equalityFieldIds against the real DeleteFile API, and there is a regression test verifying that a delete file with neither location() nor path() fails loudly.
I also added coverage for a missing equalityFieldIds() accessor.
|
Let's also give @mbutrovich an opportunity to review this PR |
sunchao
left a comment
There was a problem hiding this comment.
Re-reviewed 13a17c5efd13375649f52449b0becef9c4d868ac. The discussed serde and task-schema reflection failure paths are addressed, with no new P1/P2 findings. No local tests were run. Current-head workflows remain action_required, with no test checks recorded.
andygrove
left a comment
There was a problem hiding this comment.
All three of my inline points are addressed, and the split you landed on the equality-ids helper is better than either option I offered.
The content-type match is gone, so the block reads as what it is: the Scala side records the string and planner.rs does the validation. The suite is now genuinely thorough, and the eleven cases cover exactly the distinctions I was worried about, including the three that would previously have been silently defaulted:
- content() invocation failure propagates instead of defaulting to POSITION_DELETES
- specId() invocation failure propagates instead of defaulting to 0
- equalityFieldIds() invocation failure propagates instead of dropping equality ids
- missing content() accessor is fatal, not a default
- missing equalityFieldIds() accessor is fatal, not an empty list
- position-delete file: null equalityFieldIds() serializes with no equality ids
- equality-delete file: null equalityFieldIds() is fatal
- equality-delete file: empty equalityFieldIds() is fatal
On the second instance of finding 3, you have given the two call sites deliberately different semantics rather than one blanket rule, and I think that is right. requiredEqualityFieldIds in CometIcebergNativeScan goes through getMethod, so a missing accessor or a failed invocation is fatal, which is correct because the scan is already committed to native by then. IcebergReflection.getEqualityFieldIds now uses findMethod and maps null to empty with no catch at all, so an invocation failure propagates there too and only a genuinely absent method degrades to an empty list. Its one remaining caller is CometScanRule during planning, where falling back to Spark is still available, so degrading is the right behavior there.
That resolves my objection to Closes #5256 as well. The concern was that finding 3's second instance would be left swallowing a reflective failure; it no longer does. The absent-method case is the only one that still degrades, and I verified earlier with javap that equalityFieldIds() is declared on ContentFile in iceberg-spark-runtime 1.5.2, 1.8.1, 1.10.0 and 1.11.0, so it is unreachable on every supported version. Worth a sentence in the description saying that, since it is the one thing a reader would otherwise have to re-derive to agree that the issue is fully closed.
Both pr_build_linux.yml and pr_build_macos.yml carry the new suite, so the macOS gap I flagged is covered once the jobs run.
The branch conflicts with main now. Could you rebase? Approving in principle; I would rather see the rebase and a clean run before it merges, given the earlier macOS artifact-download failure meant the suite never actually executed there.
|
We should also make sure we throw on the |
13a17c5 to
91ed605
Compare
|
Rebased on latest |
|
I stopped CI since it looks like formatting was failing. |
mbutrovich
left a comment
There was a problem hiding this comment.
I checked this against the invariant I care about for the two Iceberg reflection sites. Failures in CometScanRule should fall back with a reason. Failures in CometIcebergNativeScan while building the proto have to be fatal, because past that point we would run a scan on guessed metadata and silently return wrong rows.
The PR gets that right for all three fields. Two things I verified beyond the diff:
extractDeleteFilesListrethrows asRuntimeExceptionat line 305 rather than returning an emptySeq, so nothing degrades on the way out ofserializeDeleteFile.- Both equality-id call sites in the serde use the strict helper, line 343 and line 1039, so the task-schema union path is covered and not just the delete-file proto. That is the second instance of finding 3 in #5256.
Leaving IcebergReflection.getEqualityFieldIds lenient is also correct. Its only remaining caller is the planning path, and CometScanRule.scala:974 catches and appends a fallback reason, so an invoke failure there becomes a fallback rather than a hard error.
Two spots in the same file still fall short of that invariant. Neither is introduced here and I do not think either belongs in this PR:
CometIcebergNativeScan.scala:477, a partition-spec JSON failure is only alogWarningand the task ships with nopartitionSpecIdx. This is thePartitionSpecParser.toJsoncase already tracked in #5257.CometIcebergNativeScan.scala:689and:1137, residual expression failures degrade toNoneand the task ships without the residual filter. This one is not in #5257's list. If Iceberg reported the filter as fully pushed, Spark will not re-apply it above the scan and we return extra rows.
Could you file an issue for the residual case so it does not get lost? Happy to link it from #5257.
| deleteBuilder.setFilePath(deletePath) | ||
|
|
||
| val contentMethod = IcebergReflection.getMethod(deleteFileClass, "content") | ||
| val contentType = contentMethod.invoke(deleteFile).toString |
There was a problem hiding this comment.
Dropping the match is right, the real validation is already on the native side at planner.rs:4364. One consequence worth covering though. contentType is now a raw pass-through of the enum's name(), and planner.rs matches those strings exactly.
I checked iceberg-spark-runtime-4.1_2.13-1.11.0 with javap. ContentFile.content() returns FileContent, and FileContent declares POSITION_DELETES and EQUALITY_DELETES without overriding toString, so the names line up today. Nothing in the suite asserts that. A constant rename in a future Iceberg would pass the new accessor-existence test and then fail every query that has delete files.
Separately, after this hunk IcebergReflection.ContentTypes.POSITION_DELETES has no callers left. A grep over spark/, common/ and native/ only finds the three occurrences this hunk deletes, and line 345 is the last use of EQUALITY_DELETES. Pinning the names in a test (see my comment on the suite) would give the constant a use again. Otherwise it should probably be removed.
| val equalityFieldIds = requiredEqualityFieldIds(deleteFileClass, deleteFile) | ||
| val isEqualityDelete = | ||
| contentType == IcebergReflection.ContentTypes.EQUALITY_DELETES | ||
| if (isEqualityDelete && equalityFieldIds.isEmpty) { |
There was a problem hiding this comment.
This is a new fatal case rather than one of the three fallbacks the PR removes. Before, an equality delete with empty ids serialized with no equality ids. Now it fails the query. I think that is correct, an equality delete with no keys cannot be applied, so failing beats guessing.
It is not in the description though. The equalityFieldIds() bullet currently says null is preserved as a legitimate no-equality-keys result. That holds for position deletes only. For an equality delete, null and empty are both fatal after this change. Could you add this to the changes list and adjust that bullet? The description becomes the squash commit message on a Closes #5256 fix, so it is worth having it match what landed.
| file.getClass, | ||
| keyMetadataMethod(file.getClass)) | ||
|
|
||
| test("required Iceberg delete-file accessors are present") { |
There was a problem hiding this comment.
This is the test that closes the version-bump gap, good addition. Could it assert the values too, not only that the methods exist?
assert(FileContent.POSITION_DELETES.toString == IcebergReflection.ContentTypes.POSITION_DELETES)
assert(FileContent.EQUALITY_DELETES.toString == IcebergReflection.ContentTypes.EQUALITY_DELETES)serializeDeleteFile now sends content().toString straight into the proto and planner.rs:4364 matches those two literals, so the string values are as load-bearing as the accessors are. The stubs all declare def content(): String, so nothing in the suite touches the real enum. org.apache.iceberg.FileContent is on the test classpath already, same jar as the DeleteFile import on line 26.
sunchao
left a comment
There was a problem hiding this comment.
Re-reviewed 91ed605 against fd8e09e, including the rebase since 13a17c5. The feedback commit is patch-equivalent. The strict helper still covers both native serialization paths, and planning-time failures still reach Spark fallback. Andy’s three earlier code and test requests remain addressed. Mike’s enum-value test and description clarification remain open. I found no new P1/P2 issue.
The execution-boundary check used maintained Spark 3.5 and 4.0 sources. No local tests ran because the isolated dependency-based probe was blocked before execution. Current CI remains failed or cancelled. Failure logs and the actual CI checkout remain unverified. This approval is based on source review, not a green validation run.
91ed605 to
3415574
Compare
|
@mbutrovich Thanks, addressed the remaining feedback. Updated the PR description to clarify the position-delete vs equality-delete |
sunchao
left a comment
There was a problem hiding this comment.
Re-reviewed 3415574e against 4d931e36. The new test checks both real Iceberg enum names, the PR description now distinguishes valid null position-delete IDs from invalid null/empty equality-delete IDs, and the residual-reflection concern is tracked in #5992.
The rebased implementation keeps strict metadata reads in both delete serialization and task-schema expansion, preserves the original exception causes, and retains the inherited deletion-vector metadata and path interning. The regression cases now live in the existing CometIcebergNativeScanSuite, already selected by both Linux and macOS CI. No new or remaining P1/P2 findings in this diff.
Validation is source-only: git diff --check passed, and I rechecked the maintained Spark 3.5/4.0 execution boundary and the pinned Iceberg API contracts. No local JVM/native tests ran. Maintained Spark 3.4/4.1/4.2 sources were unavailable, so this does not establish compatibility for those versions. Current CI, CodeQL, and title-check workflows are action_required with zero jobs. Only labeling succeeded. There is no current-head product-test result yet.
sunchao
left a comment
There was a problem hiding this comment.
Re-reviewed b710d7e0 against 5ca14992, including the merge since the previous review at 3415574e. The delete serializer, strict equality-ID helper and all 12 delete-file regression cases are unchanged. The merge retains the base's complex-type residual guard and its test, with both sets of imports preserved.
Required delete-metadata failures still propagate through both serialization paths. The merge preserves deletion-vector coordinates, record counts, key metadata and path interning. The separate residual-reflection failure boundary remains unchanged and tracked in #5992. No new or remaining P1/P2 findings in this update.
Validation is source-only. Exact source-equivalence checks and git diff --check passed. I rechecked the maintained Spark 3.5/4.0 execution boundary and refreshed the primary Iceberg API sources for the four pinned Maven dependency versions. Maintained Spark 3.4/4.1/4.2 sources remain unavailable. No local JVM/native tests or dependency-based probe ran, and CI logs were not accessed. Current Comet CI and CodeQL require approval and have zero jobs. Only labeling succeeded, so there is no current-head product-test result or executed checkout to credit.
andygrove
left a comment
There was a problem hiding this comment.
LGTM pending CI. Thanks @unikdahal
|
@andygrove @mbutrovich Could you please retrigger CI? Looks like an unrelated Gradle download failure. |
|
@unikdahal could you fix the conflicts? after that I can merge the PR. thanks |
Which issue does this PR close?
Closes #5256.
Rationale for this change
CometIcebergNativeScanserializes Iceberg delete-file metadata afterCometScanRulehas already committed the scan to native execution.Previously, reflection failures while reading required delete-file metadata could silently fall back to guessed or incomplete values:
content()failure defaulted toPOSITION_DELETESspecId()failure defaulted to0equalityFieldIds()failure could result in equality field IDs being omittedThose defaults are not safe. They can cause an equality delete to be interpreted as a position delete, associate a delete file with the wrong partition spec, or serialize an equality delete without the keys required to apply it.
At serde time there is no Spark fallback left, so required delete metadata must either be serialized correctly or fail the query.
What changes are included in this PR?
Remove the fallback to
POSITION_DELETESwhencontent()lookup or invocation fails.Remove the fallback to
0whenspecId()lookup or invocation fails.Stop swallowing lookup or invocation failures from
equalityFieldIds().Use a strict serde-side equality-field-ID helper in both:
Keep planning-time handling separate from serde-time handling. During planning, a genuinely absent accessor may still degrade where Spark fallback remains available; invocation failures propagate to the planning fallback path.
Treat a
nullequalityFieldIds()result as an empty list for position deletes. Iceberg normally represents files without equality keys this way.Treat equality deletes with either
nullor empty equality field IDs as fatal because an equality delete without keys cannot be applied correctly.Serialize
content().toStringdirectly and leave validation of the content type to the native planner.Preserve the existing delete-file serialization behavior for file-path interning, file format, deletion-vector metadata, record count, and key metadata.
Add regression coverage to the existing
CometIcebergNativeScanSuitefor the serde failure semantics and supported Iceberg API contract.Pin the real Iceberg
FileContent.POSITION_DELETESandFileContent.EQUALITY_DELETESstring values against the literals consumed by native serde so a future dependency change cannot silently break the Java/native contract.Compatibility
content(),specId(), andequalityFieldIds()are declared on Iceberg's publicContentFile/DeleteFileAPIs across the Iceberg versions currently supported by Comet.The change therefore does not alter normal behavior for supported versions. It changes what happens when reflective lookup or invocation fails after native execution has already been selected: instead of serializing guessed metadata, the query fails.
The planning path remains separate because Spark fallback is still available there.
The tests also verify that Iceberg's real
FileContent.POSITION_DELETESandFileContent.EQUALITY_DELETESstring values continue to match the values expected by the native planner.How are these changes tested?
Regression tests in the existing
CometIcebergNativeScanSuitecover:content(),specId(), andequalityFieldIds()accessors existing on the real IcebergDeleteFileAPI;nullequality field IDs serializing successfully with no equality IDs;nullequality field IDs failing;content()invocation failures propagating instead of defaulting toPOSITION_DELETES;specId()invocation failures propagating instead of defaulting to0;equalityFieldIds()invocation failures propagating instead of dropping equality IDs;content()accessor being fatal;equalityFieldIds()accessor being fatal; andThe tests live in the existing suite and run through the normal CI matrix.