fix: report row contents to native write statistics trackers - #6383
LinSimon-901101 wants to merge 3 commits into
Conversation
|
Reopening after the standard PR CI passed in my fork: CI run. |
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Native Parquet writes supplied empty rows to statistics trackers, preventing value-based statistics.
- Design approach: Preserve counting for the exact
BasicWriteTaskStatsTrackerclass and supply materialized rows to other trackers, including subclasses. - Correctness / compatibility analysis: Compared Spark’s tracker contract, projection behavior and structural type checks across Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. Checked row lifetime, null handling, callback ordering and failure cleanup.
- Key design decisions: One reusable
UnsafeProjectionper task limits allocation. Basic trackers avoid row materialization. The change stays within the existing writer and iterator abstractions. - Implementation sketch: Report rows before Arrow consumes each batch, validate compatible types before reading values, and close the batch if reporting fails.
- Behavioral changes worth calling out: Custom trackers receive row contents and must copy retained rows. Physical type mismatches fail explicitly. Field-name, nullability and Arrow-erased interval metadata differences remain accepted.
- Suggested improvements: No introduced P1/P2 issues found within this review.
Reviewed the entire two-file diff from 1628c520096e12ba3fb13988d6ee235070e77e07 to 867e5cc68feef8cfca5e55739e6945905be26ea4. The PR remains non-draft. Read the supplied discussion and linked issue. Routed skills: review-comet-pr, review-comet-ffi-pr, and review-comet-memory-pr.
Exact-head CI: Comet CI and CodeQL are action_required, with no jobs executed. The labeling check passed. The successful fork run cited by the author tested earlier head 8fe68223783dccf963abd477deec3072fe823b49.
Validation: Disposable harnesses using verified exact-head callback and conversion code passed eight row-statistics cases and four failure/cleanup cases on each of Spark 3.5.9 and 4.1.3. Additional checks using the exact-head Comet vector classes passed for nested arrays/maps, decimals, dictionary strings, interval subtypes, nulls and copied rows after closure. git diff --check passed.
Validation limits: make core reached the 120-second build limit without completing. The full writer suite, JNI integration tests and Spark SQL suite were not run locally. These isolated checks do not establish end-to-end native-write or performance results.
| if (!matchingTypes) { | ||
| throw new UnsupportedOperationException( | ||
| "Comet's native Parquet writer cannot report row statistics for a batch whose " + | ||
| "types differ from the write schema. Set spark.comet.parquet.write.enabled=false " + |
There was a problem hiding this comment.
This needs a merge with main, and the conflict lands in recordRows. #6494 renamed spark.comet.parquet.write.enabled to spark.comet.write.parquet.enabled outright with no alias, so after the merge this message would point users at a key that no longer does anything. The conflicting hunk is main's version of the old warning, which this PR deletes, so taking your side resolves it cleanly and leaves the stale key here without any visible conflict. Could you build the message from CometConf.COMET_NATIVE_PARQUET_WRITE_ENABLED.key instead, so the next rename can't strand it?
Otherwise this looks good to me. I ran the writer suite with the PR merged into current main on Spark 4.1. I also ran the same recording tracker through Spark's own WriteFilesExec and through CometWriteFilesExec on 4.1 and 4.2, and the trackers saw identical rows across about 20 column types, including NaN, -0.0, timestamps before 1970, binary, intervals and nested types.
|
@andygrove Thanks for pointing this out. I've merged main locally, resolved the conflict, and updated the error message to use I'm running CI in my fork to validate the changes. Once it passes, I'll push the updates to this PR. |
|
Here are the CI runs in my fork for PR #6383: |
Which issue does this PR close?
Closes #5307.
Rationale for this change
Native Parquet writes currently pass empty rows to every write statistics tracker. Trackers that inspect column values therefore cannot compute correct statistics. Supply actual row contents while keeping the basic row-counting path free of row materialization.
What changes are included in this PR?
How are these changes tested?
8fe682237: 20 jobs succeeded, 15 were skipped according to the PR-tier policy, andRequired Checkspassed.anyNull, nested and decimal values, copied rows after batch closure, interval subtypes, the basic counting path, and callback failures through the native writer.