feat: add native Delta Lake scan contrib module (page/row-group pruning) - #5365
dwsmith1983 wants to merge 64 commits into
Conversation
|
Update: pushed two follow-up commits extending the scan's pruning and object-store behavior.
|
888e4a7 to
7fd81aa
Compare
|
HI @andygrove, Can you review this as it adds Delta functionality? |
|
Hi @dwsmith1983 Thanks for putting this together! We are also actively looking at Delta support for Comet, and it'd be great if we can collaborate on this effort! Since #4952 is already approved and close to landing, what do you think about using it as the shared foundation for this work? Ideally, the same In addition, would it also make sense to land this in smaller pieces, for easier review and iterating? For example:
Starting to support this in Spark 4 & Delta 4 would be a useful first milestone. Curious how you see the relationship between the two PRs and whether that direction makes sense to you. Thanks. |
|
Hi @sunchao, On #4952 as the foundation: we already share more than it might look like. This PR builds on part 1 of that same breakup (#4700's CometScanWithPlanData / PlanDataInjector SPI) and keeps #4366's contrib shape, decline-gate philosophy, and test catalog, with co-authored-by credit to both earlier efforts. The remaining overlap is contrib infrastructure, and I'm glad to reconcile it once #4952 lands: adopt its contrib-delta profile and feature naming, the per-Spark delta.version matrix, the verify-gate script, and unify the proto slot (this PR is at 119, #4952 at 118). For the claim hook I'd suggest the generic CometScanRuleExtension SPI from this PR, since it keeps core free of Delta-specific code and the kernel path can register through it the same way. I do see the two read paths as different layers rather than one thing to converge on. By the time CometScanRule sees the scan, delta-spark has already done log replay, time travel, and partition pruning, so this path reuses Comet's existing native parquet scan and gets row-group pruning, page-index pruning, and filter pushdown for free. DVs become ParquetAccessPlans that DataFusion intersects with page-index pruning, so DV skips and page skips compose in one scan. As far as I know no vectorized Delta reader does all of that today, including kernel's, which has no page-index pruning. I'd want convergence to keep this as the default read path, with the kernel path covering what JVM planning can't reach (DSv2, non-Spark frontends, likely CDF and row tracking). On splitting: I'd push back on slicing by feature, for two reasons. First, the features aren't independent. Several decline gates only exist because DVs, column mapping, and Delta's own suites ran together. For example, Delta's findTouchedFiles scan looks like a plain read, and if a basic-reads slice claims it, DELETE silently rewrites files instead of writing DVs. Second, the proof is holistic: this branch runs Delta's own suites at 1156/1156 and the contrib suites at 39/39 on Spark 3.5, 4.0, and 4.1. Feature slices would decline most tables and couldn't run that meaningfully. What I can do is split along review surfaces instead: core SPI additions, native DV decode with its unit tests, the contrib module and read path, and the regression harness and CI, keeping the read path itself (DVs, column mapping, gates) as one reviewable unit. If it lands whole, Comet ships the only vectorized Delta reader with complete skipping. The Spark 4 milestone is already met, the suites are green on 4.0 and 4.1 today. Row tracking and CDF are out of scope here and seem like a natural place for the kernel work to lead. Happy to set up a chat with you and @schenksj to work out the details. |
|
Thanks @dwsmith1983 Your proposed split by review surface sounds reasonable. I agree that the reader, its safety gates, and the essential DML/fallback tests should stay together. Thanks also for being open to aligning with #4952 once it lands. We can leave row tracking and CDF for later discussions rather than expand this PR’s scope. The main additional point I’d like us to settle is keeping experimental Delta support explicitly opt-in. |
Yeah, agreed on explicit opt-in. It's mostly already set up that way. All the Delta code lives in a separate comet-contrib-delta jar that never gets bundled into comet-spark, so a stock Comet install has no Delta surface at all. If we publish that jar with releases, trying it out is just --packages and a conf, nobody has to build from source. Right now the conf defaults to on when the jar is present though, so I'll flip spark.comet.scan.delta.enabled to default false to make the opt-in explicit. The one spot where I'd differ from #4952's gate is the native binary. The Delta bits in libcomet are tiny (DV decoding plus a hand-off to the existing parquet scan, no delta-kernel dependency) and can't be reached without the jar and the conf. I'd rather keep them in the default build than make people compile their own native binary to try an experimental feature. Sound reasonable? |
|
Thanks @dwsmith1983. This makes sense to me! #4952 has just been merged. Could you rebase this PR and adapt to it? Thanks! |
|
@sunchao A few things I deliberately left for discussion rather than deciding unilaterally: unifying the two claim hooks in CometScanRule (CometScanContrib vs the CometScanRuleExtension SPI), conf naming (spark.comet.scan.delta.* vs spark.comet.scan.deltaNative.*), and Maven packaging (the -Pcontrib-delta add-source vs this module's separate jar, which is what keeps the opt-in story build-free). |
sunchao
left a comment
There was a problem hiding this comment.
Thanks for reconciling this with #4952. The generic envelope, separate source roots, and explicit runtime opt-in look like useful progress. I reviewed 9b393c15 and left eight concrete correctness and compatibility comments. The main concerns are unsafe scalar-subquery pushdown, mixed-authority file routing, and unbounded deletion-vector row-selection memory.
I checked these against Spark/Delta source and used bounded stock Spark 4.0.3 / Delta 4.0.0 and isolated Rust probes. I have not built this PR's full JNI library or run cloud-backed end-to-end tests. The Delta CI suites are green on Spark 3.5, 4.0, and 4.1. I am leaving the already-acknowledged claim-hook, naming, and packaging choices for the existing design discussion.
| let (dv_url, dv_store_path) = prepare_object_store_with_configs( | ||
| Arc::clone(&runtime_env), | ||
| dv_path.clone(), | ||
| object_store_options, | ||
| )?; |
There was a problem hiding this comment.
[P2] Avoid constructing a cold S3 store inside the DV runtime
Could we resolve the required stores before entering attach_access_plans, or make their initialization async-safe? The caller enters get_runtime().block_on(...), but an uncached S3 sidecar reaches this synchronous helper and then objectstore/s3.rs calls get_runtime().block_on(build_credential_provider(...)) again. Tokio rejects that nested Handle::block_on with a panic. A fresh executor reading a shallow clone whose data is in bucket A and whose new DV is in bucket B reaches a cold cache entry. Same-bucket tests hide the problem because the data store was created before the outer block_on. Explicit endpoint/region or static Hadoop credentials do not avoid the inner credential-provider call. Please add a test with distinct data-file and DV buckets.
There was a problem hiding this comment.
Fixed by pre-resolution: all stores (data and DV authorities) are resolved on the JNI thread before entering the runtime, and attach_access_plans no longer takes an options map or imports the store builder at all, so the async path structurally can't construct one. Your distinct data/DV bucket scenario is encoded in a new MinIO suite (CometDeltaS3Suite), but heads up that it's docker-gated and hasn't run against a live daemon yet, the contrib CI job has no docker socket so those tests cancel. First live signal needs a Docker environment.
|
Thanks @dwsmith1983. On the design topics you flagged, I’d prefer using |
sunchao
left a comment
There was a problem hiding this comment.
Rechecked 7e09e04f with five independent review scopes. One additional P2 is inline; I also followed up in the existing threads on the remaining scalar-pushdown, Azure DV store, and DV-memory issues. Verification used exact-source Spark/Delta physical-plan probes and locked-dependency Rust probes, not a full Comet/JNI or live cloud run.
|
Hey @sunchao / @parthchandra — are you looking to move this work over to @dwsmith1983’s series? We’re 2 PRs into the 10-PR series now that the contrib modules have been merged. The series fully implements all of the Delta protocols, with all 10k+ Delta test cases passing. I’m fine either way. I won’t have a huge amount of time to work on this over the next couple of months, so it could go faster with a different attendant, but the PRs are ready to roll. You can see the series here: https://github.com/schenksj/datafusion-comet/pulls (PRs #5–13). |
|
Hi @schenksj , I think your series implements Delta native scan based on the I think your series is pretty valuable and should be continued to push forward. At some point we should compare feature coverage and performance between the two. |
|
CI notes for this push: the S3 test base built its client without a region, which aborted CometDeltaS3Suite in CI's empty AWS environment before any test ran; fixed, and since GitHub mounts the Docker socket into job containers the MinIO scenarios will actually execute in CI now. They've been run live locally with all AWS env vars unset (first executions ever, both pass on Spark 3.5 and 4.1; that surfaced a missing spark-hadoop-cloud test dependency, also fixed). Heads up that inside the job container the MinIO endpoint may resolve as unreachable sibling-container networking; the suite now fails soft to canceled rather than aborting the build, and logs the resolved endpoint so the first CI run tells us whether a testcontainers host override is needed. The Spark 4.0 cell wasn't rerun locally, so CI is its first pass over these changes. The Rust 1.98 clippy fix I'd pushed got dropped in favor of #5400 from main during rebase. |
|
Reposting the two remaining P2 findings here for visibility. Both remain present at [P2] Check selected-file schemes before claiming a shallow cloneThe filesystem gate checks only the table's This was verified with a real Spark 4.0.2 / Delta 4.0.0 shallow clone that Spark successfully reads, plus the exact native store-preparation helper. Please apply the supported-scheme check to the selected data-file URIs before claiming the scan. Code · Existing discussion and reproduction details [P2] Account for the DV reader's combined-selection allocationConstruction admission and the initial reader clone are now covered. However, DataFusion 54.1 subsequently calls With the default-permitted 1,000,000 alternating deletions across 2,000,000 rows, the current attachment reserves 64,000,000 bytes, but the attached selectors plus reader-normalization allocations peak at 97,554,457 bytes and retain 65,554,432 bytes afterward. Please account for normalization and vector capacity, or avoid the additional allocation through ownership transfer. Simply changing the factor to 3 would still fall below this measured peak. This was reproduced using the unchanged attachment code and the real locked dependency conversion. These are allocator-requested bytes, not RSS or a reproduced executor OOM. Both findings were checked with focused probes and source tracing, not a full Comet/JNI integration run. |
I see value in both (even though it is extra work to maintain both paths) and in principle agree with @sunchao. Ideally, we want to converge these two. Logged an issue based on an AI generated convergence path - #5411 |
Thanks guys. I'm concerned that having 2 will create a lot of confusion when it comes to support.. Even enabling and disabling various scan features is too much to understand for most of the expert data engineers I work with every day. I'm happy to move forward initially in parallel, though like I mentioned before my time to work with this is going to be pretty sparse for the next couple of months. |
|
@schenksj Let’s see how it goes. For now, I see the In terms of your concern, I think we should aim to keep the user-facing configuration simple, perhaps with one flag to enable Delta scans and another to opt into an experimental Rust-kernel-backed path. Ideally, both approaches would share as much integration and testing infrastructure as possible. Really appreciate all your work on this! We’re planning to move quickly with the current |
|
On the macOS scans failure: pulled the hs_err from the run artifact. The crashing thread is a native thread (not a Java thread) that was exiting: the stack is pthread_start into pthread_exit into pthread TSD cleanup, then a jump through a corrupted destructor slot whose value is ASCII string bytes, at 119s elapsed, immediately after ParquetReadFromFakeHadoopFsSuite, the only suite in the group that exercises the libhdfs bridge and its JNI-attached native threads. The Delta code in this PR is structurally unreachable in those suites (native side is dispatch-gated on an operator those plans never emit, and the contrib jar is not on that build's classpath), and the Linux scans group passed on the same commit. My guess is a teardown race in the libhdfs bridge or a runner flake rather than anything this PR executes; the falsifying experiment would be rebuilding the dylib without the delta feature and re-running, since the same crash would exonerate it by construction. Could someone re-run the job? Happy to file the hs_err as an issue either way. |
sunchao
left a comment
There was a problem hiding this comment.
Thanks for addressing the earlier findings. I think the larger file-selection refactor, shared admission/schema cleanup, packaging changes, and broader deployment coverage can be tracked in follow-up PRs. I'd keep the remaining [P1] Azure safety guard, [P2] S3-authentication and AQE lifecycle fixes, and their focused regressions in this PR.
Could we replace spark.comet.scan.deltaNative.enabled with spark.comet.scan.delta.enabled consistently across both Delta contributions, keeping the default false? Please update the config definitions, tests, documentation, and dev scripts together, and use the spark.comet.scan.delta.* prefix for related settings. The intent is one consistent configuration namespace, not another enable flag.
This rename does not depend on changing the separate-JAR packaging. Broader reader-selection behavior can be discussed separately.
| override lazy val outputPartitioning: Partitioning = | ||
| UnknownPartitioning(perPartitionData.length) |
There was a problem hiding this comment.
[P2] Avoid executing adaptive pruning while inspecting partitioning
This getter forces perPartitionData, which calls InSubqueryExec.updateResult(). During AQE, that subquery can still be a non-executable adaptive broadcast placeholder.
A reduced Spark 4.0.2 / Delta 4.0.0 planning harness reproduced this through Spark's normal AQE validation: a DPP join in one UNION ALL branch and a coalescible shuffle in another caused validation to inspect this partitioning before the custom DPP rewrite. It then failed with CometSubqueryAdaptiveBroadcastExec ... does not support the execute() code path. Other operators remained on Spark, and no native Comet reader executed.
Could we return UnknownPartitioning(0) while adaptive placeholders remain and make this a non-lazy def, so the temporary value is not cached? A regression with a query-time dimension filter would help. The current DPP test filters the dimension before writing it, so it does not require dynamic pruning.
There was a problem hiding this comment.
Fixed as described: outputPartitioning is now a plain def returning UnknownPartitioning(0) while any runtime filter still holds an adaptive broadcast placeholder, so AQE validation never forces perPartitionData. Rewrote the DPP test to filter at query time and added your UNION ALL shape as a regression. That shape didn't reproduce the crash pre-fix on my Spark 3.5.9 / Delta 3.3.2 profile, so it likely needs your Spark 4.0.2 harness, but the guard matches your analysis.
|
Agreed on keeping it simple. The conf is now spark.comet.scan.delta.enabled (plus spark.comet.scan.delta.dv.maxDeletedRowsPerFile), so there's one flag to enable Delta scans, and the kernel path can add its own experimental key later. Docs updated. Fixes for the three open threads are pushed as well. |
ec2ad9b to
92ae71b
Compare
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: DSv1 Delta reads lacked this JVM-planned native scan path and its Parquet pruning capabilities.
- Design approach: Keep snapshot resolution and file planning in delta-spark, then reuse Comet’s native Parquet reader through the contrib interface.
- Correctness / compatibility analysis: Found one introduced P1 compilation failure. Compared rebasing, metadata ordering and subquery behavior against Spark 3.5.9, 4.0.4 and 4.1.3, and Delta reader behavior against the supported pairings. The previous deletion-vector consumer-count concern is addressed.
- Key design decisions: Require the contrib jar and explicit configuration, decline unsupported scan shapes, and keep rebasing specific to Delta. Shared scan builders reduce duplication. Per-file reservations now share one consumer per attachment call.
- Implementation sketch: Serialize common scan settings and partition-specific files through
DeltaSparkScan, apply deletion vectors as row selections, and resolve runtime filters before execution. - Behavioral changes worth calling out: Splits each prepare whole-file deletion-vector state. Rebasing wrappers prevent pruning on affected columns, while modern-value pass-through preserves buffers. Ancient legacy timestamps with unsupported writer zones fail explicitly. No fresh performance benchmark was run.
- Suggested improvements: Update the three stale test callers identified below, then rerun native checks and obtain the Delta integration CI verdict.
Reviewed the full 60-file diff from 4453c57249aa4bf9c47fd4ff88bd13a09cd2a9f8 to c0351d0002d5a15469a6587b50b4d5c299d35249. The PR remains open and non-draft. Read existing discussion and review threads, excluding Copilot. Routed skills: review-comet-pr, review-comet-expression-pr, review-comet-ffi-pr, and review-comet-memory-pr.
Exact-head CI: labeling passed. Comet CI and CodeQL remain action_required. There is no product test verdict and no run-delta-tests label.
Validation: the exact-head native test build fails with three E0061 errors. After supplying only the missing arguments in three test helpers in a disposable copy, 125 focused rebasing, deletion-vector, planner and Parquet tests passed. Those passes do not qualify the unmodified head. CI configuration, suite-registration and diff checks passed. Full JVM/Spark/Delta, MinIO and release-build validation was not run. The checkout remains unchanged.
| encryption_enabled: bool, | ||
| use_field_id: bool, | ||
| ignore_missing_field_id: bool, | ||
| rebase_from_file_metadata: bool, |
There was a problem hiding this comment.
[P1] Update every caller when adding these three parameters. Three helpers under native/core/src/execution/operators/dynamic_filter/join/tests/ still use the 18-argument signature: timestamp_errors.rs:86, schema_errors.rs:65, and schema_errors/partition_columns.rs:47. Building any core unit-test selection now fails because Rust compiles these helpers even when their tests are filtered out. This also blocks CI’s cargo clippy --all-targets --workspace. Supply false, "", "" at these ordinary-scan test callers, matching the other updated fixtures.
Evidence: At the unmodified reviewed head, running cargo test --locked --offline -p datafusion-comet --no-default-features --features delta datetime_rebase --no-run from native/ exits 101 with three E0061 errors: “this function takes 21 arguments but 18 arguments were supplied.” The base signature accepts 18 arguments. Adding only the three missing arguments at those test call sites in a disposable copy allows the test binary to compile.
There was a problem hiding this comment.
Fixed in 281e02f01, the merge with main that landed just after this review. timestamp_errors.rs, schema_errors.rs and schema_errors/partition_columns.rs now pass false, "CORRECTED", "CORRECTED", the same values the neighbouring ordinary-scan fixtures use, and cargo clippy --all-targets and the core test build compile again.
…scan The native scan's shared helpers carry main's Variant projection changes: existing default values go through serializeExistenceDefaultValues and the planner's shared default parsing takes main's bounds and conversion checks, so the Delta arm gets them too. The scan input gate is main's CometLeafExec check, which covers the Delta scan. Variant stays declined on the Delta path, which lacks core's Variant gates. Three join filter test call sites gain the rebase arguments the branch added to init_datasource_exec.
…scan Main's parquet field id change replaced schema_adapter's parse_field_id with parquet_support::field_id, so the datetime rebase code reads field ids through that shared helper now.
…scan Main now checks a file without field ids against a require_field_ids flag the serde computes from the required schema. The flag is set in the shared populateScanConfFlags and passed through build_parquet_scan_plan, so the Delta scan applies the check as Spark does for Delta reads, on the physical schema. The deletion vector shape now hands the physical schema to that helper, as the ordinary shape already did, so a name-mapped table does not require ids Spark strips. The reader factory runs the check before the INT96 stamp.
|
@andygrove could you approve a CI run on |
…scan Main now builds the empty-partition EmptyExec from the projected table schema, partition and constant metadata columns included, so fused operators bind the same way when a bucket is empty. That lives in the shared build_parquet_scan_plan, so an empty or fully pruned Delta partition gets it too. The native RDD keeps main's compute override for scan input metrics.
|
@andygrove could you approve a CI run on |
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: DSv1 Delta reads lacked this JVM-planned native scan path and its Parquet pruning capabilities.
- Design approach: Keep snapshot resolution and file planning in delta-spark, then reuse Comet’s native Parquet reader through the contrib interface.
- Correctness / compatibility analysis: Found one new P2: INT96-backed
TIMESTAMP_NTZreads can incorrectly fail calendar-rebase checks. Also reproduced the previously reported S3 encryption configuration discrepancy described below. Compared relevant Spark 3.5.9, 4.0.4 and 4.1.3 reader semantics and Delta 3.3.2, 4.0.1 and 4.3.1 deletion-vector formats. - Key design decisions: Shared scan builders reduce duplication. The contrib jar and configuration gate JVM use, while calendar rebasing remains specific to Delta. Storage admission mirrors Hadoop configuration semantics, but the encryption gate still uses the wrong configuration view.
- Implementation sketch: Serialize common scan settings and partition-specific files through
DeltaSparkScan, attach deletion vectors as Parquet row selections, and wrap affected columns with per-file rebase expressions. - Behavioral changes worth calling out: The default native build includes the
deltadecoder. DV plans retain reserved memory until task completion, and each split prepares whole-file DV state. Ancient timestamps with unsupported writer zones fail explicitly. No fresh performance benchmark was run. - Suggested improvements: Make timestamp rebasing respect the requested Spark type, and resolve the existing encryption gate against Hadoop’s propagated bucket configuration.
Reviewed the entire 62-file diff from ce455f32d948355e073638e81009ecc3e5dea349 to 415d7298728e1909099b4158f73366bf7f5af0d4. The PR remains open and non-draft. Read existing reviews, issue comments and review threads, excluding Copilot. Routed skills: review-comet-pr, review-comet-expression-pr, review-comet-ffi-pr, and review-comet-memory-pr.
An existing P2 remains unresolved in the SSE-C discussion, specifically the propagated-reference case already identified in review 5078484634. Set fs.s3a.custom.ref=AES256, fs.s3a.bucket.data-bucket.custom.ref=SSE-C, and fs.s3a.encryption.algorithm=${fs.s3a.custom.ref}. Current source-extracted admission helpers return None, while Hadoop’s propagateBucketOptions resolves the algorithm to SSE-C. An otherwise eligible SSE-C table can therefore reach a native client that omits the required customer-key headers. The configuration reproduction used Hadoop 3.4.2, with the same propagation-before-encryption ordering verified in supported Hadoop 3.3.4 and 3.4.1 sources. No cloud read was attempted. This is recorded as an existing blocker, without a duplicate inline finding.
Exact-head CI: labeling passed. Comet CI and CodeQL report action_required. There is no product-test verdict and no run-delta-tests label.
Validation: unmodified-head native tests passed with --locked --offline --no-default-features --features delta: 41 datetime_rebase, 64 delta_, and 28 parquet_exec tests. CI configuration, suite-registration and diff checks passed. A disposable native scan reproduction demonstrated the new failure, and stock Spark 4.1.3 successfully read the same file. Full Comet JVM/Delta integration, Spark SQL suites, MinIO, default-feature/HDFS and release-build validation were not run. The project checkout remains unchanged. Nothing was published.
| out.push(policies.date_policy(*next_leaf)); | ||
| *next_leaf += 1; | ||
| } | ||
| DataType::Timestamp(_, Some(_)) => { |
There was a problem hiding this comment.
[P2] Respect the requested TIMESTAMP_NTZ type when selecting timestamp rebase policies. On Spark 4.0/4.1, an opted-in Delta scan requesting TIMESTAMP_NTZ from an INT96 file containing 1800-01-01, with spark.sql.parquet.int96RebaseModeInRead=EXCEPTION, should return the value unchanged. Spark explicitly skips rebasing for requested NTZ timestamps, and Delta’s protocol permits INT96 storage for them. Here, INT96 is physically represented as Timestamp(_, Some("UTC")), so this branch applies the ancient-value check beneath the NTZ conversion and fails the query. Please use the physical-to-requested leaf pairing to suppress timestamp rebasing for NTZ requests, while preserving DATE-to-NTZ rebasing, and add this regression case.
Evidence: Executed a disposable test appended to an otherwise unchanged copy of the head’s parquet_exec.rs: review_ntz_int96_ignores_rebase_mode in /tmp/comet5365-ntz-review-l5jqmu1i/native/core/src/parquet/parquet_exec.rs. It writes metadata-free message m { required int96 ts; } with INT96 words (0, 0, 2378497), then calls the current scan builder with requested Timestamp(Microsecond, None), LTZ-to-NTZ conversion enabled and metadata rebasing enabled. CORRECTED returns 1800-01-01T00:00:00; EXCEPTION returns Native scan cannot rebase ancient values in column 'ts'. Stock Spark 4.1.3, reading the same /tmp/comet5365-review-int96-ntz.parquet with .schema("ts TIMESTAMP_NTZ") and the EXCEPTION setting, returned [1800-01-01T00:00]. Spark 4.0.4/4.1.3 ParquetVectorUpdaterFactory selects BinaryToSQLTimestampUpdater without rebasing for this requested type. Delta 4.3.1’s DeltaParquetFileFormat forwards the requested schema and options to that reader. Full Comet/Delta/JNI execution was not run.
There was a problem hiding this comment.
Fixed in 4b7ad02. A leaf read as TIMESTAMP_NTZ now takes no rebase policy whatever its physical type, so an INT96 or adjusted INT64 column read as NTZ returns the value unchanged under every mode. DATE read as TIMESTAMP_NTZ keeps the date policy, as DateToTimestampNTZWithRebaseUpdater does on 4.x.
Your case is timestamps_requested_as_ntz_are_never_rebased in parquet_exec.rs (INT96 words (0, 0, 2378497), requested Timestamp(Microsecond, None), EXCEPTION, CORRECTED and LEGACY all return 1800-01-01T00:00:00), and the Delta test "INT96 and adjusted INT64 timestamps read as TIMESTAMP_NTZ are never rebased" runs it through CONVERT TO DELTA on Spark 4.0.
|
This is a light fully automated review since there are so many PRs open.
|
The native Delta scan picked a timestamp leaf's rebase policy from its physical Arrow type, while Spark's ParquetVectorUpdaterFactory keys on the Spark type the column is read as. Two cases differed. An INT96 or adjusted INT64 column read as TIMESTAMP_NTZ was checked or rebased, though Spark never rebases NTZ reads, so an ancient value failed under EXCEPTION. A timezone-free INT64 column read as TIMESTAMP was never rebased, though Spark applies the datetime rebase mode whatever isAdjustedToUTC says, so EXCEPTION returned the row and LEGACY returned it unshifted. The walk that already pairs each physical leaf with its requested type now records the leaves read as TIMESTAMP_NTZ and the timezone-free leaves read as TIMESTAMP, and the policy follows that record. Dates keep the date policy, including a DATE column read as TIMESTAMP_NTZ, which Spark 4.x rebases. A metadata-free timezone-free file read as TIMESTAMP under LEGACY is now refused, as adjusted files without a recorded writer time zone already are, because the rebase needs the JVM's time zone tables.
Yes, in 4b7ad02. The walk that pairs each physical leaf with its requested type now records the timezone-free leaves read as Under The regression you suggested is |
Which issue does this PR close?
Part of #174. This PR does not close it: the delta-kernel contrib and the convergence discussion in #5411 are tracked there as well.
Rationale for this change
Adds an optional contrib module that plans Delta Lake table scans on the JVM and executes them natively, including deletion vector application inside the native scan. delta-spark has already done log replay, snapshot resolution, and partition pruning by the time CometScanRule sees the FileSourceScanExec, so there is no Delta planning to do natively: the scan reuses the existing ParquetSource path and gets row group pruning, page index pruning, and filter pushdown for free, with deletion vectors composed into the ParquetAccessPlan so DV skips and page skips intersect rather than filtering after the read.
The module is explicit opt in: the
-PdeltaMaven profile builds a separatecomet-contrib-deltajar that is never bundled intocomet-spark, andspark.comet.scan.delta.enableddefaults to false. Thedeltacargo feature (DV decoding plus the planner hand-off, no delta-kernel dependency, about 82 KB of dylib) stays in the default native build so trying the contrib needs only the jar and the config, not a custom native binary; this was agreed in review and is recorded in the Cargo.toml comment. The adjacentcontrib-deltafeature is unrelated: it gates the delta-kernel integration and default builds carry no kernel surface.Restructured after review
Core changes that previously traveled with this PR now live elsewhere:
${...}expansion, constant metadata field uniquification, and dead JNI removal: fix: expand object store option references, uniquify constant metadata names, drop dead parquet JNI #5653, now merged. The first two are prerequisites of this module and this branch is rebased on top of them.Two core-generic capabilities remain in this PR because the native read path does not have them yet and the Delta scan needs them for correctness; both are candidates to lift into core, tracked in #5662 (S3 configuration divergence for the regular native scan) and #5010 (calendar rebasing for the regular native scan):
fs.s3a.assumed.role.policy) decline outright since Hadoop sends them in the AssumeRole request and native does not.What changes are included in this PR?
contrib/delta-spark: DeltaScanSupport (scan eligibility, S3 divergence gating, DV descriptor extraction), CometDeltaNativeScan serde, service registration via the contrib scan SPI, documentation.delta_dv.rs(deletion vector decode with a full malformed input matrix, and access plan construction),delta_spark_scan.rsplanner arm,datetime_rebase.rs, proto messages for the Delta scan envelope, S3 object store helper.build_parquet_scan_plan/prepare_scan_store_and_filesextraction in the planner,object_store_url_key/prepare_object_store_with_config_hash,buildNativeScanCommonextraction,reportScanInputMetrics,hasScanInputwidening, contrib LinkageError containment.Follow-up work from review is tracked in #5655 (DV file splitting), #5656 (compressed DV decoding), #5657 (overlapping bitmap and footer reads), #5658 (shared cloud compatibility helper), #5659 (credential scoping), #5660 (v2 checkpoint coverage), #5661 (capability table), and #5662.
How are these changes tested?
--features delta(343 in the core crate), including the DV malformed input matrix (truncation at every boundary, CRC and magic corruption, size and cardinality lies, bit flip sweeps), the calendar rebase unit tests against Spark's own anchors, and end to end scan pins for per file metadata resolution; clippy and fmt clean.Benchmarks at the current head
Apple M5, JDK 17, Spark 3.5 profile, local filesystem, 120M rows in 6 files of about 490 MB (zstd), full table aggregate touching every surviving row, medians of 5 warm runs per fresh session. Results are bit identical across all modes and verified against closed form expectations.
DV decoding is negligible in every pattern; the cost center is selector expansion for alternating deletes (61 to 93 ms and about 400 MB peak per file). The default
spark.comet.scan.delta.dv.maxDeletedRowsPerFilecap (1M) declines the contiguous and alternating tables up front and falls back cleanly, which the numbers show is the better path for alternating; raising the cap without sizing the off heap pool fails tasks at the reservation by design.The calendar rebase wrapper costs 0.7 to 2.2 ns per row and is noise at scan level, but it is opaque to pruning: a selective predicate on a rebased column decoded 65x more rows than with pruning live on a sorted table. That is the tradeoff of the legacy path and only applies to files that need rebasing.
An independent run on public data (NYC taxi with a DV delete) is in the PR discussion and confirmed exact DV row removal with timing parity.