These commits are when the Protocol Buffers files have changed: (only the last 100 relevant commits are shown)
| Commit: | a2b688b | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
feat: custom Rust UDFs via arrow-ffi [experimental] (#4459) * feat: custom Rust UDFs via arrow-ffi (alternative to #4283) [skip ci] Adds custom Rust scalar UDF support using arrow-only FFI surfaces, as suggested by paleolimbot and timsaucer in feedback on #4283. This is an alternative implementation for comparison; see #4283 for the bespoke-ABI version. Two ABI flavors are provided side-by-side so reviewers can compare: 1. C ABI (sedona-style): pure C-callable struct of function pointers, parameterized only by Arrow C Data Interface (FFI_ArrowSchema / FFI_ArrowArray). Decoupled from datafusion versions; future-portable to C/C++. Modeled on apache/sedona-db's SedonaCScalarKernel header. 2. datafusion-ffi (FFI_ScalarUDF): wraps user's ScalarUDFImpl as FFI_ScalarUDF. Inherits full ScalarUDFImpl surface (variadic signatures, type coercion, metadata-aware return types) for free, at the cost of a major-version pin against datafusion-ffi. A single library may export either or both; loader walks both discovery functions. `comet-test-udfs` exposes `add_one_c` (C ABI) and `add_one_df` (datafusion-ffi) and the e2e suite drives both through Spark. Scope is intentionally minimal for comparison: scalar-only, happy-path e2e tests only (no panic/error/signature-mismatch coverage). The Scala JVM API mirrors #4283 exactly so the comparison is apples-to-apples. Pieces: - native/comet-udf-sdk: SDK with both ABI flavors + export macros - native/comet-test-udfs: cdylib exposing one UDF per ABI - native/core/src/execution/rust_udf: loader, cache, ImportedCScalarUdf - native/core/src/comet_rust_udf_bridge.rs: JNI for validateLibrary/listUdfs - native/proto: RustUdfCall message - spark/.../udf: CometRustUDF.register, registry, exception classes, JNI bridge stub - spark/.../serde/CometScalaUDF: dispatch ScalaUDF to RustUdfCall when udfName is in the registry * feat: drop datafusion-ffi flavor, contain panics at the FFI boundary, add user guide Keep a single UDF ABI built only on the Arrow C Data Interface. Carrying two public ABIs would mean two permanent compatibility promises; the datafusion-ffi flavor pins every user cdylib to Comet's DataFusion major version, which the 53 -> 54 upgrade demonstrated by removing as_any from ScalarUDFImpl. The rationale is recorded in comet-udf-sdk's crate docs. Contain panics at every extern "C" entry point. A panic escaping one aborts the process, taking down the executor JVM and every task on it. User UDF code is arbitrary and panicking is idiomatic Rust, so a panic is now converted to a query error via the existing get_last_error channel. Previously only execute was guarded; init, new_impl, the discovery entry point and the release callbacks were not. Add a user guide page describing the feature as experimental, including the executor-side library placement requirement and the trust implications of loading native code. * refactor: remove the write-only spark.comet.rustUdfs conf propagateConf wrote the registered UDF set into spark.comet.rustUdfs on every register call, and nothing ever read it. Executors resolve the library from the library_path carried in the RustUdfCall proto, so the conf implied a propagation mechanism that does not exist. Document what actually happens on the register scaladoc instead. CometRustUdfRegistry.snapshot had no other caller and goes with it. * feat: cover all non-nested Spark types, validate the declared return type Add echo_c (identity over any type) and stringify_c ((any) -> Utf8) to the test cdylib, and drive both over every non-nested Spark type with a null row: bool, the four int widths, float, double, decimal, string, binary, date, timestamp and timestamp_ntz. echo_c proves the array survives the round trip with its type and nulls; stringify_c forces the UDF to decode the values rather than hand the array straight back. No ABI change was needed, since the FFI surface is type-agnostic, but nothing had demonstrated that beyond Int64. Also validate at plan time that the return type declared through CometRustUDF.register matches what the kernel's return_field reports, naming both types and pointing at register. Registering decimal(10,2) for a column Spark had widened to decimal(11,2) previously surfaced as a bare DataFusion type assertion partway through execution. Complex types (array, struct, map) remain future work and are documented as unsupported. * feat: support complex types, compare return types ignoring nested nullability Arrays, maps and structs work over the ABI, including nesting (array of struct, struct of array), so extend the type matrix to cover them alongside the non-nested types. No ABI change was needed; the FFI surface carries child arrays already. The new plan-time return type check rejected them, and was right to look: Spark carries containsNull and field nullability inside the declared type, so array<int> with containsNull=false converts to a List with a non-nullable child, while the array actually delivered to the UDF has that child normalized to nullable. Compare with nested nullability erased, and promise DataFusion the kernel's own type rather than the declared one so its exact-match assertion compares like with like. The comparison stays strict about everything that changes how bytes are read: decimal precision and scale, timestamp unit, struct field names and order. Also document the return type model, which came up in review. A kernel has no fixed return type: return_field derives one from the argument types on every call, and echo_c serves all 19 types in the matrix. What is fixed is the type declared to register, because Spark needs a concrete DataType to plan against, and that is per-registration rather than per-kernel. * ci: run CometRustUdfSuite in the PR builds The suite cancelled itself whenever -Dcomet.test.udfs.lib was unset, and nothing set it, so it had never run in CI and the ABI had never been exercised against a Linux .so. Have the suite locate the cdylib under native/target instead of requiring the flag, keeping the property as an override. When the library is missing the suite still skips locally, but fails in CI, where its absence means the native build or the artifact upload changed rather than that someone forgot a flag. An undefined Maven property arrives in the forked JVM as the literal string "null", which is why the override needs filtering rather than a plain null check. Upload libcomet_test_udfs alongside libcomet from both native build jobs: the same cargo invocation already produces it, since the crate is a workspace default member, and the test jobs download the artifact rather than building. Register the suite in the expressions bucket of both workflows, as dev/ci/check-suites.py requires. * style: drop redundant string interpolators in the Rust UDF messages scalafix RedundantSyntax flags s-prefixed literals with nothing to interpolate. These were latent: the previous head commit on this branch carried [skip ci], so the lint jobs had never run against these files. * fix: drop a loaded UDF library after its kernels, not before LoadedLibrary declared `library: Arc<Library>` ahead of `udfs`. Struct fields drop in declaration order, so dropping a LoadedLibrary ran dlclose first and then invoked each kernel's `release` function pointer, which lives in the text of the library just unloaded. Nothing else holds a clone of that Arc. Unreachable in production, since the process-wide cache never drops a LoadedLibrary, but the loader tests build one via `load()` and drop it. Reorder the fields and say why in a comment, so a later edit does not quietly reintroduce it. * fix(ci): build the test UDF cdylib before the Rust test run The rust-test job runs `cargo nextest run` and nothing else. comet-test-udfs is `crate-type = ["cdylib"]` with no test targets, so a test build compiles the crate without ever emitting libcomet_test_udfs, and the three rust_udf tests that dlopen it failed on a clean checkout: execution::rust_udf::cache::tests::same_path_returns_same_arc execution::rust_udf::loader::tests::all_exported_kernels_are_discovered execution::rust_udf::loader::tests::load_test_udfs_succeeds It passed locally only because a previous `cargo build` had left the artifact in target/debug. Build the crate explicitly before the test run, and replace the bare `expect("load")` with a message naming that command, so the next person to hit this reads the fix instead of a dlopen error. * fix: reject deterministic = false when registering a Rust UDF `CometRustUDF.register` accepted a `deterministic` flag, installed a nondeterministic Spark stub for it, and carried it in the RustUdfCall proto, but nothing on the native side ever read it: ImportedCScalarUdf hardcodes Volatility::Immutable. A UDF the caller declared nondeterministic was therefore planned as pure and could be constant-folded, evaluated once and reused, or eliminated as a common subexpression. Honoring the flag needs the cached per-library signature to become per registration, which is #5249. Until then, fail the registration rather than let the parameter lie about what Comet does with it. Raised in review: does the proto field's value mean RustUdfs are always immutable? * docs: clarify the Rust UDF ABI's stability, purity, and path resolution Review follow-ups on the SDK and the user guide, plus two assertion tweaks: - State that the C ABI structs are specific to one Comet version and are internal under the versioning policy, while the host's own arrow and datafusion versions need not match the cdylib's. - Say on CometCScalarUdf that only immutable functions are supported, and add the same to the user guide's limitations. The guide previously told readers to pass `deterministic = false` for impure functions, which is now refused. - Describe what libraryPath actually does: an absolute path is the sensible choice, but a bare name resolves through the platform loader search path, and neither is a security boundary. - Note that CometRustUDF is deliberately not @Public, so it carries no compatibility guarantee. - Report rather than swallow a panic caught in a release callback: there is no error channel there, so it goes to stderr with a context string. - Check `release.is_some()` alongside private_data in the FFI debug asserts, which catches a call made after release rather than only an uninitialized one. * review: address mbutrovich's first pass on the Rust UDF path Drop the dead code and vacuous round-trip the review found, correct the comments that were wrong, and pin the error classification with tests. - `validateLibrary` returns void. It only ever reported back the name it was given, since the native side already filters `lib.udfs` by that name and errors when it is absent, so the JSON construct/parse cycle and the `require(described.name == name)` that followed it could not fail. `listUdfs` had no caller and comes out with it. - `classifyNativeError` matched a bare "ABI", which every loader message can contain because they all interpolate the library path: a library under a directory named `ABI` had its "failed to open" reported as an ABI mismatch. Match the message wording instead, and add a test per failure mode so a rewording on the native side fails a test rather than silently changing the exception a caller sees. - Publish the registry entry before installing the catalog stub. The stub is what makes the name resolvable to the analyzer, so the old order left a window where a query could plan the name with no registry entry and hit the stub's "not evaluated" exception. - Correct the `Mutex` comment in `ImportedCScalarUdf`. It claimed DataFusion serializes invocations of a `ScalarUDFImpl` per batch; it does not, and the planner hands every task the same instance out of the process-wide cache, so concurrent calls are the normal case. The lock is load-bearing rather than defensive, and it serializes all batches of a UDF in a process. Relevant to #5252, which would remove it. - Document `scalar_args` as reserved and always NULL in ABI v1. Nothing on the host can populate it: literals are expanded to full-length arrays before a kernel sees them. - Justify the `CometRustUdfRegistry` singleton's lifetime and bounds where the contributor guide asks for it, and record that session scoping is unresolved. - Add the tests the review asked for: a concurrent-registration test, and a name-collision test which is `ignore`d because it fails today. A Rust UDF does answer calls to an ordinary Scala UDF registered under the same name. The obvious fix does not work: `functions.udf` wraps the closure it is given, so the object in `ScalaUDF.function` is Spark's and not one Comet can recognize, and the `udf(AnyRef, DataType)` overload that would have preserved identity is gone in Spark 4. * docs: point the Rust UDF follow-up comments at their tracking issues Replaces the review-thread links added in the previous commit, and adds pointers where the code has a known gap but carried no reference: - #5294 registry session scoping - #5295 a Rust UDF answers an ordinary Scala UDF of the same name - #5296 the adapter rebuilds the impl and re-resolves the return type per batch - #5297 the library cache holds its write lock across dlopen * refactor: name the UDF mechanism for what it is, not for Rust The ABI this feature loads a UDF through is a set of `#[repr(C)]` structs of function pointers carrying Arrow C Data Interface arrays. No DataFusion type and no Rust type appears in it, the entry points are unmangled C symbols, and every allocation is freed through a callback the library supplies, so nothing about the host requires the library to have been written in Rust. The Rust SDK is one client of that ABI, not the ABI itself. The names said otherwise, and one of them is a wire format that would have been awkward to change later: - proto `RustUdfCall` -> `NativeScalarUdf`, field `rust_udf_call` -> `native_scalar_udf`, which also parallels the existing `JvmScalarUdf` for the JVM codegen dispatch path. Field number 75 is unchanged. - `CometRustUDF` -> `CometNativeUDF`, and the registry, bridge, metadata, exceptions and suite alongside it. - Rust internals move to `execution::c_udf`, matching the C-ABI names already there (`CometCScalarKernel`, `comet_c_udf_list_v1`), and the bridge file to `comet_native_udf_bridge.rs`, matching the JVM class it serves. - User-facing messages and comments say "native UDF" where they mean the mechanism, and keep "Rust" where they mean the SDK. Unchanged on purpose: `comet-udf-sdk` and `comet-test-udfs` were already neutral, and `rust_udfs.md` documents the Rust SDK specifically, so a `cpp_udfs.md` would sit beside it rather than replace it. Also narrows the claim that went with the old naming. The user guide said the ABI "is implementable from C or C++", which reads as an offer; it is true of the design but nothing supports it, since no header is published, the layouts are only defined in Rust source and are not stable across releases, and the SDK's panic guards have no automatic equivalent in another language. The docs now say that plainly. Drive-by: the `NativeScalarUdf` proto comment still described dispatch through "whichever ABI flavor (C ABI / datafusion-ffi)" the library registered under. The datafusion-ffi flavor was removed earlier in this branch. * perf: drop the per-process mutex serializing native UDF batches The adapter held one `Mutex` across the whole `invoke_with_args` body, including the user's compute. Every Spark task in an executor shares one `ImportedCScalarUdf` through the process-wide library cache, so a UDF over a wide scan ran at roughly one core no matter how many tasks were in flight. The lock is not needed. Both kernel-level callbacks the adapter reaches, `function_name` and `new_impl`, take `*const CometCScalarKernel`, and the SDK's exporter only reads through that pointer: everything mutable is the per-call `CometCScalarKernelImpl`, which never leaves the stack frame that built it. State the concurrency requirement plainly in the ABI docs, on the `CometCScalarUdf` trait and in the user guide, and pin it with a test driving eight threads through one shared adapter. * fix: stop rejecting a native UDF over positional list and map field names Two type-mapping traps, both reachable by following the user guide. The guide listed `TimestampType` and `TimestampNTZType` as the same Arrow type. They are not: Comet maps the former to `Timestamp(Microsecond, Some("UTC"))` and the latter to `Timestamp(Microsecond, None)`, and the timezone is compared exactly, so a UDF that built a `TimestampMicrosecondArray` without calling `with_timezone` was rejected at planning against a `TimestampType` registration with no hint as to which axis disagreed. The table now shows the tag, the error message names the axes that must match exactly, and `make_ts_utc_c` / `make_ts_naive_c` pin both directions. The map case was a real rejection over a spelling. Comet emits a map's children as `entries` / `key` / `value` while arrow-rs's `MapBuilder::new(None, ..)` names them `entries` / `keys` / `values`, and the comparison was strict about names at every level. Those names are positional in the Arrow columnar format and Comet's own `CometMapVector` reads them by index, so they are now normalized before comparing, as is a list's element name. Struct field names stay strict: they are part of the Spark type and are how a caller addresses the result. `echo_c` could not have caught either one, since it hands back the array Comet gave it and so inherits Comet's own tag and names. `make_map_c` builds its output through arrow-rs's default builder instead. * test: cover native UDFs in filter, join, grouping key and window Every end-to-end test called the UDF from a projection. Filters, join conditions, grouping keys and window partitioning each route through a different Comet operator, so a regression in the rule wiring for any of them would have shown up only as a silent fallback to Spark. The catalog stub is what makes these assertions load-bearing: it throws if Spark ever evaluates the UDF itself, so each test fails rather than passing on the JVM path. * refactor: expose register and get on the CometNativeUdfRegistry object Callers reached through `CometNativeUdfRegistry.instance` to get at the process-wide registry, which puts the singleton in every call site for no benefit. Forward both methods from the companion so `instance` stays an implementation detail. * test: load the concurrent adapter test's library through the cache concurrent_invocations_share_one_adapter took its adapter from a LoadedLibrary returned by `load` and let that struct drop before the threads ran. The adapter holds no reference to its library, so the drop dlclosed it. On Linux that unmaps the library, and the first `new_impl` call jumps into unmapped text. The Linux rust-test job has failed on both pushes since the mutex was removed, and the 2026-09-07 log shows this test dying with SIGSEGV. macOS keeps an image with thread-locals mapped after dlclose, so it passed locally. Go through `cache::get_or_load` instead, which is how the planner gets an adapter and which never unloads. Also correct the `library` field doc, which said the library is never unloaded when only the cache makes that true. * fix: refuse a UDF result whose type disagrees with return_field An FFI_ArrowArray carries no type, so the host imports a result as whatever return_field declared. A UDF declaring Int64 that returned an Int32Array was read as the wrong values, or past the end of the buffer. Check the result in the SDK, the last point that knows both types, and correct the crate docs that still promised cross-release compatibility and C/C++ support. * fix: keep native UDF libraries loaded and plan canonical child names The adapter now holds an Arc<Library> declared after its kernel, so it cannot outlive the library its callbacks live in. The planner promises the kernel's return type with list and map children renamed to Comet's names, and the adapter relabels each result to match, so a map built with MapBuilder defaults combines with Comet's own maps in if/CASE. Adds sub_c and tests for scalar arguments. * test: cover multi-argument, literal and if-combined native UDF calls Also removes CometNativeUdfSignatureException, which was never thrown. * docs: correct the Rust UDF guide's dependency, panic and reload advice * feat: refuse native UDF calls whose argument types differ from the registered ones The catalog stub is untyped, so Spark inserts no casts and inputTypes was only used for its length. The serde now compares the call's argument types with the registered ones, ignoring nullability, and fails at planning time naming both. Closes the open question from the self-review on what inputTypes means. * fix: release the native UDF cache's read guard before taking the write lock `get_or_load` looked up the canonical path in an `if let` scrutinee and took the write lock inside its body to record the new path. The scrutinee's read guard lives until the end of that block, so a first lookup through a path that resolves to an already-loaded library, such as a symlink, waited on its own read lock forever. Two threads racing on a library's first load through a non-canonical path could reach the same branch. The lookup is now a statement of its own. The new test looks a loaded library up through a symlink on a separate thread, so a regression fails it after 30 seconds instead of hanging the run. * fix: widen a native UDF's nested nullability before it enters the plan The planner promised DataFusion the kernel's own return type, with only list and map child names canonicalized. A kernel may report a nested field as non-nullable where the registration declares it nullable: the declared-vs- reported check disregards nested nullability, so this is accepted. An operator that takes its type from the UDF would then have to narrow a nullable field from another expression to match, and that cast fails with `Cannot cast nullable struct field 'a' to non-nullable field`. `canonicalize_child_names` becomes `promised_return_type`, which also makes every nested field nullable except a map's entries and keys, which Arrow requires to be non-null. The adapter accepts a result whose promised form matches and conforms it with a cast that only renames fields and widens their nullability, reusing the buffers. `make_struct_c` in the test library reports a non-nullable struct field. The planner and adapter tests pin the widened type. The suite runs the reported `if` query in both branch orders. Main's IF branch reconciliation (#6458) now covers that query too, so the planner change keeps other consumers consistent. * fix: ignore struct field metadata when checking native UDF argument types `checkArgumentTypes` compared `deepNullable` forms, which keep each `StructField`'s metadata. A struct argument whose field carries a comment was refused against a registration with the same field name and type, with an error naming `struct<a:int>` on both sides. `DataTypeSupport.equalsIgnoreNullability` re-derives Spark's `DataType.equalsIgnoreNullability`, which compares struct field names but neither nullability nor metadata. Spark's own is `private[sql]` on 3.4. The new test passes a struct whose field has a comment, built with `struct(col.as(name, metadata))`. * fix: discard native UDF registration results explicitly for the strict-warnings build `spark.udf.register` returns the registered `UserDefinedFunction` and `ConcurrentHashMap.put` returns the previous value. Both were discarded implicitly at the end of a `Unit` method, which `-Ywarn-value-discard` under `-Pstrict-warnings` rejects as `discarded non-Unit value`. The fatal error also cut scalac short, which is where the job's four errors in untouched shuffle files came from. * docs: pin the Rust UDF SDK to the release the docs were built for The example dependency pinned `tag = "1.1.0"`, which predates the SDK, so it could not resolve. It now uses `$COMET_VERSION`, which the docs build replaces with the release a page was built for. The development docs, whose version has no tag, say to pin `rev` instead. * fix: recognize native UDFs by their function instead of by name `CometScalaUDF` recognized a native UDF by looking the call's `udfName` up in `CometNativeUdfRegistry`, a process-wide map keyed by bare name. So an ordinary UDF registered under a native UDF's name was answered out of the native library, returning wrong values with no warning (#5295). A registration in one session also claimed the name in every other session sharing the driver JVM (#5294). `register` now installs a builder in the session's own function registry, as `spark.udf.register` does. The builder creates a `ScalaUDF` that holds a `CometNativeUdfFunction`, which carries the registration, and the serde matches on that function. An ordinary UDF registered under the same name holds its own function and stays on the JVM path. A registration is a temporary function of the session it was made in. The process-wide registry is gone. `spark.udf.register` could not do this, because `functions.udf` wraps the function it is handed, and the overload that did not is gone in Spark 4. The builder checks the argument count itself, which the fixed-arity stub used to do implicitly. `ScalaUDF` casts its function to the `FunctionN` of its arity, so there is one subclass per arity, up to the existing limit of four. The name-collision reproduction is enabled. A new test checks that another session neither sees the registration nor has its own UDF of that name answered natively. * refactor: register native UDFs as session-scoped Catalyst expressions Match the registration in #6697. A native UDF's calls now resolve to a `NativeUdfCall` expression rather than a `ScalaUDF` holding a marker function, and `CometNativeUdfCall`, a serde of its own, emits the `NativeScalarUdf`. `CometScalaUDF` and `DataTypeSupport` are back to main's versions. - Spark's analyzer checks the argument types (`ExpectsInputTypes`), disregarding nullability and field metadata, and inserts no casts. A mismatch fails analysis instead of Comet's planning, so Comet's own type check, `equalsIgnoreNullability` and `CometNativeUdfArgumentTypeException` are gone. - A call with the wrong number of arguments fails analysis with Spark's own error. - Nothing caps the arity any more. Four was the most `ScalaUDF`'s `FunctionN` casts allowed for, so `sum_c` and a test cover a five-argument call. - The codegen dispatcher declines an ordinary UDF whose arguments hold a native call, so Spark gets that operator and the failure names the real cause. - If Spark evaluates the call, it throws `CometUdfNotEvaluatedException`. `ShimSessionFunctionRegistry`, `CometUdfErrors` and `CometUdfExceptions` are byte-identical copies of #6697's, so whichever pull request lands second merges them without a conflict.
The documentation is generated from this commit.
| Commit: | d16f7bb | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
feat: support SortAggregateExec (#4565) * feat: support SortAggregateExec [skip ci] Wire SortAggregateExec through the existing CometBaseAggregate.doConvert path so that queries planned with useObjectHashAggregate=false (or with TypedImperativeAggregate functions whose buffer formats prevent hash aggregation) can run their Partial->Final pair natively instead of falling back to Spark. CollectSet was the motivating function; the same wiring covers any other natively-implemented TypedImperativeAggregate. No proto changes: DataFusion's AggregateExec::try_new auto-detects InputOrderMode::Sorted from the child's output ordering, so a sorted input naturally produces sorted output. Wrapper reuse: createExec produces a CometHashAggregateExec (matching CometObjectHashAggregateExec). CometExec.outputOrdering already defaults to originalPlan.outputOrdering, so SortAggregateExec.outputOrdering flows through unchanged for downstream operators that elided a sort against it. Cleanups landed in passing: - Lifted baseAggregateSupportLevel and adjustOutputForNativeState onto the CometBaseAggregate trait (was duplicated 3x and 2x). - Collapsed isAggregate / findCometPartialAgg's Spark-side branches / canAggregateBeConverted's shuffle guard to BaseAggregateExec. * test: add comprehensive SortAggregate coverage via SQL file tests Add a SQL file test (sort_aggregate.sql) that forces collect_set through SortAggregateExec by disabling ObjectHashAggregate, covering global and grouped aggregation, multiple grouping keys, expression keys, NULL/empty/ single-row groups, mixed aggregates, multiple collect_set, DISTINCT, HAVING, a representative data-type spread, and a clean fallback when the aggregate input is incompatible. Add a global (no grouping keys) plan-shape assertion to CometAggregateSuite alongside the existing grouped one, since the SQL framework cannot assert that a SortAggregateExec was actually planned. * fix: distinguish CometSortAggregateExec and update tests for SortAggregate support Introduce a dedicated CometSortAggregateExec wrapper instead of reusing CometHashAggregateExec for converted SortAggregateExec nodes, so the executed plan reflects whether Spark planned a hash or a sort aggregate. A shared CometBaseAggregateExec base carries the common rendering and serialization, and findCometPartialAgg matches the base type so a partial sort aggregate (e.g. collect_set) is still paired with its final aggregate instead of falling back. Update the Spark test diffs to recognize the new wrapper: ReplaceHashWithSortAgg counts CometSortAggregateExec as a sort aggregate, and the generic aggregate collectors in AdaptiveQueryExecSuite and StreamingAggregationDistributionSuite match CometBaseAggregateExec so they see through both wrappers. Refresh the affected SQL file tests now that SortAggregate runs natively: - min_max.sql: string min/max still falls back, but the reason is now the StringType limitation rather than SortAggregate being unsupported. - first_last.sql: the multi-type first/last IGNORE NULLS queries now run natively; shape test_types so each group has a single non-null value, since first/last are non-deterministic across engines with multiple non-null rows. * fix: widen post-merge aggregate guards to CometBaseAggregateExec Self-review after merging main. Three safety predicates landed on main after this branch was cut and matched only CometHashAggregateExec, so a CometSortAggregateExec Partial slipped past all of them: - CometExecRule.preserveSparkAggregateBuffers: a Celeborn exchange that falls back to Spark must restore the Comet Partial below it. Missing the sort variant let a native collect_list buffer reach a Spark Final. Added a regression test in CometCelebornShufflePlanningSuite that fails without this change. - RevertNativeForTransitionHeavyStages.hasUnsafeMixedAggregateAtStageBoundary: same class of guard at a stage boundary. - CometMetricNode.withoutAggregateMetrics: range-partition sampling would double-count sort-aggregate spill and memory metrics. Also drop the duplicate adjustOutputForNativeState that the merge left behind (main grew its own copy handling CollectList and PartialMerge; the branch's older copy handled neither), move main's aggregate metrics override onto the shared CometBaseAggregateExec base, rename COMET_EXEC_SHUFFLE_ENABLED to COMET_SHUFFLE_ENABLED, and mark SortAggregateExec supported in operators.md. * refactor: share the aggregate serde and plan node equality CometBaseAggregate is now parameterized on the Spark aggregate type and supplies enabledConfig, getSupportLevel and convert for all three aggregate serdes. Each serde object keeps only its operator-specific support check, a requiresCometShuffle flag and createExec. The shared getSupportLevel notes that the partial and final test knobs gate sort aggregates too. equals and hashCode move into CometBaseAggregateExec and require the same class, so the two wrappers keep only withNewChildInternal. * fix: walk through sorts when pairing sort aggregate buffers Spark puts a SortExec between each SortAggregateExec and the exchange below it. The passes that pair a buffer-consuming aggregate with the Partial that produces its buffers stopped at that sort, so a native collect_list Partial could stay native under a Spark Final or merge stage. Spark then read Comet's array buffer where it expected binary, and a distinct collect_list chain slipped past the issue #4724 guard and failed natively. findPartialAggInPlan, revertUnsafePartialAggregates and hasUnrepairedNativeBuffer now walk through SortExec and CometSortExec. * fix: keep sort aggregate output in grouping-key order natively CometSortAggregateExec reports SortAggregateExec's grouping-key output ordering, and Spark may remove sorts above it on that basis. DataFusion emits groups in input order only when it sees the input sorted on them. When the ordering comes from outside the native plan, as with a cached sorted relation read through an unordered scan, its vectorized grouping can emit groups out of order: NULL and an empty array hash alike, so an ORDER BY over array keys returned [1] before []. A new HashAggregate field, ordered_by_grouping_keys, is set for sort aggregates. The native planner then sorts the aggregate output on the grouping columns unless the aggregate's output ordering already satisfies them. It also normalizes float grouping keys the way native sorts do, so that DataFusion sees the sort below a float-keyed final aggregate and streams it instead of falling back to the hash path. * test: run the collect_set fixtures through SortAggregateExec collect_set.sql and collect_set_floating_fallback.sql gain a ConfigMatrix over useObjectHashAggregateExec, so their whole fixtures run through both ObjectHashAggregateExec and SortAggregateExec. sort_aggregate.sql keeps only the shapes those files do not cover (several grouping keys, an expression key, hashable aggregates riding along) and adds the decimal SUM fallback. Its floating-point fallback query, which expected a fallback that Spark 4.2 no longer has, moves to the matrix. A new fixture checks that SortAggregateExec stays in Spark, with the shuffle reason, when Comet shuffle is disabled. min_by.sql and max_by.sql drop the claim that Comet never converts SortAggregate and check the type fallback without the TypedImperativeAggregate too. The two native collect_set Scala tests become one test over both queries. * docs: note how first and last can differ under sort aggregation Spark plans a sort aggregate for first and last over a column whose buffer cannot use hash aggregation, such as a string, and Comet now runs it natively. Within a group the row order is not defined, and Spark and Comet can order rows with equal grouping keys differently, so a group with more than one candidate value can return a different, equally valid, value. * test: run the mixed-engine collect_list tests with sort aggregates Each mixed-engine collect_list test now also runs with useObjectHashAggregateExec=false, so both halves of the split are covered for SortAggregateExec, where a sort sits between each half and its exchange. The tests count every Comet aggregate wrapper rather than only CometHashAggregateExec, and check which aggregate operator Spark planned. * fix: keep decimal AVG with a maximum-precision sum out of sort aggregates Spark's Average keeps its decimal sum at precision p + 10, capped at 38. When a sort aggregate's buffer has a field that is not mutable, such as a string FIRST, Spark holds that sum in a generic row where it is unbounded, so 0.6 + 0.6 - 0.4 over DECIMAL(38,38) recovers to 0.8. The native AVG records the overflow at 1.2 and returns NULL, or raises under ANSI. These queries stayed in Spark before sort aggregates converted. CometSortAggregateExec now declines a decimal AVG or TRY_AVG whose sum is at maximum precision, next to the existing decimal SUM check. The compatibility guide notes both fallbacks. * test: expect the distinct collect_list sort aggregate chain to run natively #4727 added native support for collect_list and collect_set alongside a distinct aggregate, so the PartialMerge stages of the distinct rewrite no longer force the chain back to Spark (#4724). Under sort aggregates the whole chain now runs natively and matches Spark. Check that, and keep the fallback case by checking that the chain stays in Spark, through the sort below each stage, when its final cannot convert.
| Commit: | a0e8b5c | |
|---|---|---|
| Author: | Unik Dahal | |
| Committer: | GitHub | |
feat: support native MergeSummary on Spark 4.1+ (#6640) * feat: support native MergeSummary on Spark 4.1+ * Preserve MERGE summary eligibility through AQE and strengthen parity coverage * Use monotonic cumulative snapshots in the Spark 4.2 retry test CometMetricNode.set adds positive deltas and keeps the highest reported value, so reporting count + 10 and then count made each successful attempt contribute 11. Report the cumulative count twice instead, keeping the failed-attempt injection.
| Commit: | 207fb6f | |
|---|---|---|
| Author: | Unik Dahal | |
| Committer: | GitHub | |
feat: add native support for MergeRowsExec (row-level MERGE INTO) (#5318) * feat: add native MergeRows execution * Address MergeRows review feedback * Review Comments Addressed * Address Review Feedback
| Commit: | 1dcea02 | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
feat: add a native Comet version of RangeExec (#6553) * feat: add an experimental Comet version of RangeExec Add CometRangeExec, which writes the values of spark.range and SQL range() straight into Arrow batches on the JVM and feeds them to the native plan, so the operators above a range can run natively without converting Spark rows. Partition bounds and values follow Spark's generated code for RangeExec, including its behavior where the arithmetic overflows. Off by default behind spark.comet.exec.range.enabled. When the Spark-to-Arrow conversion is also enabled, CometRangeExec takes precedence for Range. Adds CometRangeExecSuite and CometRangeBenchmark. * docs: list the experimental Comet RangeExec on the operators page * feat: generate Range values in native code Add a native RangeExec to the operators crate, planned from a new RangeScan protobuf message. Each task computes its partition's bounds from the partition index in i128 arithmetic and walks them with the loop from Spark's generated code for RangeExec, so the values match Spark's default path, overflow cases included. On the JVM, CometNativeRangeExec is a CometLeafExec that reports Spark's partitioning, so the native plan runs one task per slice with no JVM input. The internal spark.comet.exec.range.native.enabled (default true) switches between this and the JVM generator. CometRangeExecSuite runs every test against both, and CometRangeBenchmark has an arm for each. * refactor: keep only the native Range generator Drop the JVM generator (RangeArrowReader) and the internal switch between the two, and name the native operator CometRangeExec. The native generator was faster on every benchmark case. * fix: keep Range on Spark where its interpreted path disagrees with codegen With whole-stage codegen disabled, Spark runs the interpreted RangeExec.doExecute, which returns different rows from its generated code when that code's arithmetic overflows. Comet follows the generated code, so such ranges now fall back when codegen is off. Also list CometRange in the operator reference, mark originalPlan transient as the other native leaves do, and test a global aggregate over an empty range. * fix: fail a Range with no slices as Spark does Spark does not validate a range's slice count. A non-empty range with fewer than one slice fails when Spark runs it, while Comet would plan no tasks and return no rows, so such ranges now stay on Spark. The Spark-to-Arrow conversion now steps aside for Range only where CometRangeExec supports the range, so a range it declines can still use that conversion. * refactor: simplify CometRangeExec and generalize leaf precedence - Build the native range on DataFusion's LazyMemoryExec with a LazyBatchGenerator, as DataFusion's own generate_series does, instead of a hand-written ExecutionPlan and stream. This also gives it cooperative yielding. The loop now reads like Spark's doProduce, and each batch is sized to the values left. - In CometExecRule, a leaf with an enabled Comet operator now tries that operator before the Spark-to-Arrow conversion, which becomes the fallback for a leaf the operator declines. This replaces the Range special case in shouldApplySparkToColumnar and records the decline reason. - Share the codegen-disabled check, including the NO_CODEGEN factory mode on Spark 3.5 and later, between CometRangeExec and the hash aggregate. - Restore Spark's RangeExec where no native operator consumes the range, since generating it natively would only add a conversion back to rows. - Drop overrides CometExec already provides, and duplicate tests and benchmark validation runs. * refactor: keep overflowing ranges on Spark and generate the rest directly Spark's generated code for RangeExec and its interpreted path, which runs when whole-stage codegen is off or the generated code is too large, disagree only for a range whose arithmetic overflows. Such ranges now stay on Spark whatever the codegen setting, so Comet no longer has to predict which path Spark takes. The check is exact: an element count that does not fit in a long, or a batch span that does not. Every range Comet runs now has the values start + i * step in both of Spark's paths, so the native generator drops the emulation of Spark's 1000-value batches and computes them directly. The codegen-disabled helper shared with the hash aggregate is no longer needed, so the aggregate is back to main's version. The config and docs no longer call the feature experimental. It stays off by default because it can be slower than Spark under cheap operators.
| Commit: | 48ec1c0 | |
|---|---|---|
| Author: | Bhargava Vadlamani | |
| Committer: | GitHub | |
feat: wire native existence join support (#4587) * implement_native_existence_joins * implement_native_existence_joins * implement_native_existence_joins_fix_plans * implement_native_existence_joins_fix_plans * implement_native_existence_joins_fix_plans * init_commit * make_doc_changes * make_doc_default_changes * address_review_comments_update_benches * make native join stricter and address review comments * make stricter fallbacks and add further tests * make stricter fallbacks and add further tests lint * make stricter fallbacks and add further tests * make stricter fallbacks and add further tests * fix: use Long literals in CometExistenceJoinBenchmark to satisfy strict Scala warnings * address review: distinct fallback reasons, IN/direct-marker/string-key SQL coverage, plain conf, docs - CometConf: switch existenceJoin flag to a plain conf(...) with a doc listing what still falls back (matches sortMergeJoinWithJoinFilter precedent) - operators.scala: give the flag-off ExistenceJoin case its own fallback reason - operators.md: note native existence support on BHJ/SHJ and what falls back - existence_join.sql: add direct-marker projection, IN-with-OR, and string-key cases; fix the stale 'disabled by default' header - CometJoinSuite: add flag-off fallback test and existence-SMJ-under-forceSHJ fallback test - CometExistenceJoinBenchmark: set spark.comet.exec.onHeap.enabled so both arms measure the native join * remove dead ExistenceJoin arm from CometSortMergeJoinExec.producedAttributes Existence sort-merge joins always fall back to Spark (the SMJ serde has no ExistenceJoin case), so a CometSortMergeJoinExec is never built with that join type. The override only ever returned AttributeSet.empty, which is already the inherited default for a join, so drop it entirely. The reachable hash-join and broadcast-hash-join overrides are unchanged. * make stricter fallbacks and add further tests
| Commit: | 40369a0 | |
|---|---|---|
| Author: | Ping Zhang | |
| Committer: | GitHub | |
feat: push local TopK thresholds into Parquet readers (#6263) * feat: connect local TopK thresholds to native Parquet readers * test: cover TopK reader filtering lifecycle and fallback * perf: benchmark and document TopK reader pruning
| Commit: | 12e8557 | |
|---|---|---|
| Author: | dustin | |
| Committer: | GitHub | |
fix: key each scan's planning data by its plan node in a native block (#6297) Scans in one native block have their planning data collected under a key and merged with toMap, so two scans with the same key both read the last one's files. The Iceberg key (metadata location plus the SparkScan hash) and the Parquet key (source plus a hash of schema, data filters and projection) leave out what can tell two scans of one table apart: storage partitioned join grouping, partition filters and DPP. A storage partitioned self-join with partially clustered distribution returned duplicate rows, and a bucketed Parquet self-join without pushed filters read the wrong partitions. Both keys now include the scan operator's plan_id, which is the id of the Spark plan node the scan was converted from. The merge in findAllPlanData keeps one copy of identical data under a key and fails on different data instead of dropping one side. Closes #6278.
| Commit: | c59021b | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
feat: support Spark HyperLogLog sketch functions (hll_sketch_agg, hll_union_agg, hll_sketch_estimate, hll_union) (#4802) * feat: add DataSketches HLL wrapper and Spark compatibility spike Wrap the pure-Rust datasketches crate's HLL_8 sketch/union behind SparkHllSketch/SparkHllUnion, hashing inputs via the crate's hash_value wrappers (raw_bytes, sign_extend) so MurmurHash3-x64-128 input matches DataSketches-Java. Verified cross-engine: Comet-produced sketches are byte-identical to Spark hll_sketch_agg output for HLL-array mode and mutually readable for low-cardinality List/Set mode. * feat: add version-specific aggregate serde registration hook * feat: add HllSketchAgg and HllUnionAgg proto messages Adds the two protobuf messages needed to serialize Spark's hll_sketch_agg and hll_union_agg aggregate functions, wired into the AggExpr oneof as field numbers 19 and 20. * feat: add native hll_sketch_agg accumulator and planner arm Wire up HllSketchAgg as an AggregateUDFImpl backed by SparkHllSketch, accepting Int8/16/32/64, Utf8, and Binary inputs and returning a serialized HLL sketch as Binary. Null groups evaluate to NULL, matching Spark's HllSketchAgg. Wires the new AggExprStruct::HllSketchAgg arm into the native planner's create_agg_expr. Also add a placeholder HllUnionAgg planner arm returning a clear "not yet supported" error, since that oneof variant already exists in the proto but its native accumulator lands in a follow-on task. * feat: add hll_sketch_agg Scala serde for Spark 4.x * feat: add hll_sketch_estimate scalar function * feat: mark HLL expressions incompatible and add opt-in end-to-end test [skip ci] * feat: add native and serde for hll_union_agg [skip ci] * feat: add hll_union scalar function [skip ci] * test: add hll_union end-to-end and error tests, apply formatting [skip ci] * fix: HLL empty-group returns empty sketch and size() accounts for sketch heap [skip ci] * ci: run CI for HLL sketch functions * fix: decode compact HLL sketches and keep empty union partials NULL * fix: honour a NULL allowDifferentLgConfigK and reject undecodable HLL_4 hll_union is a TernaryExpression, so a NULL third argument makes the whole call NULL. Separately, an updatable HLL_4 sketch carrying auxiliary-map exceptions is now rejected: the bundled decoder reads the compact aux layout in both forms, so those sketches would decode to a wrong estimate. * fix: document HLL support and address review notes Remaining review feedback on the HLL sketch functions. `docs/source/user-guide/latest/expressions.md` still listed the `hll_*` family under "Not currently planned", so the page contradicted the feature. Narrow that bullet to the families that stay unplanned, add the four rows to the `agg_funcs` and `misc_funcs` tables, and explain why this is the one sketch family Comet accelerates. The Implementation cells were checked against a `generate-docs` run on spark-4.0 and spark-4.1 rather than filled in by hand. Minor notes from the same review: - Drop the dead `update_i32`/`update_i16`/`update_i8` helpers. The accumulator widens with `as i64`, which sign-extends identically. - Downcast once per batch in `HllSketchAccumulator::update_batch` instead of building a `ScalarValue` per row, which copied every string and binary value onto the heap only to hash it and drop it. `every_input_type_hashes_the_same_as_a_direct_update` pins the new dispatch to the old bytes exactly, per input type. - Invalid sketch bytes are user data, not a Comet invariant, so `from_bytes` reports `Execution` rather than `Internal`. This is also what the "surfaces as a plain Comet execution error" incompatibility note already promised. - Complete `getUnsupportedReasons()`: the lgConfigK range case was missing from `CometHllSketchAgg`, and `CometHllUnion` / `CometHllUnionAgg` documented no unsupported cases at all despite both returning `Unsupported` for a non-foldable flag. - Pin `datasketches` to `=0.3.0`. The compact-decode workaround and the HLL_4 aux guard are both written against that release's array code, so the version should not move without re-checking those two tests. Also record the HLL_4-with-aux-entries rejection as an incompatibility on the three serdes that read sketch bytes: Comet errors there where Spark returns an estimate, which belongs in front of anyone opting in. Add the end-to-end test comphead asked for on the NULL `allowDifferentLgConfigK` fix. Spark marks `HllUnion` null-intolerant, so `NullPropagation` folds a foldable NULL flag away before Comet sees it; the test excludes that rule so the native kernel is what answers, and asserts `hll_union` survives into the plan so it cannot pass on a folded or fallen-back plan. * test: keep the HLL lgConfigK rejection test native on Spark 4.2 Spark 4.2 replaces MergeScalarSubqueries with MergeSubplans, which also merges non-grouping Aggregate nodes. That folded the two hll_sketch_agg branches of the test's UNION ALL into a single CTE projecting a struct with two fields both named 's'. Comet does not support CreateNamedStruct with duplicate field names, so the whole plan - hll_union_agg included - fell back to Spark and the test saw Spark's HLL_UNION_DIFFERENT_LG_K instead of Comet's native error. Materialize the lgConfigK=10 and lgConfigK=12 sketches into Parquet first and run the union over a plain scan, a shape that stays fully native on 3.4 through 4.2, and assert that with checkCometOperators. * fix: charge HLL accumulators for the sketch they hold, not the dense maximum `size()` on both HLL accumulators charged the full 2^lgConfigK register array as soon as a group existed. A sketch stays in an 8-slot coupon list or a small coupon hash set until it has seen enough distinct values, so at lgConfigK 21 each singleton group was charged 2 MiB for 32 bytes. 64 such groups exhausted a 16 MiB pool in the Final aggregation even with spilling, where Spark runs the query fine. The wrapper now reports the representation the sketch actually holds. The datasketches crate keeps the mode private, but its estimate is never below the coupon count in LIST and SET mode and seeds the HIP accumulator on promotion, and a union with a dense input is tracked from that input's preamble. A test walks every SET resize and the promotion and compares against the layout the crate serializes: never below it, and exact at all but a handful of points. --------- Co-authored-by: test <a@b.c>
| Commit: | ee3f239 | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
fix: reject a file without field ids at any depth whether or not id matching is on (#6116) (#6266) * fix: reject a file without field ids at any depth whether or not id matching is on Spark's `ParquetReadSupport.getRequestedSchema` raises when the requested schema carries a Parquet field id at any depth and the file carries none at any depth, unless `spark.sql.parquet.fieldId.read.ignoreMissing` is set. The check runs on every read, for both readers, and does not consult `spark.sql.parquet.fieldId.read.enabled`. The native scan ran that check only with the read flag on, only over root fields, and only inside the name remap, which a case-sensitive session with the flag off never reaches. A read schema with ids over a file without ids returned rows where Spark raises, and a file whose ids sit only on nested fields was rejected where Spark reads it and null-fills the unmatched fields. The check now runs at the top of the expression adapter factory with a recursive predicate on both schemas. Id matching itself stays gated on the read flag and root ids, as Spark gates it per struct level. The Iceberg scan opts out, since its reader resolves columns by id and supplies the ids itself. * fix: answer both halves of the missing-field-id check from what Spark reads Spark checks the pruned required schema for ids and the raw Parquet schema for their absence. The requested half now comes from `required_schema` at plan time, since the logical schema DataFusion hands the adapter is the full read schema, so an id on a column the query never projects no longer raises. The file half moves into the eager page index reader, which walks the raw schema from the footer. The Arrow schema the adapter sees has lost container metadata after INT96 coercion and never carries ids on repeated list or key value groups, so a file whose only id sat on a struct holding a timestamp raised where Spark reads. The JNI error conversion unwraps the Parquet external error so the Java side keeps the same exception, and the Iceberg scan no longer needs an opt-out since it does not use that reader. Tests cover the pruned projection, the timestamp struct, an id only on a list or key value group, a directory mixing files with and without ids, nested schema pruning, id zero, count(*), and the exception class on each version. * test: find the missing field ids error anywhere in the cause chain On Spark 4.0 and 4.1 the RuntimeException Spark raises for a read schema with field ids over a file without any arrives wrapped once, in the FAILED_READ_FILE SparkException that collect raises, and the helper assumed one more layer above it. It now walks the cause chain and requires a plain RuntimeException carrying Spark's message, for Spark and for Comet alike, so the layer count no longer matters. * fix: let the serde say whether the file must carry field ids The requested side of the missing field ids check is Spark's own hasFieldIds over the required schema, so the serde now sends one flag, require_field_ids, in place of ignore_missing_field_id, and the native Arrow walk is gone. The reader factory takes that one flag through a builder and the parquet options carry nothing for it. The error names the file, and the module doc leads with the lasting reason the check reads the raw footer: ids on repeated groups and on the message root never reach Arrow. The Rust scan tests that had Scala twins are gone, the repeated group cases moved to Scala where they compare against Spark, and the Scala tests fold into the existing port and one pruned schema test with Spark's own assertion form. (cherry picked from commit 863f11de5eb78bacf11737032cdbcf496ad520b5) Co-authored-by: dustin <dwsmith1983@users.noreply.github.com>
| Commit: | 863f11d | |
|---|---|---|
| Author: | dustin | |
| Committer: | GitHub | |
fix: reject a file without field ids at any depth whether or not id matching is on (#6116) * fix: reject a file without field ids at any depth whether or not id matching is on Spark's `ParquetReadSupport.getRequestedSchema` raises when the requested schema carries a Parquet field id at any depth and the file carries none at any depth, unless `spark.sql.parquet.fieldId.read.ignoreMissing` is set. The check runs on every read, for both readers, and does not consult `spark.sql.parquet.fieldId.read.enabled`. The native scan ran that check only with the read flag on, only over root fields, and only inside the name remap, which a case-sensitive session with the flag off never reaches. A read schema with ids over a file without ids returned rows where Spark raises, and a file whose ids sit only on nested fields was rejected where Spark reads it and null-fills the unmatched fields. The check now runs at the top of the expression adapter factory with a recursive predicate on both schemas. Id matching itself stays gated on the read flag and root ids, as Spark gates it per struct level. The Iceberg scan opts out, since its reader resolves columns by id and supplies the ids itself. * fix: answer both halves of the missing-field-id check from what Spark reads Spark checks the pruned required schema for ids and the raw Parquet schema for their absence. The requested half now comes from `required_schema` at plan time, since the logical schema DataFusion hands the adapter is the full read schema, so an id on a column the query never projects no longer raises. The file half moves into the eager page index reader, which walks the raw schema from the footer. The Arrow schema the adapter sees has lost container metadata after INT96 coercion and never carries ids on repeated list or key value groups, so a file whose only id sat on a struct holding a timestamp raised where Spark reads. The JNI error conversion unwraps the Parquet external error so the Java side keeps the same exception, and the Iceberg scan no longer needs an opt-out since it does not use that reader. Tests cover the pruned projection, the timestamp struct, an id only on a list or key value group, a directory mixing files with and without ids, nested schema pruning, id zero, count(*), and the exception class on each version. * test: find the missing field ids error anywhere in the cause chain On Spark 4.0 and 4.1 the RuntimeException Spark raises for a read schema with field ids over a file without any arrives wrapped once, in the FAILED_READ_FILE SparkException that collect raises, and the helper assumed one more layer above it. It now walks the cause chain and requires a plain RuntimeException carrying Spark's message, for Spark and for Comet alike, so the layer count no longer matters. * fix: let the serde say whether the file must carry field ids The requested side of the missing field ids check is Spark's own hasFieldIds over the required schema, so the serde now sends one flag, require_field_ids, in place of ignore_missing_field_id, and the native Arrow walk is gone. The reader factory takes that one flag through a builder and the parquet options carry nothing for it. The error names the file, and the module doc leads with the lasting reason the check reads the raw footer: ids on repeated groups and on the message root never reach Arrow. The Rust scan tests that had Scala twins are gone, the repeated group cases moved to Scala where they compare against Spark, and the Scala tests fold into the existing port and one pruned schema test with Spark's own assertion form.
| Commit: | 590b669 | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
feat: positional round robin shuffle keyed on a row ordinal (#6095) * feat: positional round robin shuffle keyed on a row ordinal Comet implements Spark's round-robin shuffle as hash partitioning over every column of every row. On a wide nested schema that dominates the shuffle write: `create_murmur3_hashes` recurses into every struct child per row, and the resulting row-level scatter forces `interleave_record_batch` to walk every column and child again on flush. It is also not round robin. Placement is a pure function of a row's contents, so a column of one repeated value lands entirely on one reducer where Spark's round robin spreads it. Adds `RoundRobinStrategy::RowGroups`, which places rows the way Spark does: the row at task-global ordinal i goes to `(mapPartitionId + i / groupRows) % numPartitions`. The counter is over rows, not batches, and carries across batch boundaries, so placement does not depend on how the reader frames its input. That matters because no Spark contract covers framing: `DETERMINATE` promises the same rows in the same order and says nothing about chunking, so an operator that spills can reframe under different memory pressure while still honouring it. Keying on a row ordinal reduces the residual assumption to exactly the one Spark's own round robin makes. Positional placement is used only where that assumption is established, in two independent places that both have to hold: * `CometShuffleExchangeExec.replaysRowsInOrder` walks the native subtree fused into the writer, which the RDD graph cannot see because the subtree collapses into one `CometNativeShuffleInputRDD`. Short allowlist rather than a denylist: a native scan under nothing but projections and filters. * `CometNativeShuffleInputRDD.getOutputDeterministicLevel` applies Spark's own `isOrderSensitive` rule to everything below that RDD, reporting INDETERMINATE over a non-determinate parent so the DAGScheduler rolls the stage back instead of re-running one task into consumed output. Anything else keeps `HashAll`, which is safe to re-execute whatever its input does. Both are behind `spark.comet.shuffle.native.partitioning.roundrobin.positional.enabled`, default false. The index representation gains a second shape: positional placement records `(batch, start, len)` runs instead of one `(batch, row)` pair per row, which is smaller against the spill reservation and lets the flush copy whole ranges. `RunIterator` builds a chunk by slicing and concatenating runs, and passes a run covering an entire buffered batch straight through without copying. A schema containing `Utf8View` or `BinaryView` falls back to `HashAll`, because those are the one family whose buffers the IPC writer does not truncate for a slice. * fix: honour a positional group size larger than the batch size Capping `groupRows` at `batch_size` silently ignored what the user asked for. A longer group is meaningful and works as written: it sends several consecutive input batches to the same output partition, which is a legitimate way to trade balance for fewer, larger shuffle blocks. * refactor: simplify positional round robin plumbing Resolve the positional group size into the partitioning itself instead of a parallel field, give each round-robin strategy its own match arm, and buffer runs without moving the scratch vector in and out. Skip the unused partition_starts scratch under positional placement, pre-size the run iterator's chunk scratch, and compute the driver-side positional decision once per exchange so the RDD and the writer read the same value. * fix: rebuild the shuffle writer per iteration in end-to-end benches A ShuffleWriterExec publishes its partition offsets through a OnceLock, so re-executing one exec across criterion iterations fails on the second run with "partition offsets were already published". Build a fresh exec per iteration with iter_batched, keeping construction out of the timing. * bench: add a partitioning-only benchmark for round robin placement The end-to-end shuffle benches are dominated by IPC encoding and the file write, so they cannot separate two placement strategies. Add a group that stops before both: one arm times placement and per-partition index buffering alone, the other adds the flush through a writer that discards its batches, which is where a run-indexed partitioner diverges from a row-indexed one. benches/ compiles as its own crate, so reaching the partitioners needs a pub seam. Rather than export MultiPartitionShuffleRepartitioner and the PartitionWriter trait, add one opaque handle in a doc(hidden) module and leave the rest crate-private. The nested fixture grows a per-row fill. The existing one repeats a single value down every leaf, so every row hashes alike and a hash strategy would place the whole input on one output partition, never performing the scatter that the gather exists to undo. The encoding benches keep the constant fill and their current numbers. * fix: scramble the positional round robin start partition Starting each map task at its own partition id made consecutive tasks start on consecutive partitions. A task walks ceil(rows / groupRows) consecutive partitions from its start, so adjacent starts overlap and the partitions past numMapTasks + groupsPerTask get nothing: ten tasks of 5,000 rows into 200 partitions at the auto group of 64 leave 112 reducers empty. That is the correlation SPARK-21782 fixed. Compute the start the way Spark does, XORShiftRandom(partitionId) .nextInt(numPartitions) + 1, in CometNativeShuffleWriter where the TaskContext is in scope, and pass it in the proto. Still a pure function of the map partition, so a retried task reproduces its own placement, and the + 1 matches Spark's pre-increment so groupRows = 1 places rows exactly where Spark's round robin would. The planner no longer reads self.partition, which removes the jni_api partition-0 caveat. * docs: correct the positional round robin balance and alignment claims Three wording fixes from review, none of them behavioural. The groupRows bound is per map task. A reducer sees the sum over all of them, which only evens out when each task emits many more groups than there are output partitions, so say that in the config doc rather than promising a stage-wide bound. MIN_AUTO_GROUP_ROWS is a cap on how finely a batch is cut, not an alignment guarantee. A run only starts on a byte boundary when the batch starts on a group boundary, and row_seq counts across batches, so after a filter every run in a batch is offset. The start has to be decorrelated across mappers, not merely distinct, which is what XORShiftRandom is for. Also update the round robin item in the shuffle review skill, which still said a positional round robin was a bug by construction. * refactor: place positional runs straight into the run index positional_runs now returns an iterator of (partition, rows) that the repartitioner consumes directly into BufferedRuns, which drops the PositionalRun struct, the scratch vector holding a batch's runs, and the second pass over it. PartitionIndices::empty_like takes no argument, since its one caller passed the index's own partition count, and the gather timer is named interleave_time again, after the metric it feeds and that the Spark UI and the docs show. RunIterator's zero-copy branch compared a whole buffered batch against batch_size with >=, but insert_batch slices every batch to at most batch_size. It is == now, and the two paths' comments agree that every chunk but a partition's last is batch_size rows. The view-type fallback built RoundRobinStrategy::default(), silently dropping a configured maxHashColumns. RowGroups carries max_hash_columns now, and the fallback, factored into partitioning_for_schema so it can be tested, hashes with it. Nothing spilled under the run-indexed shape in any test. The spill metrics and heterogeneous-batching tests run it alongside the row-indexed shape, and positional_placement_survives_spilling requires placement under both spill triggers, the buffer limit and a pool that refuses to grow, to equal the unspilled placement exactly. Three unit tests the end-to-end placement tests subsumed are gone, and the start-partition test is folded into the walk test, which now asserts every group's partition, a wrapping start included. The rustdoc and proto comments shrink to what each item does and point at native_shuffle.md for the argument. * fix: keep positional round robin off under Celeborn Positional placement would be the first path to hand the Celeborn push writer sliced batches, and an indeterminate stage's rollback has not been worked through for push shuffle, so positionalRoundRobinSpec now declines under the Celeborn shuffle manager. The planning suite checks it over a bare native scan, so the manager is the only thing ruling it out. A missing TaskContext now throws instead of starting the task at XORShiftRandom(0), which would put every task at the same partition: the correlation the scrambled start exists to prevent. usesPositionalRoundRobin moves from a public companion predicate to a package-private method on the exec, with shuffleType folded into the decision, since nothing consulted it for a columnar exchange. The two tests that exercised only the placement arithmetic are replaced by one that runs ten real map tasks of 5,000 rows into 200 reducers through the shuffle and requires none of them empty. With the start reverted to the bare map partition id it reports exactly the 112 the simulation predicted. The exchange, RDD and config docs now say what Spark's round robin requires. By default it sorts each map partition first, so a retry only has to produce the same rows; positional placement never sorts and needs them in the same order, which makes replaysRowsInOrder the whole safety argument. The RDD-level isOrderSensitive check is described as the defence in depth it is: a native scan leaf contributes no input RDD, so under today's allowlist it cannot fire. * docs: correct the Spark round robin comparison and keep the argument in one place Spark's round robin sorts each map partition before assigning positions unless spark.sql.execution.sortBeforeRepartition is off, so by default a retry only has to produce the same rows. Positional placement never sorts and needs them in the same order. The guide said the two made the same assumption, which undersold how much rests on replaysRowsInOrder, and it said groupRows = 1 matches Spark, which holds only with the sort off. The rationale was written out in about seven places and had already needed a three-way fix. The round robin section of native_shuffle.md now holds it, opening with what Spark does in both modes, and the review skill's item is a pointer plus the checks a reviewer should make. The interleave_time metric description covers the run copy as well. * bench: build each partitioning fixture batch separately Both partitioning fixtures were one batch cloned eight times. A clone shares its buffers, so the reservation charged seven of the eight nothing, and the gather kept rereading one cache-resident batch. Each batch is now built from where the previous one left off. On the same machine, the cloned fixture reproduces the numbers first posted for this branch, while distinct batches roughly triple the nested gather wherever rows are copied: HashAll's place+gather goes from 37.7 ms to 106 ms and RowGroups(auto)'s from 8.5 ms to 20.4 ms. The zero-copy RowGroups(8192) arm, which copies nothing, is unchanged. * fix: gate positional round robin on deterministic expressions and freeze its group size Admit a native project or filter under positional placement only when its expressions are deterministic, since a nondeterministic UDF can reorder or re-filter rows on a retry while the RDD still reports DETERMINATE. Resolve the default group size on the driver and carry it in the shuffle dependency, so a map stage re-run after spark.comet.batchSize changes places rows exactly as the attempt it replaces. The native planner now rejects a non-positive group size instead of deriving one from the executor's batch size. Qualify the "wraps once per batch" claim for partition counts past the 64-row floor. --------- Co-authored-by: Oleks V <comphead@users.noreply.github.com>
| Commit: | ef7bda5 | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
feat: support listagg / string_agg aggregate (Spark 4.0+) (#4816) * feat: support listagg / string_agg aggregate (Spark 4.0+) Adds a `SparkListAgg` UDAF and a `CometListAgg` serde so that Comet can natively execute the simple form of Spark 4.0's `LISTAGG(child, delimiter)` / `string_agg` on `StringType` inputs with a literal delimiter. `WITHIN GROUP (ORDER BY ...)`, `BinaryType` inputs, non-literal delimiters, and non-default collations fall back to Spark. `DISTINCT` falls back because Comet already rejects multi-column distinct aggregates. The native accumulator returns `Utf8` but keeps its intermediate state as `Binary` to match Spark's `TypedImperativeAggregate` buffer schema so the Comet shuffle layer does not insert a `Utf8` → `Binary` cast the merge side cannot read back. Scaffolding produced by the `implement-comet-expression` skill. * refactor: hoist CometListAgg reason strings and use `.foldable` Consolidate the reason strings into `private val`s so `getSupportLevel` and `getUnsupportedReasons` reference the same source of truth, and replace the manual `case _: Literal` delimiter check with the standard `.foldable` gate used elsewhere (e.g. `CometPercentile`). * refactor: address listagg review feedback Add a GroupsAccumulator fast path to SparkListAgg for grouped aggregation, remove the dead `datatype` proto field and unreachable serialization branch, drop the never-emitted `distinctReason`, and add a multi-partition merge test. Keep the `Binary` intermediate state: Spark's `ListAgg` is a `TypedImperativeAggregate` whose `aggBufferAttributes` is `BinaryType`, and the shuffle Exchange between the partial and final aggregate uses that Spark-declared schema even though partial and final run in the same engine. The module doc and `state_fields` comment are rewritten to explain this. * test: cover listagg collation fallback and align docs with the Binary state - Add a UTF8_LCASE collation fallback case to listagg.sql; the aggregate serde's collation guard is reachable even though the scan falls back. - Drop the comments claiming the collation guard is unreachable. - Describe DISTINCT as falling back via the multi-column distinct check in aggExprToProto, not Spark's multi-stage rewrite (serde, native, proto). - Update the audit doc: intermediate state is Binary to match Spark's TypedImperativeAggregate buffer at the Exchange.
| Commit: | 5442c93 | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
feat: hook native Parquet writes into Spark's WriteFilesExec seam on Spark 4.0+ (#5763) * feat: hook native Parquet writes into Spark's WriteFilesExec seam on Spark 4.0+ Native writes replace the whole DataWritingCommandExec, which means InsertIntoHadoopFsRelationCommand.run never runs. Everything that method does has to be re-implemented inside CometNativeWriteExec: a hardcoded SQLHadoopMapReduceCommitProtocol (so spark.sql.sources.commitProtocolClass is ignored), dynamicPartitionOverwrite pinned to false, a hand-ported copy of the SaveMode logic, a bespoke commit-message accumulator, and its own commitJob call. Most of the open native-writer issues are symptoms of that one decision rather than independent defects. On Spark 4.0+, V1WritesUtils.getWriteFilesOpt matches the WriteFilesExecBase trait (introduced in 4.0 precisely for this), so a Comet node that extends it gets driven through FileFormatWriter.executeWrite -> SparkPlan.executeWrite -> doExecuteWrite, and Spark keeps ownership of everything above the per-task write. Spark 3.x has no such trait: getWriteFilesOpt matches the concrete WriteFilesExec case class, a Comet node there would not be found, and Spark would silently take FileFormatWriter's non-planned, row-based branch. So the new seam is additive. CometDataWritingCommand and CometNativeWriteExec are kept unchanged and remain the 3.4/3.5 path; CometExecRule picks the path by version and the two never both fire. The legacy path goes away with Spark 3.x support. Add: - CometWriteFilesExec, overriding doExecuteWrite and mirroring FileFormatWriter.executeTask for the parts Comet must do itself: build the TaskAttemptContext, ask the commit protocol for a path, run the native writer, drive the stats trackers, commit or abort. Plus the CometWriteFiles serde and a two-line ShimCometWriteFilesExec in spark-4.x / spark-3.x. - File paths come from FileCommitProtocol.newTaskTempFile and are used verbatim, so names match Spark's part-<id>-<uuid>-c000.<codec>.parquet and committers that track individual files (S3A magic, streaming manifest) work. - Column names, nullability and field IDs come from WriteJobDescription.dataColumns rather than the query output, so INSERT INTO t SELECT a+1 writes the target column's name. - Byte and row counts come from BasicWriteTaskStatsTracker, which stats files through the FileSystem API and is therefore correct on HDFS. - ParquetWriter proto: work_dir is now optional. When set (3.x) the native writer derives the file name as before; when unset (4.0+) output_path is the exact file to write and is used verbatim. - On 4.0+ the opt-in moves to spark.comet.operator.WriteFilesExec .allowIncompatible, with the old DataWritingCommandExec key kept as a deprecated alternative. isOperatorAllowIncompat now resolves alternatives, which the planner's by-name lookup previously bypassed. AQE re-plans the write command's child and re-inserts a WriteFilesExec above the node Comet already converted; leaving DataWritingCommandExec in place means Comet no longer has to guard against the resulting nested native writes. * refactor: take a single deprecated alternative for operator incompat configs Matches the reviewed form on #5293: ConfigBuilder mutates in place, so the Seq destructuring was rebinding the same object. Only one operator has an alternative and there is no reason to expect more. * fix: escape the output path before parsing it as a URI The output path reaches both write serdes as `Path.toString`, which decodes percent escapes: a directory containing a space or a literal `%` yields a string that is not a valid URI, and `URI.create` throws on it. On Spark 4.0+ that exception escaped `CometWriteFiles.convert` and failed the query. On 3.x `CometDataWritingCommand.convert` caught it and silently handed the write back to Spark, so the native writer was never used for those paths. Round-trip through `Path` instead, which re-escapes. Only the scheme and authority reach `extractObjectStoreOptions`, but parsing has to succeed to get at them. Also assert that the INSERT INTO visibility test's write actually went native, rather than inferring it from the read-back. * fix: decline HDFS writes whose output path needs URI escaping * fix: also decline Unicode HDFS write destinations The raw/decoded URI comparison only catches what java.net.URI had to escape, and java.net.URI leaves non-ASCII path characters alone, so an hdfs://ns/cafe<U+0301>/output destination was admitted. percent_encoding's should_percent_encode is !byte.is_ascii() || set.contains(byte), so the native parser escapes every non-ASCII byte regardless of the encode set and the writer creates caf%C3%A9 outside Spark's staging directory. The guard now also declines any character the native parser rewrites. The ASCII half of that set was determined against the locked url 2.5 crate by parsing hdfs://ns/pre<c>post/output for every printable ASCII c: space, ", #, <, >, ?, backtick, { and } are rewritten and the rest survive, so partition directories and Spark's _temporary attempt paths still qualify. The comment no longer claims the Java comparison detects the divergence on its own; both conditions are kept because the Java one still catches a literal % that the native parser leaves alone. Tests add accented (precomposed and combining), CJK, emoji and nested non-ASCII cases plus the remaining escaped ASCII characters, built from code points since scalastyle forbids non-ASCII source. Disabling the new condition makes the accented case fail, so the Java comparison alone demonstrably does not cover it. * fix: address review feedback on the native WriteFilesExec seam Correctness: - Decline HDFS writes whose *file names* would diverge, not just their directory. `mapreduce.output.basename` is caller-controlled and reaches every committed name through `HadoopMapReduceCommitProtocol.getFilename`; a basename holding `?` or `#` makes the native URL parser truncate, so every task writes the same name and they overwrite each other at commit. Adds an execution-time backstop over the complete `newTaskTempFile` path, which a custom commit protocol owns and planning cannot predict. - Read the compression option case-insensitively, as Spark's `ParquetOptions` does. `option("Compression", "lz4_raw")` used to fall through to the SQLConf default, so the unsupported-codec guard was bypassed and Comet wrote SNAPPY into a file Spark had named `.lz4raw.parquet`. The codec is now also re-derived per task from `CodecConfig.from(taskAttemptContext)`, the same place the file extension comes from, so the name and the contents agree by construction. The shared helpers move to `NativeWriteUtils`, which fixes the identical bug on the Spark 3.x path. - Use `Utils.tryWithSafeFinallyAndFailureCallbacks` / `tryWithSafeFinally` in `executeTask` and `writeNatively`, matching `FileFormatWriter`: a failure while aborting or closing the iterator is attached as a suppressed exception instead of replacing the failure that caused it. `statsTrackers` moves inside the guard so a throwing `newTaskInstance` still reaches `abortTask`. Planning and reporting: - Convert `WriteFilesExec` from its enclosing `DataWritingCommandExec` rather than from a separate tag pre-pass, so the output path comes straight from the command that owns it and neither `withNewChildren` copying tags nor "nothing hands us a bare WriteFilesExec" has to hold. - Only skip the fallback reason on `DataWritingCommandExec` when its child really was converted; a write with no native child now says why it fell back. - Drop the node's duplicate `files_written`/`bytes_written`/`rows_written`. `BasicWriteJobStatsTracker` is authoritative here, and the native `bytes_written` reads 0 on HDFS. Tests and docs: - Rust `url_path_rewritten_characters` pins the `url` crate's path encode set, which the JVM guard mirrors; a crate upgrade can no longer reopen the hole with a green build. - New JVM coverage: mixed-case `compression` (honored, and declined when unsupported), the #3426 nested-name INSERT, an empty non-zero partition writing no file, the third-party `WriteTaskStatsTracker` warning, and the basename/committer-path guards. - The abort test no longer claims to demonstrate task retry or speculation. - `installation.md` says which Spark versions its EXPLAIN output applies to. * fix: drop a redundant string interpolator flagged by scalafix RedundantSyntax * fix: decline percent-bearing HDFS basenames at planning mapreduce.output.basename=part%foo passed the planning guard and then failed checkNativeWriteDestination at execution, aborting the job instead of falling back to Spark's writer. needsNativeUrlEscaping deliberately excludes % because the url crate leaves it alone, so only Java's raw-vs-decoded comparison sees it. The basename now goes through the same hdfsPathDivergence predicate the task guard uses, applied to the path the basename produces. That closes the gap for % and for the other characters java.net.URI escapes but the url crate keeps ([, ^, |), and makes the two guards agree by construction rather than by two character sets staying in sync. Also updates CometEmptyRelationParquetWriterSuite, which main added while this branch was open: a native empty relation is the zero-partition input CometWriteFilesExec swaps a single-partition RDD in for, so on 4.0+ that write is accelerated rather than declined. * fix: drop imports left unused by moving the write assertions into the base
| Commit: | 7472561 | |
|---|---|---|
| Author: | Peter Lee | |
| Committer: | GitHub | |
perf: spill every shuffle partition of a task into one file (#5916) * perf: return native shuffle partition offsets over JNI instead of an index file The native shuffle writer knew every partition offset by the time it finished a map task, but handed them to the JVM through a temporary file: LocalPartitionWriter::finish_all created an index file and wrote num_output_partitions + 1 little-endian i64 offsets into it, and CometNativeShuffleWriter read the whole file back with Files.readAllBytes, converted the offsets to lengths, deleted it, and passed the lengths to IndexShuffleBlockResolver.writeMetadataFileAndCommit, which writes Spark's real index file. The temp file existed only to move an array of longs across the JNI boundary, and cost every map task a create, write, read and unlink on top of the index file Spark commits anyway. The parse also allocated an intermediate array and a ByteBuffer per partition (item 5 of #5198), which goes away with the file. The offsets are now published in memory through a PartitionOffsets slot shared by the writer and its ShuffleWriterDestination, and read back over JNI by Native.getShufflePartitionOffsets. The index path no longer travels in the plan, so LocalPartitionWriter.output_index_file and the legacy ShuffleWriter.output_index_file are removed and their field numbers reserved. The offsets have to be read while the native plan is still alive. CometExecIterator closes itself when its stream reaches the end, and close releases the execution context that owns the writer, so reading after drainAndClose returned freed memory and produced garbage lengths. The iterator instead captures the offsets at end of stream, before close, when built with capturePartitionOffsets, which only the local destination sets: RSS reports its partition lengths through its pusher. Partition lengths are derived from effectivePartitionCount, the output partition count, not the numParts constructor argument, which is the input partition count. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * review: address feedback on naming and layering - The shuffle crate does not know about JNI, so PartitionOffsets and the local writer describe the handover in terms of the caller driving the plan rather than the JVM reading over JNI. - ShuffleWriterExec::try_new says what it writes where: partition data to a local file, offsets in memory. - CometExecIterator had three lookalike names for one thing. The read is now a named method, readPartitionOffsetsBeforeClose, whose name and doc carry the constraint that made the placement surprising: the offsets live in the native execution context, close releases it, and hasNext closes as soon as the plan runs out of output, so the final hasNext is the last point they can be read. The field is partitionOffsets and the constructor flag is documented. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * fix: size shuffle partition lengths from the offsets the writer returned Two CI failures, both from assuming more than the writer guarantees. The proto crate's own tests still referenced ShuffleWriter.output_index_file and LocalPartitionWriter.output_index_file, so datafusion-comet-proto failed to compile its test target. I had only checked the shuffle and core crates locally rather than the whole workspace. The round-trip tests now assert that a new plan carries no index path, and that a plan still carrying the retired tag 4 decodes cleanly because the tag is reserved rather than reused. partitionLengths was sized by effectivePartitionCount, which is not what the writer produces. isSinglePartitioning serializes a range partitioning whose sampled bounds came out empty as SinglePartition, so native writes one partition while the declared output partitioning still reports several, and the require failed with "returned 2 partition offsets for 10 output partitions". The index file was always sized by what the writer produced, so deriving the length count from the returned offsets restores the previous behaviour exactly. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * review: drop the reserved declarations for the retired index path Requested in review. The plan proto is built by the JVM and consumed by native in the same process from the same artifact, so no old plan ever meets a new reader and there is nothing for a reserved tag to protect against. Decoding is unaffected either way: an undeclared tag is skipped as an unknown field, and reserved only stops protoc from later reusing the number. The proto round-trip test covering a plan that still carries the retired tag 4 keeps passing, and its comment no longer credits reserved for that. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * review: trim comments and close the gap left by the retired field number Comments cut back to what they need to say, per review: the JNI entry point, the PartitionOffsets type and its set/get, the destination field, try_new, the partition_offsets accessor, the zero-offset test, the finish_all note, and the two comments in CometNativeShuffleWriter. ShuffleWriter field numbers 5 through 11 shift down to 4 through 10, closing the gap the retired output_index_file left. Both sides of the plan are generated from this file and ship together, so no encoded plan outlives the change. That does mean tag 4 now belongs to codec, and the LegacyShuffleWriter test struct claimed it for a string. Decoding a plan carrying it would be a wire type mismatch rather than a skipped unknown field, so the struct drops that field. The test still covers a legacy plan decoding without a partition writer. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * perf: spill every shuffle partition of a task into one file The native shuffle writer spilled each output partition to its own temporary file, created on that partition's first spill and held open until the task finished. A task with P partitions that spilled held up to P spill files open at once, each with a create and an unlink, and every spill scattered its bytes across P files. A task now spills to a single file. Each write appends the partition's blocks and records the range they occupy, and finish_partition copies a partition's ranges into the output in write order. Correctness does not depend on the order partitions are written in, which PartitionWriter leaves unspecified. A single spill round gives contiguous ascending ranges, so the merge reads sequentially. A write that fails partway can leave uncounted bytes in the shared file, which would shift every later range, so the spill refuses further writes and range reads after a failure. The merge also checks each copy's length, so a spill file shorter than its ranges fails the partition instead of writing it short. Closes #3859. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * perf: merge short spill ranges with one positional read Copying a spill range with io::copy costs an lseek, two statx calls and a copy_file_range. Ranges that fit in the write buffer are now read with one read_exact_at into a scratch buffer; longer ranges still use io::copy. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * perf: buffer spill writes across partitions Each partition's spill write ended with a flush, so a spill round issued one write syscall per partition. The spill file's writer is now buffered across partitions and flushed before the merge reads it. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * review: reserve the retired index file tags and drop stale index file docs Requested in review: ShuffleWriter keeps its fields on tags 5 through 11 with reserved 4, and LocalPartitionWriter gets reserved 2, as operator.proto does for other retired fields. The proto test again decodes a plan carrying tag 4. native_shuffle.md steps 6 and 7, and the LocalPartitionWriter, shuffle_write and shuffle_bench output_dir comments, no longer describe an index file. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5 <noreply@anthropic.com>
| Commit: | 8c45f3e | |
|---|---|---|
| Author: | dustin | |
| Committer: | GitHub | |
perf: cache parsed plan data across a stage tasks (#5615) * perf: cache parsed plan data across a stage tasks Every task of a stage deserializes byte-identical plan bytes, yet each one parsed the full operator tree, re-derived the scan source key by stringifying its schema and filter lists, and re-parsed the scan common message before injecting its partition data. Three bounded per-executor caches now share that work: the parsed base plan keyed on content, the parsed NativeScanCommon, and a source-key memo that hits protobuf reference-identity fast path once the base plan is shared. Injection itself stays per task, since partition data genuinely differs, and the injected tree is never cached. Per-task overhead drops from roughly 300us to 60us on a 100-column scan plan and about 3x on a 1000-column plan. Cache misses compute outside any lock so unrelated stages never serialize behind one parse, and racing threads on a cold key adopt a single instance so reference sharing holds. * perf: carry the scan key in the plan and scope prepared data to the base plan entry The executor derived each native scan's lookup key by stringifying its schema and filter lists, memoized in an LRU probed under the map monitor, and parsed scan commons through a second per-scan LRU that a single plan with 17 scans churned to zero reuse. The driver already derives the key once in CometNativeScanExec, so carry it inside the NativeScan proto and read it back on every injection path, including the native shuffle writer's, which previously missed the fast paths entirely. Prepared commons now live inside the base plan's own cache entry (scoped to the shuffleId on the shuffle-write path), so a plan and its scans form one eviction unit. Entries pin the finalized common bytes because scalar-subquery data filters resolve per execution: equal base plan bytes do not guarantee equal finalized commons, and a stale entry is replaced rather than served. The base plan cache keys on a stored hash computed once per task outside the monitor instead of a raw ByteBuffer that rescanned the bytes on every probe. * fix: release prepared shuffle data when a shuffle is unregistered or the manager stops The shuffle-scoped prepared-commons store is a JVM singleton keyed by shuffleId, and shuffle ids restart at zero for every SparkContext, so a local or embedded caller that stops and recreates its context kept stacking new scan keys under ids the previous context had used. Both CometShuffleManager and CometCelebornShuffleManager now drop a shuffle's entry in unregisterShuffle (reached from Spark's ContextCleaner via BlockManagerStorageEndpoint) and clear the store in stop(), so a recreated context starts empty and long-lived contexts release each shuffle's data with the shuffle. * perf: fingerprint the plan bytes on the driver and key the base plan cache by it Every task of a stage built its PlanKey by hashing the full plan bytes with Arrays.hashCode, a scalar loop that took about a quarter of the cached path on JDK 17 for a wide plan. The driver now computes one XXH64 fingerprint where it serializes the plan and carries it on CometExecRDD, so the executor probe hashes nothing. Equality still compares the bytes, so a fingerprint collision cannot serve another stage's plan. The reference check in PlanKey.equals goes too, since HashMap already performs it before calling equals. * refactor: tighten the plan data injector SPI and release both caches on stop Type the injector's two halves through an abstract Prepared member so an implementation's prepareCommon and inject must agree, and move the single remaining cast to the memo boundary. Prefix memo keys with the injector class: every contrib scan arrives as the same CONTRIB_SCAN envelope, so two injectors agreeing on a key and byte-identical commons would otherwise share one prepared object. Drop the Iceberg injector's inner cache and the injectPlanData overload with no production caller, since both real entry points now memoize. The shuffle store gets its own maxCachedShuffles bound, the shuffle managers clear the base plan cache together with the shuffle store on stop, and the lifecycle suite asserts which stage shapes land in which store: a map-only stage in the base plan cache, a scan fused into a native shuffle under the shuffle id. * perf: ship only the source key hash in the NativeScan proto The transported key was the scan's source followed by a hash, and the source already sits in the common of the same message. The proto now carries the int32 hash alone as an optional field, so presence is explicit and a zero hash is still a transported value, and the executor rebuilds the key from the source next to it. The field is unreleased, so it changes in place. * docs: plan bytes are parsed once per stage, not per task QueryContextInterner described the plan as re-parsed per task by CometExecRDD.compute; parsing now happens once per stage on that path, while the per-task re-serialization it also mentions still holds. * docs: state the JVM-wide scope of the plan data stores and the memo key's classloader assumption * fix: adopt the winning prepared common when two tasks race a cold scan key prepareShared put its freshly prepared common unconditionally, so tasks racing the first access to a scan could each keep their own instance. It now merges into the memo atomically and returns whichever equal-bytes entry landed first, on both the base-plan and shuffle-scoped paths. Barrier tests cover both. --------- Co-authored-by: Andy Grove <agrove@apache.org>
| Commit: | e893b43 | |
|---|---|---|
| Author: | Xuanyi Li | |
| Committer: | GitHub | |
feat: add Lance contrib build gate (#5728)
| Commit: | db79067 | |
|---|---|---|
| Author: | Matt Butrovich | |
| Committer: | GitHub | |
feat(iceberg): Iceberg table format V3: apply deletion vector on reads (#5853) * upmerge main, resolve delta * serialize file type, clean up reflection * dedupe delete file path * add missing scala file changes * add new test for deduping puffin files, update docs * address PR feedback --------- Co-authored-by: Matt Butrovich <mbutrovich@apache.org> Co-authored-by: Andy Grove <agrove@apache.org>
| Commit: | 149a1db | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
feat: support max_by and min_by aggregate expressions (#4817)
| Commit: | 99237e1 | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
feat: support native aggregate function `mode` (#4782) * feat: support native mode aggregate function Add native support for the Spark `mode` aggregate, the most frequent value within a group. Spark breaks ties on the default `mode(col)` form non-deterministically (the chosen value depends on JVM hash-map iteration order), so the function is registered as Incompatible and opt-in via allowIncompatible; Comet resolves ties deterministically by returning the smallest tied value. NULLs are ignored, empty input returns NULL, and float keys are normalized to match Spark. The deterministic-flag and WITHIN GROUP forms fall back to Spark. Closes #3970 * fix: make mode's -0.0 key handling match the Spark version, plus review cleanups Addresses review feedback on #4782. The main item was the `-0.0` normalization in `normalize_key`. The review asked to stop collapsing `-0.0` into `0.0`, because Spark keys the frequency map on `java.lang.Double.equals` (via `OpenHashSet`), which distinguishes the two, and `NormalizeFloatingNumbers` never touches aggregate arguments. That is correct for Spark 3.4 through 4.1, but Spark 4.2.0 changed it: SPARK-57329 treats the split `-0.0`/`0.0` counts as a correctness bug and normalizes the key at update time. The fix landed in branch-4.2 after v4.2.0-rc1, so released 4.2.0 has it. Correct behaviour is therefore version-dependent, and neither always-collapsing nor never-collapsing is right across the profiles this repo builds. Added a `normalize_neg_zero` flag to the `Mode` proto message, set from `isSpark42Plus` in the serde (same pattern as `BloomFilterVersion` and `setIsSpark4Plus`), and gated the fold on it natively. `NaN` canonicalization stays unconditional, since `doubleToLongBits` collapses `NaN` on every supported version. `ScalarValue`'s `PartialEq`/`Hash` for floats are both bit-based, so once `NaN` is canonical the zeros stay distinct without further work. Test fixtures used `CAST(-0.0 AS DOUBLE)`, which does not produce a negative zero: an unsuffixed `-0.0` is a DecimalType literal and Decimal has no signed zero, so the column contained only `+0.0` and the existing signed-zero coverage was vacuous. Switched to `-0.0D`. Verified the new fixture is non-vacuous by forcing the flag to the wrong value and confirming it fails. Also in this commit: - correct the doc comment, which claimed the old behaviour matched `NormalizeFloatingNumbers`, and record which Spark comparison path governs `mode` versus `max_by`/`min_by` so the two are not "fixed" to match each other - drop the redundant `default_value` override - count key heap bytes in `size()` so string/binary/decimal modes do not under-report to the memory pool - assert the non-empty-groups invariant at both grouped emit sites - note that the `ScalarValue` frequency map is intentionally type-generic - add a `timestamp_ntz` compared query - add mode_within_group.sql pinning the Spark 4.x ordered forms as fallbacks, including that `ModeBuilder` rewrites `mode(col, false)` to the plain form so it still runs natively - TODO recording that the ASC WITHIN GROUP form could be Compatible * fix: account for group slot storage in mode's grouped accumulator size() --------- Co-authored-by: test <a@b.c>
| Commit: | 8c7b706 | |
|---|---|---|
| Author: | Ping Zhang | |
| Committer: | GitHub | |
feat: native dynamic filter pushdown for hash join into Parquet scans (#5699) * feat: push native join runtime filters into Parquet scans * test: fix right outer join hint for Spark 3.4 * fix: evaluate join runtime filters on the probe key * test: extract dynamic filter tests into a child module * refactor: move join dynamic filter planning into planner * fix: adapt join dynamic filters to DataFusion 55
| Commit: | 2da3291 | |
|---|---|---|
| Author: | Peter Lee | |
| Committer: | GitHub | |
fix: prevent silent overflow when reading Parquet TIMESTAMP_MILLIS values (#5177) * refactor: use Arrow casts for temporal conversions * fix: preserve Spark temporal cast semantics * test: match Spark micros-to-millis cases * test: link Spark micros-to-millis cases * test: use Spark tag in source link * Use imported arity kernel for Spark timestamp downscaling * fix: enforce Spark temporal conversion semantics * andy's 3rd review * address review: drop temporal.rs refactor, improve adapter error message - Revert cast_date_to_timestamp to main's version to avoid conflicting with #5457, which rewrites the same function as a safety fix - Document the Int32 input guarantee and zero-copy reinterpret path in date_from_unix_date - Make the microsecond-physical/millisecond-target rejection message actionable (report link + scan fallback workaround) Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * chore: remove unused imports in parquet_support Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * test: exercise dictionary-encoded pages in TIMESTAMP_MILLIS overflow regression A single row falls back to PLAIN even with dictionary encoding enabled, so the dictionary leg never covered a dictionary read. Write 16 repeated rows and assert hasDictionaryEncodedPages matches the writer setting. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * fix: preserve timestamp pruning by rewriting predicates to the millisecond domain The predicate over a TIMESTAMP_MILLIS file column was wrapped in CometCastColumnExpr, which DataFusion's pruning analyzer cannot see through, so row groups Spark prunes from millisecond statistics were read and hit the checked millis->micros conversion. Rewrite predicate comparisons into the millisecond domain (exact integer rescaling of the literal, mirroring Spark's ParquetFilters) and unwrap IS NULL checks, so pruning works and predicate evaluation never converts file values. The scan output conversion stays checked. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * fix: cover IN, null-safe equality, and nested predicates in the millis-domain rewrite Extend the millisecond-domain predicate rewrite to InListExpr and IsDistinctFrom/IsNotDistinctFrom, both of which DataFusion's pruning predicate can analyze, so row-group pruning protects flat TIMESTAMP_MILLIS columns for those forms too. Nested-field predicates can be neither pruned (PruningPredicate has no nested-field support) nor evaluated as Parquet row filters (struct columns are classified non-pushable), so no rewrite can keep the checked conversion from failing queries Spark answers via nested statistics pruning. Scans whose data filters reference nested fields fall back to the safe cast (overflow -> NULL, the pre-existing behavior), and the checked conversion is scoped to top-level columns for the same reason. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * test: cover empty NOT IN with row filters * fix: retain dropped parquet filters * fix: retain scalar subquery filters --------- Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
| Commit: | 165e2c6 | |
|---|---|---|
| Author: | Peter Lee | |
| Committer: | GitHub | |
feat: carry VariantType identity through schema serialization (#5631) * feat: carry VariantType identity through schema serialization * ci: register CometVariantTypeSuite * fix: address Variant type identity review feedback
| Commit: | 81d637b | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
feat: implement regr_slope, regr_intercept, regr_r2, regr_sxx, regr_syy, regr_sxy aggregates (#4775) * feat: implement regr_slope, regr_intercept, regr_r2, regr_sxx, regr_syy, regr_sxy aggregates Add native support for the six simple linear regression aggregates that previously fell back to Spark. regr_avgx, regr_avgy and regr_count were already accelerated because Spark rewrites them to Average/Count. The native accumulators are composed from Comet's existing Spark-compatible covariance and variance accumulators so the partial aggregation state matches the buffer layout Spark's planner expects between partial and final stages: RegrReplacement (regr_sxx/regr_syy) -> 3 fields, Covariance (regr_sxy) -> 4, PearsonCorrelation (regr_r2) -> 6, and the slope/intercept composite -> 7. regr_r2 matches Spark's behavior of returning 1.0 when the dependent variable is constant but the independent variable varies (a perfect horizontal fit), which differs from DataFusion's regr_r2. * fix: match Spark's per-version regr_slope/intercept/r2 semantics The regr aggregates diverged from Spark in three ways that surfaced as CI failures across Spark versions: - regr_r2 degenerate cases were inverted. Spark 3.4/3.5/4.0 return null when the dependent variable is constant and 1.0 when the independent variable is constant; Comet had these swapped. - Spark 4.1 swapped that degenerate handling again (constant dependent -> 1.0, constant independent -> null). Route the behaviour through a new r2_constant_dependent_is_perfect_fit proto flag set from isSpark41Plus. - regr_slope/regr_intercept compute VariancePop(x) over both-non-null pairs on Spark 3.5+, but over every x-non-null row on Spark 3.4. Route this through a new filter_var_by_pair_nulls proto flag set from isSpark35Plus. Also evaluate regr_r2 as corr = ck / sqrt(m2_y * m2_x); corr * corr to mirror Spark's exact float rounding, so the golden-file postgres aggregates tests match bit-for-bit. * fix: regr_r2 degenerate-case swap (SPARK-55969) applies from Spark 3.5, not 4.1 SPARK-55969 swapped which regr_r2 degenerate case returns NULL and which returns 1.0. It shipped in Spark 3.5.9, 4.0.3 and 4.1.0, so every Spark version Comet builds against except 3.4 has the new behavior. Comet was gating it on isSpark41Plus, which made regr_r2 disagree with Spark on the 3.5 and 4.0 profiles.
| Commit: | 5b332b9 | |
|---|---|---|
| Author: | Jordan Epstein | |
| Committer: | GitHub | |
feat: implement native Iceberg V2 writer via iceberg-rust (#5361) Co-authored-by: Jordan Epstein <jordan.epstein@imc.com>
| Commit: | a2c6bd4 | |
|---|---|---|
| Author: | Ping Zhang | |
| Committer: | GitHub | |
feat: add partition writer destinations to native shuffle plans (#5476)
| Commit: | 5baa6b0 | |
|---|---|---|
| Author: | Oleks V | |
| Committer: | GitHub | |
feat: support `WindowGroupLimitExec` (#4870)
| Commit: | 6716318 | |
|---|---|---|
| Author: | Scott Schenkein | |
| Committer: | GitHub | |
feat: build gate + inert wiring for contrib Delta scans [Delta contrib split, part 2] (#4952) 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>
| Commit: | e4ec630 | |
|---|---|---|
| Author: | Chao Sun | |
| Committer: | GitHub | |
fix: preserve Catalyst nullability and field IDs in native Parquet writes (#5369)
| Commit: | 3df58d5 | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
perf: intern QueryContext SQL text into a per-plan pool (up to 20x smaller serialized plans for TPC-DS) (#5204)
| Commit: | 60413eb | |
|---|---|---|
| Author: | Parth Chandra | |
| Committer: | GitHub | |
feat: support Iceberg metadata columns _pos, _spec, _file, and _partition (#4752) * feat: [iceberg] support iceberg metadata columns _pos, _spec, _file, _partition
| Commit: | 1ce9df1 | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
feat: native uuid() implementation compatible with Spark (#5034) * feat: native uuid() implementation compatible with Spark * refactor: share MersenneTwister, use uuid crate for formatting Extract SparkMersenneTwister into internal/mersenne.rs so shuffle and uuid both depend on it rather than uuid reaching into the shuffle module. Format UUIDs via the uuid crate encoded into a pre-sized StringBuilder, removing per-row String allocation and the hand-rolled hex formatting. * test: address review feedback on uuid tests Fills the five coverage gaps flagged on #5034: - Multi-partition bit-for-bit: uuid(42) FROM ... DISTRIBUTE BY id exercises partitionIndex != 0, which a single-partition test cannot catch (both engines would silently agree even if partitionIndex were ignored). - Aliased projections: uuid(0) = uuid(0) (seeded) and uuid() = uuid() (unseeded, spark_answer_only) pin the stateful-alias semantics that freshCopyIfContainsStatefulExpression relies on. - Empty batch: a Rust test that evaluate() on a zero-row batch does not advance RNG state, guarding against a stray next_uuid() outside the row loop. - Golden fixture: five seeds x five UUIDs captured from Commons Math3's MersenneTwister (the exact RNG behind Spark's RandomUUIDGenerator) lock all 128 bits and cover negative / MIN / MAX seeds; catches next_long regressions that shuffle's next_int tests would miss. - No-scan path: length(uuid(0)) pins the OneRowRelation shape with a deterministic assertion alongside the existing SELECT uuid(0). * test: make uuid partition and seed coverage non-vacuous The DISTRIBUTE BY query never exercised partitionIndex != 0. `... FROM t DISTRIBUTE BY id` parses as RepartitionByExpression on top of the Project, so uuid ran on the scan's partitions, and five tiny files pack into a single FilePartition. Moved the projection above the exchange: SELECT uuid(42) FROM (SELECT id FROM test_uuid_parts DISTRIBUTE BY id) Add CometUuidExpressionSuite, which builds `Uuid(Some(seed))` directly through the existing version-shimmed `getColumnFromExpression`. This gives 3.4 and 3.5 a real cross-engine assertion for the first place: the SQL `uuid(seed)` form is 4.0+, so uuid_with_seed.sql is skipped there and the Rust golden constants were the only guard, with no way to confirm they came from Commons Math3 rather than from the Rust code they are meant to test. One test covers five seeds including negative/MIN/MAX; the other repartitions and asserts getNumPartitions > 1 before comparing, so the partition-offset coverage is checked rather than assumed. Verified by mutation: dropping `expr.seed.wrapping_add(planner.partition())` in UuidBuilder::build fails both new Scala tests. Restored, all pass on spark-3.5 and spark-4.1. Also: - Drop `spark_answer_only` from `uuid() = uuid()` in uuid.sql. The values are random but the comparison is not -- distinct fresh seeds mean false on every row -- so default query mode keeps the value check and adds the nativeness assertion that the guard depends on. - Record the real reason typeof(uuid(0)) cannot be used: Spark's TypeOf.doGenCode interpolates `child.dataType.catalogString` unquoted into generated Java, emitting `UTF8String.fromString(string)`. That code is byte-identical on 3.5, 4.0 and master, so it is not a 4.1 quirk as previously recorded. It is normally hidden because TypeOf.foldable is true and ConstantFolding removes the node; this suite excludes ConstantFolding, which exposes it. Captured as a harness constraint: no SQL fixture here can use typeof on a non-foldable input. - Note on the Rust golden test that CometUuidExpressionSuite is now its backstop, and that regenerating the constants from the Rust side would make it circular. No partition-count assertion is added to the SQL fixture: that suite compares Comet against Spark, so if the plan collapsed to one partition both engines would collapse identically and any such check would still pass. It is only assertable from Scala, which is where it now lives. * ci: register CometUuidExpressionSuite in PR build workflows The check-suites.py preflight requires every *Suite.scala to be listed in both pr_build_linux.yml and pr_build_macos.yml. Add the new uuid suite to the expressions bucket.
| Commit: | d2a61bd | |
|---|---|---|
| Author: | Matt Butrovich | |
| Committer: | GitHub | |
fix: disambiguate Iceberg scans that share a metadata_location (#5180)
| Commit: | c21fe12 | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
fix: throw ARITHMETIC_OVERFLOW for Long.MinValue div -1 under ANSI mode (#5084) (#5146) * fix: throw ARITHMETIC_OVERFLOW for Long.MinValue div -1 under ANSI mode * refactor: reuse Spark's checkDivideOverflow and dedupe overflow check * test: use ARITHMETIC_OVERFLOW condition name in integral divide expect_error Matches the convention used by the other expect_error queries in arithmetic_ansi.sql. Verified against the spark-3.4, spark-3.5, spark-4.0, spark-4.1, and spark-4.2 profiles. (cherry picked from commit 38c3f8bce35642479117ba49f44a02fee31e6158)
| Commit: | e6577ed | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
feat: support SampleExec natively for sampling without replacement (#5110) (#5144) * feat: support SampleExec natively for sampling without replacement Adds a native Comet operator for Spark's SampleExec when sampling is performed without replacement, covering DataFrame.sample, SQL TABLESAMPLE, and DataFrame.randomSplit. The native operator ports Spark's BernoulliCellSampler on top of the existing XorShiftRandom, drawing one value per row and seeding per partition with seed + partitionIndex, so it selects the same rows as Spark for a given seed. Sampling with replacement reports Unsupported and falls back to Spark. * refactor: address review feedback on native sample operator - Build the selection mask with BooleanBuffer::collect_bool, which packs bits without allocating an all-ones null buffer - Precompute the empty-range check in BernoulliCellSampler instead of repeating it per row, and mark sample() inline - Drop the redundant schema() override and the single-use sample_batch helper - Use checkSparkAnswerAndFallbackReason in the fallback tests so they assert the reason, not just the absence of the Comet operator - Trim documentation that restated the same paragraph in four places - Point the contributor guide at the directory operator serdes actually live in * fix: shim the SampleExec seed for Spark 4.2 Spark 4.2 changed SampleExec.seed from Long to Option[Long] and resolves an absent seed into a resolvedSeed field, which is what the operator samples with. Read the seed through CometSampleShim, following the CometCollectShim precedent for the same 4.2 divergence. * test: expand sample operator test coverage and add benchmark Add tests for ordering preservation above a sort, sampling above an aggregate, SQL TABLESAMPLE, and empty batch interleaving in the native operator. Add a Sample case to CometExecBenchmark comparing against Spark. (cherry picked from commit 4b091acc9acc442c301754876348480f49b7321d)
| Commit: | 4b091ac | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
feat: support SampleExec natively for sampling without replacement (#5110) * feat: support SampleExec natively for sampling without replacement Adds a native Comet operator for Spark's SampleExec when sampling is performed without replacement, covering DataFrame.sample, SQL TABLESAMPLE, and DataFrame.randomSplit. The native operator ports Spark's BernoulliCellSampler on top of the existing XorShiftRandom, drawing one value per row and seeding per partition with seed + partitionIndex, so it selects the same rows as Spark for a given seed. Sampling with replacement reports Unsupported and falls back to Spark. * refactor: address review feedback on native sample operator - Build the selection mask with BooleanBuffer::collect_bool, which packs bits without allocating an all-ones null buffer - Precompute the empty-range check in BernoulliCellSampler instead of repeating it per row, and mark sample() inline - Drop the redundant schema() override and the single-use sample_batch helper - Use checkSparkAnswerAndFallbackReason in the fallback tests so they assert the reason, not just the absence of the Comet operator - Trim documentation that restated the same paragraph in four places - Point the contributor guide at the directory operator serdes actually live in * fix: shim the SampleExec seed for Spark 4.2 Spark 4.2 changed SampleExec.seed from Long to Option[Long] and resolves an absent seed into a resolvedSeed field, which is what the operator samples with. Read the seed through CometSampleShim, following the CometCollectShim precedent for the same 4.2 divergence. * test: expand sample operator test coverage and add benchmark Add tests for ordering preservation above a sort, sampling above an aggregate, SQL TABLESAMPLE, and empty batch interleaving in the native operator. Add a Sample case to CometExecBenchmark comparing against Spark.
| Commit: | 38c3f8b | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
fix: throw ARITHMETIC_OVERFLOW for Long.MinValue div -1 under ANSI mode (#5084) * fix: throw ARITHMETIC_OVERFLOW for Long.MinValue div -1 under ANSI mode * refactor: reuse Spark's checkDivideOverflow and dedupe overflow check * test: use ARITHMETIC_OVERFLOW condition name in integral divide expect_error Matches the convention used by the other expect_error queries in arithmetic_ansi.sql. Verified against the spark-3.4, spark-3.5, spark-4.0, spark-4.1, and spark-4.2 profiles.
| Commit: | 8a9473c | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
feat: native randstr implementation compatible with Spark (#5035) * feat: native randstr implementation compatible with Spark * refactor: gate randstr shape restrictions in getSupportLevel, avoid per-row UTF-8 scan Move the literal-length/literal-seed/non-negative-length restrictions from convert into getSupportLevel + getUnsupportedReasons so they surface in EXPLAIN and the compatibility docs, matching the rand/randn serde pattern. Build each string in a reused String buffer instead of validating a byte buffer per row. * test: add golden-value, partition-index, and filter coverage for randstr Address review feedback on the randstr PR: - Rust golden-value test asserting bit-for-bit equality with Spark 4.1.1 across positive, zero, and negative seeds, independent of the SQL tests. - Rust partition-index test exercising the seed + partition_index arithmetic. - Broaden the zero-length test across a wider seed set (seed is irrelevant when length is zero). - Rust large-length smoke test guarding the builder capacity math, and saturating_mul on the capacity hint to avoid overflow on huge lengths. - SQL test asserting an Int seed and its Long literal produce identical output (serde toLong sign extension). - SQL test running randstr through a native projection feeding a filter.
| Commit: | 69df0f1 | |
|---|---|---|
| Author: | Peter Lee | |
| Committer: | GitHub | |
feat: add CalendarIntervalType support (#4898) * Add CalendarIntervalType Arrow support * Support interval vectors in UDF and shuffle codegen * remove stale test * spotless apply
| Commit: | 2e907e5 | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
feat: native collect_list / array_agg aggregate (#4720) * feat: native collect_list / array_agg aggregate Wires Spark's CollectList aggregate to datafusion-spark's SparkCollectList. array_agg, registered as a SQL alias of CollectList in FunctionRegistry, is also covered. Closes #2524. * fix: fall back collect_list/collect_set in multi-stage distinct aggregates A distinct aggregate combined with collect_list/collect_set produces a multi-stage plan (Partial -> PartialMerge -> Final). CollectList/CollectSet declare a BinaryType buffer in Spark but produce a native ArrayType state, so Comet cannot read a Spark-produced Binary buffer, nor round-trip its own ArrayType buffer across the intermediate PartialMerge stages. Both led to native crashes ("could not cast Binary to List" / "cast List to Binary"). Force these multi-stage aggregates to fall back to Spark consistently: - tag the feeding Partial when a PartialMerge stage of CollectList/CollectSet is present (CometExecRule.tagUnsafePartialAggregates), and - fall back a PartialMerge stage whose buffer was produced by a Spark partial (CometBaseAggregate.doConvert). Two-stage collect_list/collect_set continue to run natively. Patch the upstream SPARK-22223 plan-shape test to disable Comet, since native collect_list removes the ObjectHashAggregateExec it asserts on. Enabling fully-native multi-stage execution is tracked in #4724. * feat: fall back collect_list/collect_set on Spark 4.2 RESPECT NULLS Spark 4.2 adds an ignoreNulls field to CollectList and CollectSet, and collect_list(x) RESPECT NULLS sets it to false, keeping null elements. The native path delegates to SparkCollectList/SparkCollectSet, which always drop nulls, so it would silently return a different result from Spark. Add a per-version CometCollectShim that reads ignoreNulls (always true on Spark 3.4 through 4.1, where the field is absent) and fall back to Spark in getSupportLevel when it is false. Also rename QueryPlanSerde.hasIncompatibleBufferAgg to hasNativeArrayBufferAgg to describe what it detects: an aggregate whose native ArrayType state cannot round-trip Spark's declared BinaryType buffer. * test: accept Comet operators in SPARK-22223 instead of disabling Comet The SPARK-22223 ObjectHashAggregate test asserts on the executed plan. With Comet enabled, collect_list runs natively as CometHashAggregateExec, so ObjectHashAggregateExec is no longer present. Rather than disabling Comet, update the operator assertion to also accept CometHashAggregateExec, matching the pattern already used elsewhere in the diffs. The exchange assertion already matches ShuffleExchangeLike, which CometShuffleExchangeExec implements, so the single-shuffle check still holds. * style: reflow adjustOutputForNativeState doc comment Adding CollectList to the comment pushed a line past the column limit during the apache/main merge; reflow to satisfy spotless. * fix: coerce collect_list/collect_set nested field nullability collect_list and collect_set build their result list with all element fields marked nullable, but SparkCollectList/SparkCollectSet derive the return type from the child, preserving non-nullable nested fields. When the child is a nested type with a non-nullable inner field (e.g. a struct field built from non-nullable columns), the declared aggregate output disagrees with the array the accumulator produces, and the grouped native AggregateExec fails validating its output batch with "column types must match schema types". Cast the collect child to the all-nullable variant of its type so the declared and produced types stay consistent. * review: drop unreachable buffer-source block; consolidate 3.4/3.5 shim; add FILTER + map tests operators.scala: drop the CollectList/CollectSet PartialMerge fallback block. It is unreachable now that the general missingCometProducer + aggsNotSupportingMixedExecution guard just above returns None for the same case. CometExecRule.scala: add a comment on the collect-specific tagging block explaining it is separate from the tagging block just above because canAggregateBeConverted skips the child-native check, so an all-native distinct collect chain would otherwise slip past. QueryPlanSerde.scala: narrow the hasNativeArrayBufferAgg doc comment to describe what the code actually matches and note that Percentile has the same shape but is not matched here. Consolidate CometCollectShim: move the identical spark-3.4 and spark-3.5 copies to spark-3.x. Add collect_list FILTER (WHERE ...) test and a map-input test hitting make_all_fields_nullable's Map arm (size-only via spark_answer_only, since sort_array cannot order a MapType). * test: cover collect_list/collect_set RESPECT NULLS fallback on Spark 4.2 Add a Spark 4.2-gated SQL file test asserting that collect_list and collect_set with RESPECT NULLS fall back to Spark with Spark-identical results, and expand the getSupportLevel comments to explain that the fallback branch is only reachable on 4.2+ (a no-op on 3.4 through 4.1).
| Commit: | bc53ad5 | |
|---|---|---|
| Author: | Matt Butrovich | |
| Committer: | GitHub | |
feat: Iceberg table format V3: native table decryption, fall back for other V3 features (#4991)
| Commit: | 0761e54 | |
|---|---|---|
| Author: | Matt Butrovich | |
| Committer: | GitHub | |
perf: dedupe Iceberg residuals and delete files in native scan serde (#4982)
| Commit: | f06aa31 | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
feat: support approx_count_distinct aggregate expression (#4819) * feat: support approx_count_distinct aggregate expression Add native support for Spark's approx_count_distinct, a faithful port of Spark's HyperLogLogPlusPlus / HyperLogLogPlusPlusHelper. Each non-null input is hashed with Comet's Spark-compatible XxHash64 (seed 42, floats normalized first), and the HyperLogLog++ registers are stored in Spark's exact packed-Long buffer layout (10 six-bit registers per word). The cardinality is estimated with the same linear-counting and bias-correction tables Spark uses, so results are bit-identical to Spark and the partial-aggregation state matches Spark's aggBufferSchema. Includes a vectorized GroupsAccumulator, SQL file tests comparing against Spark across a range of cardinalities, native unit tests, benchmark coverage, and documentation updates. * feat: address review feedback on approx_count_distinct - restrict decimal inputs to precision <= 18 (Spark hashes wider decimals through BigDecimal, which the native i128 path does not match) - set supportsMixedPartialFinal=true; the register buffer matches Spark's aggBufferSchema, enabling mixed Comet/Spark partial and final aggregation - add getUnsupportedReasons so the Compatibility Guide reflects the type limits - derive Hash, add assert/debug_assert invariants and a float-order comment, reuse the shared normalize_float, drop a redundant field, make bias correction lazy, and reuse a per-accumulator hash scratch buffer - compact BIAS_DATA layout with rustfmt::skip (values unchanged) - add tests: decimal boundary and wide-decimal fallback, non-UTC timestamp, collated-string fallback, and mixed partial/final plan shape
| Commit: | 1981091 | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
feat: add spark.comet.shuffle.maxBufferBytes to cap native shuffle writer memory (#4989)
| Commit: | c65a5ee | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
feat: support gzip compression in native Parquet writes (#4930) * feat: support gzip compression in the native Parquet writer Add Gzip to the CompressionCodec proto enum and introduce a Parquet-specific ParquetCompression enum in parquet_writer.rs so the native writer can honor gzip without affecting the shuffle codec. The planner maps the new proto variant to ParquetCompression::Gzip, which writes with parquet's default GzipLevel (6), matching parquet-mr's zlib default. The shuffle writer path is unchanged and still rejects Gzip via its catch-all error arm. * feat: write Parquet natively when the compression codec is gzip * fix: request parquet-rs codec features explicitly for native Parquet writer The native writer's gzip, snappy, lz4, and zstd support only worked because Cargo feature unification happened to pull in parquet-rs's flate2/snap/lz4/zstd features via other workspace dependencies. Declare these features directly on the parquet dependency so the native writer does not silently lose codec support if that unification changes. * fix: honor parquet.compression precedence and uncompressed codec alias parseCompressionCodec only checked the compression option and the SQLConf default, skipping the parquet.compression option that Spark's own ParquetOptions treats as the middle rung of precedence. Fix the lookup to match Spark's compression, parquet.compression, SQLConf order, and accept "uncompressed" as an alias for "none" since Spark does the same. Add coverage for the parquet.compression option taking precedence over the SQLConf default, for uncompressed as a none alias, and for an unsupported codec (brotli) causing the write to fall back to Spark's own writer instead of CometNativeWriteExec. * refactor: inline single-call-site compression_to_parquet delegate compression_to_parquet was a one-line wrapper around ParquetCompression::to_parquet with a single caller; call it directly instead. * test: replace brittle brotli fallback test with lz4_raw round trip The unsupported-codec fallback test relied on Spark's write failing with ClassNotFoundException for BrotliCodec, so it only passed because the environment lacks a Brotli codec class rather than because Comet routed the write through the fallback path. Use lz4_raw instead, which Spark can write directly via parquet-mr and Comet does not support, so the test can assert a real successful round trip and the correct footer codec instead of intercepting an unrelated failure. Revert the allowFailure plumbing added to captureWritePlan for that test since it is no longer needed. * test: skip lz4_raw fallback test on Spark 3.4 Spark 3.4's PARQUET_COMPRESSION config validates against a fixed set of codecs (brotli, uncompressed, lz4, gzip, lzo, snappy, none, zstd) that does not include lz4_raw, so withSQLConf threw before the fallback path ran. Guard the test with assume(isSpark35Plus, ...) so it runs on versions where lz4_raw is a valid codec value. * test: cover full three-tier codec precedence for Parquet writes Adds a test that pins the `compression` > `parquet.compression` > `spark.sql.parquet.compression.codec` precedence: SQLConf is `zstd`, `parquet.compression` is `snappy`, and `compression` is `gzip`; the written file must report GZIP. A leak from either lower layer would surface as a codec mismatch. The existing test still covers `parquet.compression` beating the SQLConf default when `compression` is absent.
| Commit: | 88039f1 | |
|---|---|---|
| Author: | Oleks V | |
| Committer: | GitHub | |
chore: use Datafusion `substring` (#4161)
| Commit: | eb5b761 | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
feat: support approx_percentile / percentile_approx aggregate (#4801)
| Commit: | a282d29 | |
|---|---|---|
| Author: | Matt Butrovich | |
| Committer: | GitHub | |
chore: [branch-17] backport #4760 (#4829) * fix: size Iceberg delete files in the native scan to avoid dropping deletes (#4760) (cherry picked from commit b70e529ae945393ac24cb2647d560de1cd747f2a) * run prettier
| Commit: | a4b6bc9 | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
feat: support shuffle array expression (#4797) * feat: support shuffle array expression Add Comet support for Spark's `shuffle` array expression with full Spark compatibility. Comet reproduces Spark's exact random permutation by porting Apache Commons Math3's MersenneTwister and the inside-out Fisher-Yates algorithm from RandomIndicesGenerator, combining the resolved random seed with the partition index like Spark does. The native ShuffleExpr is a stateful PhysicalExpr that carries the RNG state across batches within a partition, mirroring the existing Rand/Randn implementation. * review: address feedback on shuffle expression - rename RNG to PRNG in doc comments - drop unreachable DataType::Null arm from ShuffleExpr::evaluate
| Commit: | f994b23 | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
feat: support PreciseTimestampConversion for native batch time-window grouping (#4784) * feat: support PreciseTimestampConversion for native batch time-window grouping Wire Spark's internal PreciseTimestampConversion expression, which the analyzer emits when resolving window()/session_window() grouping. It is a pure reinterpret between the timestamp types and Long, so it maps to an Arrow cast between microsecond Timestamp and Int64. Also support the KnownNullable tagging expression that the window resolution wraps around window bounds. Together these let batch tumbling and sliding time-window aggregations run natively. * test: convert time-window tests to a Comet SQL file test; simplify serde Convert the window() coverage from Scala tests to a SQL file test at expressions/datetime/window.sql. The default query mode still asserts native execution (checkSparkAnswerAndOperator), and a ConfigMatrix over session timezones confirms timestamp-window results match Spark in every zone. Simplify CometPreciseTimestampConversion.convert to use the existing optExprWithFallbackReason helper instead of a hand-rolled if/else. * docs: point session_window support entry at dedicated issue #4785
| Commit: | 464afe1 | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
feat: support interval types and make_ym_interval / make_dt_interval (#4541) * feat: support interval types and make_ym_interval / make_dt_interval [skip ci] Implements the type-support prerequisite from issue #4540: add Spark YearMonthIntervalType and DayTimeIntervalType as physical types that round-trip through Comet's Arrow FFI, and route the make_ym_interval / make_dt_interval constructors through the JVM codegen dispatcher so they execute natively and match Spark exactly. Type plumbing: - proto: add YEAR_MONTH_INTERVAL (18) and DAY_TIME_INTERVAL (19) to DataTypeId. - native serde.rs: map them to Arrow Interval(YearMonth) and Duration(Microsecond) respectively. DayTime stores microseconds in an int64, which matches Duration(Microsecond) rather than the lossy Interval(DayTime) {days, millis}. - Utils.toArrowType / fromArrowType: same mapping on the JVM side. - QueryPlanSerde.serializeDataType: emit the new type ids. - CometBatchKernelCodegen: accept the two interval types in isSupportedDataType and resolve IntervalYearVector / DurationVector; emit primitive set() writes. Expressions: - make_ym_interval -> CometMakeYMInterval, make_dt_interval -> CometMakeDTInterval, both via CometCodegenDispatch, registered in temporalExpressions. CalendarIntervalType and interval arithmetic remain follow-ups under #4540. Tests: SQL file tests for both constructors assert answer parity and native execution (checkSparkAnswerAndOperator), covering column and literal inputs, defaults, negatives, and nulls. * test: add overflow cases for make_ym_interval and make_dt_interval Confirm the codegen-dispatch path propagates Spark's arithmetic-overflow exception identically. The expect_error pattern uses the lowercase word overflow so it matches every Spark version: 4.x raises INTERVAL_ARITHMETIC_OVERFLOW while 3.x raises a raw ArithmeticException.
| Commit: | d2d976e | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
feat: support exact percentile and median aggregates natively (#4542) * feat: support exact percentile aggregate natively [skip ci] Wire Spark's exact `Percentile` aggregate (and the `percentile_cont` ANSI form, which Spark rewrites to `Percentile`) to DataFusion's `percentile_cont` aggregate. DataFusion uses the same `index = p * (n - 1)` linear interpolation as Spark, so results match for the common single-percentage form. - proto: add `Percentile` AggExpr message (child, percentage, datatype). - native planner: map it to `percentile_cont_udaf()` with [child, percentile]. - CometPercentile serde: Compatible for a single literal double percentage, default frequency, and numeric input; the child is cast to double so the native result is DoubleType. Array-of-percentages, a non-default frequency argument, and interval inputs fall back to Spark. - operators.adjustOutputForNativeState: map Percentile's TypedImperativeAggregate Binary partial buffer to the native List<Float64> state (ArrayType(DoubleType)), mirroring CollectSet, so the partial/shuffle/final exchange schema is correct. Codegen dispatch is not applicable: aggregates (TypedImperativeAggregate) cannot run in the per-row scalar kernel, so native is the only path. Tests: SQL file test covering global, grouped, integer-input, all-null, exact and interpolated percentiles, plus fallback assertions for the array and frequency forms. No new regressions in the SQL suite. * bench: add percentile cases to CometAggregateExpressionBenchmark [skip ci] * feat: guard percentile DESC fallback, broaden tests and docs Address audit findings on the native percentile aggregate: - getSupportLevel now falls back for the descending WITHIN GROUP form (percentile_cont/disc WITHIN GROUP ... ORDER BY ... DESC on Spark 4.0+), where Percentile.reverse=true. The native percentile_cont always interpolates ascending, so the descending form would return a wrong answer. - Extract the fallback reason strings into shared private vals and add a getUnsupportedReasons override so they reach the compatibility guide. - Expand percentile.sql with long/float/decimal/smallint/tinyint inputs, negative values, and median() coverage (median rewrites to percentile). - Add percentile_within_group.sql (Spark 4.0+) covering the ascending native path and the descending fallback. - Mark percentile, percentile_cont, and median as supported in the expression reference and record the cross-version audit in agg_funcs.md. - Document the DataFusion 6-decimal interpolation quantization (#4719). * fix: mark percentile aggregates Incompatible by default (#4719) DataFusion's percentile_cont quantizes the linear interpolation weight to 6 decimal places, so a deeply-interpolated percentile can differ from Spark by up to roughly (upper - lower) * 1e-6. Gate the otherwise-supported percentile / median / percentile_cont form as Incompatible so it falls back to Spark by default and is opt-in via spark.comet.expression.Percentile.allowIncompatible=true. Add the allowIncompatible config to the percentile SQL file tests so the native path stays covered, and update the expression reference and audit doc to reflect the opt-in status. * style: wrap percentile precision comment under 100 chars
| Commit: | b70e529 | |
|---|---|---|
| Author: | Matt Butrovich | |
| Committer: | GitHub | |
fix: size Iceberg delete files in the native scan to avoid dropping deletes (#4760)
| Commit: | ba6429f | |
|---|---|---|
| Author: | Oleks V | |
| Committer: | GitHub | |
feat: extend native windows support (#4209)
| Commit: | 1ff9555 | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
fix: reject Parquet INT96 as TimestampNTZ on Spark 3.x (#4357) * fix: reject Parquet TimestampLTZ as TimestampNTZ on Spark 3.x for native_datafusion Pre-Spark-4 (SPARK-36182) rejects reading a Parquet TimestampLTZ column as TimestampNTZ; native_datafusion previously did not, and silently returned the UTC instant. Plumb a per-Spark-version flag from ShimCometConf through the NativeScan proto into SparkParquetOptions, and gate a new rejection arm in the schema adapter on it. INT96 remains a gap because DataFusion's coerce_int96 strips the source timezone before the schema adapter runs, so it is indistinguishable from a true TIMESTAMP_NTZ source. Compatibility guide updated to describe the correctness implications. * style: cargo fmt * fix: reject Parquet INT96 as TimestampNTZ on Spark 3.x for native_datafusion Closes the INT96 gap left by the parent commit. INT96 columns previously surfaced as Timestamp(us, None) on the Rust side because DataFusion's coerce_int96 stripped the timezone, making them indistinguishable from a true TimestampNTZ source. With the new coerce_int96_tz option (DataFusion PR apache/datafusion#22318) we ask DataFusion to coerce INT96 to Timestamp(us, Some("UTC")), restoring the LTZ signal the schema adapter already pattern-matches against. Comet-side change is small: set coerce_int96_tz = "UTC" alongside the existing coerce_int96 = "us"; unskip the INT96 + native_datafusion variant of ParquetTimestampLtzAsNtzSuite. The schema adapter's existing Timestamp(_, Some(_)) -> Timestamp(_, None) rejection now fires for INT96 reads as well. EXPERIMENTAL: pinned via [patch.crates-io] to an andygrove/datafusion fork branch. Cannot be merged until apache/datafusion#22318 ships in a release.
| Commit: | 175a77d | |
|---|---|---|
| Author: | Bhargava Vadlamani | |
| Committer: | GitHub | |
Revert "feat: Native Broadcast nested loop join support (#4429)" This reverts commit 4c88f5d4863f55c370b05cc104dcd0950b3ad2bd.
| Commit: | 4c88f5d | |
|---|---|---|
| Author: | Bhargava Vadlamani | |
| Committer: | GitHub | |
feat: Native Broadcast nested loop join support (#4429) * native_support_broadcast_nested_loop_join
| Commit: | 9e86dd9 | |
|---|---|---|
| Author: | Matt Butrovich | |
| Committer: | GitHub | |
perf: replace CometBatchIterator FFI input path with the Arrow C Stream Interface (#4572)
| Commit: | aa6be27 | |
|---|---|---|
| Author: | Matt Butrovich | |
| Committer: | GitHub | |
perf: avoid FFI import/export between native subtree and ShuffleWriter (#4507)
| Commit: | a08cb4e | |
|---|---|---|
| Author: | Matt Butrovich | |
| Committer: | GitHub | |
feat: vendor-pluggable S3 credentials for native scans (#4309)
| Commit: | b23b760 | |
|---|---|---|
| Author: | Parth Chandra | |
| Committer: | GitHub | |
feat: implement make_time and to_time (#4256)
| Commit: | 0ca37e1 | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
feat: add support for `posexplode` and `posexplode_outer` (#4270)
| Commit: | dc08a96 | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
fix: complete native_datafusion Parquet schema-mismatch rejections (#4229)
| Commit: | 8119b1e | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
feat: add JVM UDF framework for native execution (#4232)
| Commit: | cf06ffb | |
|---|---|---|
| Author: | Matt Butrovich | |
| Committer: | GitHub | |
feat: support Parquet field ID matching in native_datafusion scan (#4216)
| Commit: | fbadc91 | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
fix: support Spark 4.1 BloomFilter V2 format and bit-scattering (#4196)
| Commit: | 7d5884f | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
fix: [Spark 4.1.1] preserve parent struct nullness when all requested fields missing in Parquet (#4190)
| Commit: | da187f2 | |
|---|---|---|
| Author: | Oleks V | |
| Committer: | GitHub | |
feat: support `PartialMerge` aggregation mode (#4003)
| Commit: | 0d49389 | |
|---|---|---|
| Author: | Liang-Chi Hsieh | |
| Committer: | GitHub | |
feat: support regular BuildRight+LeftAnti hash join (#4073)
| Commit: | 2bd01af | |
|---|---|---|
| Author: | hsiang-c | |
| Committer: | GitHub | |
feat: Support Spark expression: arrays_zip (#3643) * Define ArraysZip expr proto * Create ArraysZip SerDe * Register SerDe to arrayExpressions * Add SQL test * Register expression to planner * Rust wrapper around DF's arrays_zip * Null checks * Update supported Spark expressions doc
| Commit: | d6d5f09 | |
|---|---|---|
| Author: | Oleks V | |
| Committer: | GitHub | |
feat: support `collect_set` (#3954)
| Commit: | 9a7e616 | |
|---|---|---|
| Author: | ChenChen Lai | |
| Committer: | GitHub | |
feat: Support Spark expression hours (#3804) * feat: Add Spark V2 partition transform `Hours` to calculate hours since epoch from timestamps.
| Commit: | 9d3e166 | |
|---|---|---|
| Author: | Parth Chandra | |
| Committer: | GitHub | |
fix: Make cast string to timestamp compatible with Spark (#3884) * fix: Make cast string to timestamp compatible with Spark Add addtional formats and handle edge cases. Update compatibility guide Spark version specific behaviour for cast string to timestamp
| Commit: | 9b2f1b1 | |
|---|---|---|
| Author: | Liang-Chi Hsieh | |
| Committer: | GitHub | |
feat: support LEAD and LAG window functions with IGNORE NULLS (#3876) * feat: support LAG window function with IGNORE NULLS - Add ignore_nulls field to WindowExpr proto message - Serialize Lag window function with its ignoreNulls flag in CometWindowExec - Extend find_df_window_function to also look up WindowUDFs (not just AggregateUDFs) - Pass ignore_nulls to DataFusion's create_window_expr - Enable previously-ignored LAG tests Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com> * fix: make LAG window tests deterministic by adding secondary sort key ORDER BY b alone has ties, causing Spark and DataFusion to produce different but both-valid row orderings. Add c as a secondary sort key so tie-breaking is deterministic and results are comparable. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com> * fix: enable WindowExec allowIncompatible for LAG tests CometWindowExec is marked Incompatible by default. Add allowIncompatible=true config so LAG tests actually run via Comet and checkSparkAnswerAndOperator can verify native execution. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com> * fix: skip partition/order validation for offset window functions LAG/LEAD (FrameLessOffsetWindowFunction) support arbitrary partition and order specs. The existing validatePartitionAndSortSpecsForWindowFunc check (which requires partition columns == order columns) is only needed for aggregate window functions, not offset functions. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com> * feat: add LEAD support and IGNORE NULLS test for LAG window function - Handle Lead case in windowExprToProto alongside Lag - Add comment explaining hasOnlyOffsetFunctions guard - Add test for LAG IGNORE NULLS - Enable LEAD tests with allowIncompatible config and deterministic ORDER BY Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com> * test: add LEAD with IGNORE NULLS test case Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com> * test: add SQL tests for LAG/LEAD window functions with IGNORE NULLS Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com> --------- Co-authored-by: Claude Sonnet 4.6 <noreply@anthropic.com>
| Commit: | e5e452a | |
|---|---|---|
| Author: | Liang-Chi Hsieh | |
| Committer: | GitHub | |
feat: support SQL aggregate FILTER (WHERE ...) clause in native execution (#3835)
| Commit: | 8bab4a5 | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
perf: stop using FFI in native shuffle read path (#3731)
| Commit: | f57e54a | |
|---|---|---|
| Author: | Parth Chandra | |
| Committer: | GitHub | |
feat: [ANSI] Ansi sql error messages (#3580) * feat: [ANSI] Ansi sql error messages
| Commit: | 1a4eef6 | |
|---|---|---|
| Author: | Matt Butrovich | |
| Committer: | GitHub | |
chore: bump iceberg-rust dependency to latest [iceberg] (#3606)
| Commit: | d0aa1ff | |
|---|---|---|
| Author: | Matt Butrovich | |
| Committer: | GitHub | |
perf: Add Comet config for native Iceberg reader's data file concurrency (#3584)
| Commit: | a6741e8 | |
|---|---|---|
| Author: | Matt Butrovich | |
| Committer: | GitHub | |
feat: CometNativeScan per-partition plan serde (#3511)
| Commit: | 28e13dd | |
|---|---|---|
| Author: | Matt Butrovich | |
| Committer: | GitHub | |
feat: CometExecRDD supports per-partition plan data, reduce Iceberg native scan serialization, add DPP for Iceberg scans (#3349)
| Commit: | e9dafd0 | |
|---|---|---|
| Author: | Kazantsev Maksim | |
| Committer: | GitHub | |
Feat: to_csv (#3004)
| Commit: | ec5df97 | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
feat: Add support for round-robin partitioning in native shuffle (#3076)
| Commit: | 077005c | |
|---|---|---|
| Author: | Parth Chandra | |
| Committer: | GitHub | |
perf: [iceberg] Use protobuf instead of JSON to serialize Iceberg partition values (#3247) * perf: Use protobuf instead of JSON to serialize Iceberg partition values
| Commit: | f538424 | |
|---|---|---|
| Author: | Kazantsev Maksim | |
| Committer: | GitHub | |
Experimental: Native CSV files read (#3044)
| Commit: | e4a0142 | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
feat: Add support for `unix_timestamp` function (#2936)
| Commit: | aff07d0 | |
|---|---|---|
| Author: | Oleks V | |
| Committer: | GitHub | |
feat: Comet Writer should respect object store settings (#3042)
| Commit: | 2bf2835 | |
|---|---|---|
| Author: | B Vadlamani | |
| Committer: | GitHub | |
feat: Support ANSI mode avg expr (int inputs) (#2817)
| Commit: | fd53edb | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
feat: Add partial support for `from_json` (#2934)
| Commit: | 53e4092 | |
|---|---|---|
| Author: | Matt Butrovich | |
| Committer: | GitHub | |
perf: [iceberg] Deduplicate serialized metadata for Iceberg native scan (#2933)
| Commit: | 5ec12d4 | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
feat: Make shuffle writer buffer size configurable (#2899)
| Commit: | fd0ab64 | |
|---|---|---|
| Author: | B Vadlamani | |
| Committer: | GitHub | |
feat: Support ANSI mode SUM (Decimal types) (#2826)
| Commit: | 0bda9d2 | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
feat: Add support for `explode` and `explode_outer` for array inputs (#2836)
| Commit: | 1b3354b | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
feat: Partially implement file commit protocol for native Parquet writes (#2828)
| Commit: | 1ec3563 | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
feat: Add experimental support for native Parquet writes (#2812)
| Commit: | 937cacd | |
|---|---|---|
| Author: | Matt Butrovich | |
| Committer: | GitHub | |
feat: [iceberg] Native scan by serializing FileScanTasks to iceberg-rust (#2528)
| Commit: | fc3e6e9 | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
feat: Add support for `abs` (#2689)
| Commit: | bd3235f | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
chore: Remove code for unpacking dictionaries prior to FilterExec (#2659)
| Commit: | acfd03c | |
|---|---|---|
| Author: | B Vadlamani | |
| Committer: | GitHub | |
feat:support ansi mode rounding function (#2542)
| Commit: | c23dc25 | |
|---|---|---|
| Author: | Matt Butrovich | |
| Committer: | GitHub | |
feat: Parquet Modular Encryption with Spark KMS for native readers (#2447)