feat: build gate + inert wiring for contrib Delta scans [Delta contrib split, part 2] - #4952
feat: build gate + inert wiring for contrib Delta scans [Delta contrib split, part 2]#4952schenksj wants to merge 1 commit into
Conversation
|
@parthchandra @andygrove this is part 2 of the Delta contrib split — the build gate + inert wiring — following on from part 1 (#4700, now merged). Thank you both for the reviews on part 1! 🙏 Sorry it took a while to get this next phase out there — I was heads-down getting https://github.com/capitalone/vulnhunter to market. Back on the Delta series now. This part is deliberately inert/gated (zero Delta surface on default builds), so it should be a fairly self-contained review of the wire format, the build machinery, and the gate-enforcement script. Would appreciate your eyes when you have a chance. |
| * and the lookups resolve, dispatching the call into the contrib helpers. | ||
| * | ||
| * Keeping this bridge as one small file in core lets the Delta detection block in `CometScanRule` | ||
| * and the serde dispatch in `CometExecRule` stay ~10 lines each -- exactly the shape Parth's |
There was a problem hiding this comment.
The change to CometExecRule is small but this itself is pure Delta code in core. We really need to follow the same pattern we established of discovering using the ServiceLoader mechanism. PlanDataInjector (in operators.scala:72-127, merged in Part 1 / #4700) already uses ServiceLoader discovery. PR #4633 (Lance) also uses the same pattern with CometScanContrib
I would recommend we replace this with a generic CometScanContrib trait in a new file spark/src/main/scala/org/apache/comet/rules/CometScanContrib.scala (in core). This trait can cover both V1 and V2 scans:
trait CometScanContrib {
/** V1 scan hook. Return Some(plan) to claim, None to pass. */
def tryTransformV1(
plan: SparkPlan,
session: SparkSession,
scanExec: FileSourceScanExec,
relation: HadoopFsRelation): Option[SparkPlan] = None
/** V2 scan hook. Return Some(plan) to claim, None to pass. */
def tryTransformV2(scanExec: BatchScanExec): Option[SparkPlan] = None
}
object CometScanContrib extends Logging {
private lazy val contribs: Seq[CometScanContrib] = {
// Built-in contribs (Parquet, Iceberg) can be registered here.
// Contrib-gated ones (Delta, Lance) discovered via ServiceLoader.
val discovered = try {
ServiceLoader.load(classOf[CometScanContrib], getClass.getClassLoader).asScala.toSeq
} catch {
case NonFatal(e) =>
logWarning("Failed to load contrib CometScanContrib services", e)
Seq.empty
}
discovered
}
def tryTransformV1(...): Option[SparkPlan] = {
contribs.view.flatMap(_.tryTransformV1(...)).headOption
}
def tryTransformV2(...): Option[SparkPlan] = {
contribs.view.flatMap(_.tryTransformV2(...)).headOption
}
}
This is format-agnostic code.
The Delta contrib then has:
contrib/delta/src/.../DeltaScanRuleContrib.scalaimplementingCometScanContribcontrib/delta/resources/META-INF/services/org.apache.comet.rules.CometScanContribnaming it
This also subsumes PR #4633's CometScanContrib — the Lance PR would implement tryTransformV2 on the same trait rather than defining a separate one.
| // Delta Lake scan. Wire format used by `contrib/delta/`. Only decoded when | ||
| // core is built with `--features contrib-delta`; in default builds the | ||
| // dispatcher arm is `#[cfg]`-stubbed out so the contrib has zero runtime cost. | ||
| DeltaScan delta_scan = 118; |
There was a problem hiding this comment.
Currently both this PR and PR #4633 claim field number 118 (delta_scan = 118 vs lance_scan = 118). More fundamentally, adding a new oneof variant for every contrib scan type means the core proto file must change every time a new format is added — violating the goal that core doesn't know about contrib implementations.
Spark Connect solves this at relations.proto:109 with google.protobuf.Any extension = 998;. We should do the same.
We can replace this change with a permanent extension point:
// One-time addition to the oneof, never changes again:
google.protobuf.Any contrib_scan = 200;
Then DeltaScan and LanceScan message definitions still live in operator.proto (or preferably in separate proto files under contrib/) — they're just packed into Any on the JVM side and unpacked by type_url on the Rust side.
There was a problem hiding this comment.
If this gets tricky to implement we can consider assigning field ids for contrib scans, but I feel that the approach outlined here is feasible.
| case scan if !CometConf.COMET_NATIVE_SCAN_ENABLED.get(conf) => | ||
| withFallbackReason(scan, "Comet Scan is not enabled") | ||
|
|
||
| // V1 scans go through `transformV1Scan` which itself first delegates to any |
There was a problem hiding this comment.
With a generic CometScanContrib.tryTransformV1(), this reordering is unnecessary. The outer transformScan match keeps its current order (metadata-colum guard before FileSourceScanExec). Each contrib's tryTransformV1 can decide internally whether it can handle metadata columns.
| // vanilla scan path. When the Delta classes are on the classpath, the contrib | ||
| // either claims the scan (returning a CometScanExec marker) or declines via | ||
| // its own `withFallbackReason` fallback message. | ||
| DeltaIntegration.transformV1IfDelta(plan, session, scanExec, r) match { |
There was a problem hiding this comment.
This will change to CometScanContrib.tryTransformV1(plan, session, scanExec, r). We can also remove the re-applied metadataCols check — it's now redundant because the outer guard still runs first and the Delta implementation decides whether to handle metadata columns or not.
| // activated. The marker wraps the original, link-bearing scan, so the produced exec's | ||
| // originalPlan keeps its logicalLink with no workaround. If conversion declines, the marker | ||
| // itself falls back to the vanilla Spark Delta scan, so leaving it in the plan is safe. | ||
| case scan if DeltaIntegration.isDeltaScanMarker(scan) => |
There was a problem hiding this comment.
We can clean this up too - Define a marker trait and Delta can implement it
/** Marker trait for contrib scan nodes that carry their own serde handler. */
trait CometContribScanMarker { this: SparkPlan =>
def scanHandler: CometOperatorSerde[_ <: SparkPlan]
}
then the match becomes
case marker: CometContribScanMarker =>
convertToComet(marker, marker.scanHandler).getOrElse(marker)
|
Thanks for the thorough review, @parthchandra — the "core must stay format-agnostic" framing is What changed (per thread)
Two deliberate deviations1. Hand-rolled envelope instead of 2. One subtlety on the metadata-column guard (threads at
|
|
@schenksj did you forget to push the code ? |
…ack] Addresses @parthchandra's review on apache#4952. Every thread had the same theme: core must not name a specific contrib format. Replaces the Delta-specific core touchpoints with generic extension points, mirroring the ServiceLoader SPI established in part 1 (apache#4700, `PlanDataInjector`) and shared with the Lance PR (apache#4633). - Delete `DeltaIntegration.scala` (reflective bridge with cached `MODULE$` / `getMethod` lookups). Replaced by `CometScanContrib`: a `trait` with `tryTransformV1` / `tryTransformV2` (both defaulting to `None`) plus a ServiceLoader-backed object, discovered exactly like `PlanDataInjector`. Default builds ship no `META-INF/services` entry, so the registry is empty and both hooks are inert. Both hooks are wired for real -- `tryTransformV1` at the top of `transformV1Scan`, `tryTransformV2` at the top of `transformV2Scan` -- so a V2 contrib (Lance) is consulted too; this trait subsumes the one apache#4633 was defining. - Add `CometContribScanMarker`, a marker trait carrying its own `scanHandler: CometOperatorSerde[_ <: SparkPlan]`. `CometExecRule` is now a plain type test instead of a class-name match plus a reflective handler lookup. It `extends SparkPlan` rather than using a `this: SparkPlan =>` self-type: a self-typed trait value is not a `SparkPlan`, so `convertToComet(marker, ...)` and `getOrElse(marker)` would not typecheck. - Proto: replace the contrib-specific `DeltaScan delta_scan = 118` oneof variant with a single permanent `ContribScan contrib_scan = 200` envelope (`type_url` + packed `value`), and `reserved 118`. Core's oneof never grows per-contrib again, and the 118 collision with apache#4633's `lance_scan` is gone. The envelope is hand-rolled rather than `google.protobuf.Any` because Comet compiles this .proto with two toolchains and the Maven `protoc-jar` plugin cannot resolve the bundled well-known types (`includeStdTypes` NPEs inside the plugin). Field layout is identical to `Any`, so the JVM can populate it from `Any.pack(...)`. - Native: `OpStruct::ContribScan` is routed by `type_url` to the gated `delta_scan::try_plan_contrib_scan`, which claims only its own type and decodes `DeltaScan` itself -- core names no contrib type. A default build reaching a `contrib_scan` gets a clear, `type_url`-identifying error. - `CometScanRule`: outer `transformScan` match order restored to match main, and the redundant re-applied metadata-column guard dropped. Verification: default + `contrib-delta` cargo builds, clippy both feature states, `dev/verify-contrib-delta-gate.sh` (default libcomet: 0 Delta symbols), JVM compile on spark-3.4/Scala 2.12 and spark-3.5/Scala 2.13, spotless and scalastyle -- all green. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
|
Apologies @parthchandra — you're right, that one's on me. My reply above walked through the refactor, but the push silently didn't go through on my end. It's up now: commit d1783f2 ("make core contrib-scan wiring format-agnostic") on this branch, which contains everything in that walkthrough. Thanks for the nudge. 🤖 This reply was drafted with Claude Code. |
parthchandra
left a comment
There was a problem hiding this comment.
Thanks Scott, this is great. I have some more comments, but we are getting close.
| * than swallowing silently) keeps an unexpectedly-declining contrib diagnosable. `NonFatal` | ||
| * deliberately lets `LinkageError`/`OOM`-class failures through. | ||
| */ | ||
| private def firstClaim(hook: CometScanContrib => Option[SparkPlan]): Option[SparkPlan] = |
There was a problem hiding this comment.
This is nice. Can we add a small test suite similar to CometScanWithPlanDataSuite?
Some cases to consider (all running on the default build) -
- an empty registry returns None from tryTransformV1/tryTransformV2
- a stub contrib registered through a URLClassLoader that claims a scan is returned
- a stub that throws is logged and declined so the next contrib still gets a look
Also, can two contribs claim the same scan?
There was a problem hiding this comment.
There is a corner case here. A future V2 contrib that has a table name like files, snapshots etc (iceberg reserved names) would pass the isIcebergMetadataTable check and the contrib will never get called. We could either call transformV2Scan first or explicitly check for an iceberg scan in isIcebergMetadataTable ?
| * | ||
| * Both hooks default to `None` so a contrib overrides only the scan kind(s) it handles. | ||
| */ | ||
| trait CometScanContrib { |
There was a problem hiding this comment.
Can we add a doc note that an implementation MUST return None for a scan it does not own. This is to prevent two contribs from competing for the same scan.
| case Some(handled) => return handled | ||
| case None => // proceed with vanilla logic | ||
| } | ||
| if (metadataCols(scanExec).nonEmpty) { |
There was a problem hiding this comment.
Can we add a test case that covers metadata for V1 scans? say selecting _metadata.file_path for a parquet source. The default build should fall back with this reason.
| // ===================================================================================== | ||
|
|
||
| // Per-scan invariants. Lives at the head of every Delta scan operator payload. | ||
| message DeltaScanCommon { |
There was a problem hiding this comment.
Can we log a follow up issue to move this out of core and into contrib/delta/proto? We will probably need to add a (manual) pipeline for the contribs.
| // itself is unconditional so a default build that receives a contrib-shaped plan | ||
| // from a misconfigured driver gets a clear error instead of a "no match" decode | ||
| // failure. | ||
| #[cfg(feature = "contrib-delta")] |
There was a problem hiding this comment.
I know this is based on the previous review comment but on deeper thought it may be possible to remove this from core as well by using a generic handler (similar to the jvm side).
Could you log a follow up issue for this as well? We have two paths we can consider -
- a true service loader type system for dynamically discovering and loading an extension. There may be dragons along this path.
- a statically linked version which builds each crate independently but may need some refactoring at the crate level to avoid circular dependencies.
There was a problem hiding this comment.
Rust dynamic loading ref: https://nullderef.com/blog/plugin-dynload/
…b split, part 2]
Part 2 of the Delta contrib split: the build gate and the inert core wiring an
out-of-tree scan contrib plugs into. Nothing here is reachable on a default
build -- no contrib is registered, no contrib class is compiled, and the native
library carries zero contrib symbols.
Core gains two format-agnostic extension points, both discovered at runtime so
core holds no compile-time reference to any contrib:
- `CometScanContrib`, a ServiceLoader-discovered hook (mirroring
`PlanDataInjector`) that lets a contrib claim a V1 or V2 scan before
Comet's built-in handling runs, plus `CometContribScanMarker` so
`CometExecRule` can route a contrib's scan node to the contrib's own serde
handler by a plain type test.
- `ContribScan contrib_scan = 200`, a single permanent `Any`-shaped proto
envelope (`type_url` + packed `value`) dispatched by `type_url` on the
native side. Core's oneof never grows per-format, so independent contrib
PRs cannot collide on a field number -- as `main` taking field 118 for
`Sample` has since demonstrated.
Plus the build machinery: the `contrib-delta` Maven profile and Cargo feature,
and `dev/verify-contrib-delta-gate.sh`, which asserts a default build compiles
no contrib classes, packages no contrib `META-INF/services` files, and links no
contrib symbols.
Where the hooks sit, and why. Both run *before* Comet's built-in guards for
their scan kind, because a contrib may support things the built-in scan does
not -- the Delta contrib synthesises `_metadata.*` in its own reader, and a
contrib's table name may end in `files`/`snapshots` like an Iceberg metadata
table. Applying those guards first would decline such a scan before the contrib
was ever offered it. So `transformV1Scan` consults the contrib ahead of the
metadata-column guard, and the Iceberg metadata-table check moves out of the
outer `transformScan` match into `transformV2Scan`, after its hook. Core's
per-path metadata handling is otherwise untouched: `main` serves
`fileConstantMetadataColumns` natively in V1 and the Iceberg metadata columns
in V2, and both keep doing so.
Ownership contract. An implementation MUST return `None` for a scan it does not
own: contribs are offered a scan one at a time and the first claim wins, so a
contrib claiming another format's scan hides it from the contrib that could
have read it, with the outcome depending on unspecified ServiceLoader ordering.
"Own but cannot handle" is a distinct, expressible case -- claim the scan and
terminate it with `withFallbackReason` rather than declining. Core cannot
arbitrate competing claims (a claim is opaque; the only way to know a second
contrib would also have claimed is to ask it, which is what claiming prevents),
so the contract carries it.
Tests. `CometScanContribSuite` covers the registry contract on a default build:
no contribs registered (asserted against raw ServiceLoader discovery, not just
the registry -- `contribs` swallows a ServiceConfigurationError, so "empty"
alone is ambiguous), a stub discovered through a URLClassLoader whose claim is
returned, decline-passes-through, first-claim-wins with later contribs not
consulted, throw-is-a-decline, and LinkageError still propagating.
`CometScanRuleSuite` gains a V1 case asserting the fallback *reason* for
`_metadata.row_index`; verified red with the guard removed.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
d1783f2 to
eebc232
Compare
|
Thanks @parthchandra — and apologies for being hard to reach lately. Hoping we can get our meeting All six threads are addressed, plus a rebase onto current The rebase (worth reading first)This branch was 141 commits behind. Rebasing changed two substantive things: 1. Field 118 is now 2. That changes the shape of the fix for the metadata thread — see below. Per-thread
I also documented "own but cannot handle" as a distinct, expressible case: return Can two contribs claim the same scan? In principle yes, and core cannot detect it: a claim is
Test added to
One thing that test found: asserting "the default build registers no contribs" against the
ValidationOn the rebased branch:
🤖 This reply was drafted with Claude Code. |
Part 2 of the Delta Lake contrib PR breakup. Part 1 (#4700, the core SPI for contrib leaf scans) is merged. This part establishes the
contrib-deltabuild gate and the inert wiring that lets a gated build compile and link end to end — while the default build stays byte-for-byte unchanged (zero Delta surface). It ships no real Delta read logic: a Delta read that reaches native returns a clean "not implemented" error and falls back to vanilla Spark. The full sequence and dependency graph live in the tracking umbrella, #4366.The whole point of this part is that everything it adds is gated or inert, so it is safe to land on
mainwell ahead of the read path, and reviewers can review the wire format, the build machinery, and the gate enforcement once, in isolation, before any Delta code shows up.Changes
Build gate
contrib-deltaprofile (spark/pom.xml) with a per-Sparkdelta.version(3.5 → 3.3.2, 4.0 → 4.0.0, 4.1 → 4.1.0) and anadd-sourceofcontrib/delta/src. Defaultdelta.versionfloor inpom.xml. The defaultspark.versionstays 4.1.2 — the delta-spark 4.1.1 pin is a separate, deferred decision (raised for later parts).contrib-deltafeature on core (optional path dep oncomet-contrib-delta);native/Cargo.tomlexcludes../contribfrom the workspace so non-Delta committers never build it.dev/verify-contrib-delta-gate.shproves the default cargo tree, Maven dependency set, compiled classes, andlibcometall carry zero Delta surface, and that the gated build pulls the right deps (delta-spark per Spark profile,comet-contrib-deltain the cargo tree). Wired into a minimaldelta_build_gate.ymlCI job. The full test-suite and regression workflows land in later parts.Inert wiring
Delta*messages anddelta_scan = 118(117 isBroadcastNestedLoopJoin). One-time wire-format review here, early.OpStruct::DeltaScanarm that returns a not-compiled-in error on default builds, and a feature-gateddelta_scanshim that calls the contrib; exhaustive-match arms inoperator_registry/jni_api;convert_spark_types_to_arrow_schemapromoted topub(crate).contrib/delta/native):plan_delta_scanreturnsDataFusionError::NotImplemented— just enough to satisfy the core shim's contract so--features contrib-deltalinks. This makes the exact core↔contrib contract visible in a small PR.DeltaIntegration(reflective; every lookup returnsNoneuntil the contrib classes exist), theCometExecRuleDelta-marker hook (the CDF hook is deferred to a later part), theCometScanRuleDelta delegation + metadata-column reorder, and the leafDeltaConf.What this part deliberately does NOT do yet
plan_delta_scanis a stub. Log replay, predicate pushdown, deletion vectors, and the kernel read land in parts 3a/3b.DeltaIntegration's lookups resolve toNone, so the marker always falls back. The claim/decline layer and native exec land in parts 4a/4b.CometExecRuleCDF hook andDeltaIntegration's CDF members are held back to part 5.Why it is safe on default builds
Without
-Pcontrib-delta: the cargo feature is off (nocomet-contrib-delta, nodelta_kernelin the tree), the Maven profile is inactive (noio.delta:*, nocontrib/deltasources compiled), and the nativeDeltaScanarm is#[cfg]-stubbed to an error that is never reached because nothing emits the proto message.DeltaIntegrationis the only always-present class and every one of its reflective lookups returnsNone. The gate script asserts all of this mechanically: defaultlibcomethas 0 Delta symbols and is the same size as before.The
CometScanRulechange reorders the metadata-column guard so V1 scans reachtransformV1Scan(which delegates to any V1 contrib) before the generic metadata-column rejection — but for a non-contrib V1 scan the guard is re-applied insidetransformV1Scan, so vanilla behavior is unchanged.Verification
Run against this branch rebased on current
main:--features contrib-delta): green.cargo clippyboth feature states: clean.dev/verify-contrib-delta-gate.sh: all checks pass (0 Delta symbols in the default dylib).-Pcontrib-delta) and default JVM compile (spark-3.4 / Scala 2.12): both compile.Roadmap
Parts still to come: Rust driver-side planning (3a), Rust executor-side read path (3b), Scala claim/decline (4a), Scala execution — end-to-end native reads (4b), Change Data Feed (5), test battery + regression harness (6), and docs (7). Each is gated behind
-Pcontrib-delta, so every intermediate state onmainis safe for default builds. Tracking umbrella: #4366; part 1: #4700.🤖 AI disclosure: this PR was prepared with assistance from Claude Code (Claude Opus 4.8), under the submitter's review and direction.