Skip to content
CAI
Software that uses CAICheck a score

apache/datafusion-comet

76.6

Strong · 27 September 2026

192.1k

lines of production code

Rust

with Scala

4

measurements over time

CAI band scale
CAI trend line
CAI lens gauges

What this system is

Apache DataFusion Comet is a high-performance acceleration engine for Apache Spark that offloads execution to a native Rust backend built on Apache DataFusion. It provides native implementations for a wide range of Spark expressions, operators, and shuffle mechanisms to replace slower JVM paths, while supporting direct scanning of Parquet, Iceberg, and CSV data from cloud storage. The system includes comprehensive tooling for benchmarking, release management, and local development to ensure compatibility and performance across multiple Spark versions.

How it got here

2024 — Initial project scaffolding and native engine foundation

25 changes.

This period established the Apache DataFusion Comet project through initial repository scaffolding, build configuration, and the definition of a dedicated exception hierarchy. It focused on laying the groundwork for the native execution engine by implementing core Protobuf schemas, native expression builders, and essential operators for Parquet and CSV scans. The work also introduced critical developer tooling, release automation scripts, and performance benchmarks to support the upcoming 1.1.0 release.

2025 — Native expression implementation and engine refactoring

24 changes.

This period focused on expanding native Spark compatibility by implementing a wide range of expressions, including string, datetime, array, map, JSON, and aggregate functions, alongside critical behavioral fixes for nullability and ANSI mode compliance. The underlying engine was significantly refactored through a modular registry-based planner, organized expression modules, and improved memory pool architecture to enhance maintainability and extensibility. Supporting infrastructure was also strengthened with local HDFS development tools, Azure Blob and S3 credential support for Parquet scans, and robust CI validation scripts.

2026 — Native shuffle infrastructure and benchmarking

24 changes.

This period focused on establishing a robust native shuffle engine, introducing a new JNI bridge and common library to handle complex data types, schema reconciliation, and error reporting. Significant effort was dedicated to performance optimization, including reusable compression contexts and remote shuffle service support, alongside the creation of comprehensive benchmark suites for PySpark, TPC, and micro-benchmarks. The work also laid the groundwork for future integrations by adding build gates for Delta and Lance, while enhancing native execution with features like JVM UDF support and dynamic filter pushdown.

Features

Add TPC-DS benchmark queries to repository

Added a collection of SQL benchmark queries (q1 through q37) derived from the TPC-DS standard to the \benchmarks/tpc/queries/tpcds\ directory. These files provide the specific workload definitions required to run the CometBench-DS performance benchmarks.

benchmarks/tpc/queries · high confidence

Add native RSS partition writer with task-owned JNI callbacks

Introduces a new native RSS (Remote Shuffle Service) partition writer that encodes shuffle blocks into Arrow IPC format and pushes them to a task-owned remote pusher via JNI callbacks. This implementation enforces JVM partition limits (i32), validates frame sizes against configurable reservation limits, and ensures block boundaries are preserved for remote concatenation, enabling efficient map-side shuffle push to external shuffle services.

native/shuffle/src/writers/rss · high confidence

Add native support for Spark UnsafeArray and UnsafeMap types in shuffle

This change introduces new Rust modules (\list.rs\, \map.rs\) within the \spark\_unsafe\ crate to handle Spark's native \UnsafeArray\ and \UnsafeMap\ data structures during shuffle operations. The implementation provides the necessary infrastructure to read array and map elements from JVM-allocated memory and append them to Arrow builders, enabling the native shuffle engine to correctly process nested complex types (lists and maps) that were previously unsupported or handled differently.

_native/shuffle/src/spark\unsafe · high confidence

Add standalone shuffle benchmark tool

A new standalone binary tool (\shuffle\_bench\) has been added to profile Comet shuffle write performance independently of Spark. It reads input directly from Parquet files and allows users to configure partitioning schemes (hash, single, round-robin), compression codecs, memory limits, and concurrent task simulation to measure write throughput and metrics.

native/shuffle/src/bin · high confidence

Added JNI utility for Parquet schema deserialization

A new \jni\ module has been added to the Parquet utilities, providing a \deserialize\_schema\ function that allows the native layer to parse Arrow IPC byte streams into Arrow schemas. This enables Java/Python bindings to efficiently transfer schema definitions from the native Rust core without requiring full serialization round-trips.

native/core/src/parquet/util · high confidence

Added build gate stub for native Lance integration

A new native Rust module for the Lance contrib component has been introduced, establishing a build-time gate for the feature. This location provides the core library structure and a planner stub that currently returns a 'not implemented' error, serving as the foundational placeholder for the native Lance read path within the DataFusion Comet integration.

contrib/lance/native · high confidence

Added build-gate stubs and configuration for native Delta scan integration

This change introduces the initial scaffolding for the native Delta scan feature within the \contrib/delta\ module. It adds a Rust build-gate stub (\contrib/delta/native\) that defines the \plan\_delta\_scan\ entry point, which currently returns a 'not implemented' error to allow compilation without the full implementation. Additionally, it registers new Spark configuration options under the \spark.comet.scan.deltaNative\ prefix (such as \enabled\, \fallbackOnUnsupportedFeature\, and \dataFileConcurrencyLimit\) to control the behavior of the future native Delta read path, ensuring these settings are discoverable in documentation and do not affect default builds.

contrib/delta, contrib/delta/native · high confidence

Added local HDFS development environment configuration

Users can now run a local Hadoop Distributed File System (HDFS) cluster for development and testing purposes. This change introduces a Docker Compose setup (\hdfs-docker-compose.yml\) and its associated environment configuration (\hadoop.env\) within the \kube/local\ directory. The configuration defines a cluster consisting of a NameNode, three DataNodes, a ResourceManager, and a NodeManager, allowing developers to spin up a complete HDFS environment locally using the specified BDE2020 Hadoop images.

kube/local · high confidence

Comet integration for Spark 3.4, 3.5, 4.0, and 4.1

This change integrates the Apache DataFusion Comet engine into Spark versions 3.4.3, 3.5.9, 4.0.4, and 4.1.3. It adds the \comet-spark\ dependency to the build, modifies \SparkSession\ to automatically load the \CometSparkSessionExtensions\ when enabled, and updates \SparkPlanInfo\ to include metadata for Comet scans. Additionally, it disables Comet in specific test suites (e.g., fallback storage, explain plans) where it is incompatible or causes precision differences.

dev/diffs · high confidence

Consolidated TPC benchmark suite with Iceberg and profiling support

The TPC benchmarking scripts have been consolidated into a unified runner (\run.py\) that supports TPC-H and TPC-DS benchmarks across multiple engines (Spark, Comet, Comet-Iceberg, Gluten). This update introduces native Iceberg table support, allowing users to benchmark Comet's native Iceberg scan acceleration by converting Parquet data to Iceberg format and running queries against the resulting tables. Additionally, the benchmark suite now includes built-in profiling capabilities via Java Flight Recorder (JFR) and async-profiler, enabling detailed performance analysis of Java and native code during execution. The suite also features consistency checks with result hashing to detect regressions and a new script to generate comparative speedup charts from benchmark results.

benchmarks/tpc · high confidence

Docker-based infrastructure for TPC benchmarks

Added Dockerfiles and Docker Compose configurations to containerize the TPC benchmark environment. The setup includes a base benchmark image with Java 8/17, Python, and async-profiler, a dedicated builder image for compiling Comet native libraries, and two Docker Compose profiles (laptop and cluster) to orchestrate Spark standalone clusters for running benchmarks.

benchmarks/tpc/infra · high confidence

Initial project scaffolding and repository configuration

The repository has been initialized with the foundational structure for Apache DataFusion Comet, including the Apache License 2.0, NOTICE file, and a Maven wrapper (version 3.2.0) to ensure consistent builds. A Makefile is provided to manage the native Rust and JVM build lifecycle, alongside configuration files for code formatting (scalafmt, prettier) and linting (scalafix). The project is configured as an Apache Software Foundation repository via \.asf.yaml\, which enforces a squash-merge workflow and sets up branch protection for release branches. Additionally, developer guidelines are established through \AGENTS.md\ and \CONTRIBUTING.md\, and the README is updated to reflect the project's Arrow-native acceleration capabilities.

(repo-wide) · high confidence

Introduce PySpark benchmark suite for shuffle performance comparison

A new PySpark benchmark framework has been added to the \benchmarks/pyspark\ directory, enabling users to compare shuffle performance across Spark, Comet JVM, and Comet Native implementations. The suite includes a data generation script (\generate\data.py\) that creates realistic test datasets with complex nested schemas, a base benchmark class (\benchmarks/base.py\) for standardized timing and configuration reporting, and a registry (\benchmarks/\\init\\_.py\) that currently supports \shuffle-hash\ and \shuffle-roundrobin\ benchmarks. Users can run individual benchmarks or execute the full suite via \run\_all\_benchmarks.sh\, which automatically resolves the Comet JAR and configures the necessary Spark settings for each mode.

benchmarks/pyspark · high confidence

Introduce micro benchmark runner for standalone execution

Added a new Python-based micro benchmark runner (\benchmarks/micro/run.py\) and its documentation, enabling users to execute Comet micro benchmarks on a dedicated machine (such as an EC2 instance) without cloning the entire repository. The script automates the full workflow—including environment setup, building Comet in release mode, running specific benchmark suites, collecting results, and publishing them via pull requests—while discovering test suites directly from the source tree.

benchmarks/micro · high confidence

Introduce native columnar-to-row conversion and aggregate merge support

This change adds the native implementation for converting Arrow columnar data into Spark's UnsafeRow format (columnar\_to\_row.rs), enabling faster row-based operations in Spark. It also introduces a MergeAsPartial wrapper (merge\_as\_partial.rs) that bridges DataFusion's aggregate modes to support Spark's PartialMerge aggregation semantics, allowing intermediate state merging to be executed natively.

native/core/src/execution · high confidence

Introduce native common library for error handling, schema reconciliation, and UTF-8 decoding

A new \native/common\ crate has been added to centralize shared native logic. It introduces a comprehensive \SparkError\ enum that maps native failures to specific Spark error codes (such as \CAST\_INVALID\_INPUT\, \ARITHMETIC\_OVERFLOW\, and \MALFORMED\_VARIANT\) to ensure consistent exception reporting. The library includes schema reconciliation helpers (\cast\_and\_stamp\_schema\, \widen\_nested\_nullability\) that normalize Arrow type drift—particularly nested nullability mismatches—before stamping schemas onto record batches. It also provides UTF-8 decoding utilities (\decode\_string\_arrays\, \decode\_utf8\_spark\_lossy\) that validate and decode string data at the JVM-to-native FFI boundary, matching the JDK's replacement behavior for invalid sequences. Additionally, the crate contains utilities for pushing struct null masks into children to prevent silent data corruption and a tracing recorder for performance profiling.

native/common/src · high confidence

Introduce official release process scripts for Apache DataFusion Comet

The \dev/release\ directory now contains a complete set of scripts to automate the official Apache Software Foundation source release process. This includes \build-release-comet.sh\ for building native binaries in Docker containers, \create-tarball.sh\ for generating signed tarballs and drafting the release vote email, \publish-to-maven.sh\ for uploading artifacts to a Nexus staging repository, and \verify-release-candidate.sh\ for validating release candidates. Supporting tools include \generate-changelog.py\ for automated changelog generation, \run-rat.sh\ and \check-rat-report.py\ for license compliance checking, and \release-tarball.sh\ for promoting approved candidates to the final distribution area.

dev/release · high confidence

Native Iceberg partitioning transforms and CSV struct serialization

This release adds native support for Iceberg's partitioning system functions—\iceberg\_bucket\, \iceberg\_truncate\, \iceberg\_years\, \iceberg\_months\, \iceberg\_days\, and \iceberg\_hours\—ensuring that hidden partitioning logic (hash distribution, local sort, and row-level filters) produces identical results to the Java implementation. It also introduces a native \to\_csv\ expression that converts struct columns to CSV strings, configurable via \CsvWriteOptions\ (delimiter, quote, escape, null value, and whitespace handling).

native/spark-expr/src · high confidence

Native Parquet scans now support Azure Blob Storage and location-scoped S3 credentials

Native Parquet scans can now read from Azure Blob Storage by translating Hadoop ABFS configurations (e.g., \fs.azure.\*\) into \object\_store\ options, supporting authentication via environment variables, account-scoped keys, and SAS tokens. For S3, the system now supports location-scoped credentials, allowing different paths within a bucket to use distinct credentials via a \LocationScopedObjectStore\ that routes requests based on policy locations. Additionally, S3-compliant filesystem aliases (like \blob://\) can be configured to route through the native S3 object store, with support for hostless URLs where the bucket is in the path.

native/core/src/parquet/objectstore · high confidence

Native S3 credential bridge for pluggable authentication

A new native module in the core cloud layer provides a JNI bridge to the JVM's CometS3CredentialDispatcher SPI, enabling the native scan paths (both raw Parquet via object\_store and Iceberg via iceberg-rust) to retrieve S3 credentials from pluggable Java providers. This change introduces per-location credential support and allows the native engine to use the same credential dispatch logic as the JVM side, ensuring consistent authentication behavior across scan implementations.

native/core/src/cloud · high confidence

Native Spark array functions now support insertion, position lookup, slicing, overlap checks, zipping, and flattening

The native execution engine now implements several Spark-compatible array functions in \native/spark-expr/src/array\_funcs\. Users can now insert items into arrays (\array\_insert\), find the position of an element (\array\_position\), slice arrays (\array\_slice\), check for overlapping elements between two arrays (\arrays\_overlap\), zip multiple arrays into structs (\arrays\_zip\), and flatten nested arrays (\flatten\). These additions also include extracting fields from struct arrays (\get\_array\_struct\_fields\) and extracting elements by index (\list\_extract\). The implementations handle Spark-specific behaviors such as three-valued null logic for overlaps, correct handling of negative indices for slices, and proper null propagation for nested structures.

_native/spark-expr/src/array\funcs · high confidence

Native Spark-compatible Bloom Filter implementation

Added native Rust implementations of the \bloom\_filter\_agg\ aggregate and \might\_contain\ scalar functions, along with the underlying \SparkBloomFilter\ and \SparkBitArray\ data structures. This enables native execution of Bloom filter operations with full compatibility to Spark's binary serialization formats, including support for both V1 (Spark ≤4.0) and V2 (Spark 4.1+) formats, correct big-endian byte ordering, and proper handling of the V2 seed field and bit-scattering algorithm.

_native/spark-expr/src/bloom\filter · high confidence

Native Spark-compatible map lookup and sorting functions

This change introduces native Rust implementations for Spark's \GetMapValue\ (map lookup) and \MapSort\ expressions in the \map\_funcs\ module. The new \SparkMapExtract\ function replaces the default DataFusion behavior for map lookups (\element\_at\ and \m\[k\]\) by returning the matched value directly instead of a one-element list, and uses vectorized Arrow equality checks to significantly improve performance for constant-key lookups. Additionally, the \spark\_map\_sort\ function provides a native implementation to sort map entries by key in ascending order, featuring optimized paths for single-entry maps and string-keyed maps to reduce overhead. These additions allow Comet to execute these map operations natively without falling back to Spark, improving both correctness alignment with Spark's semantics and execution speed.

_native/spark-expr/src/map\funcs · high confidence

Native dynamic filter pushdown for hash joins into Parquet scans

This change introduces native dynamic filter pushdown for hash joins into Parquet scans. The \DynamicFilterJoinExec\ operator connects a hash join's completed build domain to its probe input, allowing a direct Parquet reader to use the same live predicate for pruning. This optimization is applied within the \native/core/src/execution/operators/dynamic\_filter\ module, specifically through the \join.rs\ and \parquet\_reader.rs\ files, which handle the runtime filter wiring and schema adaptation. The implementation ensures that Parquet conversion errors are preserved during join filtering, as verified by the new tests in \join/tests/schema\_errors.rs\ and \join/tests/timestamp\_errors.rs\. This feature enhances query performance by reducing the amount of data read from Parquet files during join operations.

_native/core/src/execution/operators/dynamic\filter · high confidence

Native expression builders for arithmetic, bitwise, temporal, and random functions

The \native/core/src/execution/expressions\ module now provides native DataFusion expression builders for a broad set of Spark-compatible operations. Arithmetic and bitwise operations (Add, Subtract, Multiply, Divide, BitwiseAnd, BitwiseOr, BitwiseXor, BitwiseShiftLeft, BitwiseShiftRight) are handled via typed builders, with arithmetic operations wrapped in \CheckedBinaryExpr\ to catch and wrap Spark errors with query context. Comparison and logical operators (Eq, Neq, Lt, LtEq, Gt, GtEq, EqNullSafe, NeqNullSafe, And, Or, Not) are also implemented. New temporal builders support Hour, Minute, Second, UnixTimestamp, TruncTimestamp, and HoursTransform. Random and ID generators (Rand, Shuffle, Uuid, RandStr, Randn, SparkPartitionId, MonotonicallyIncreasingId) are wired to native implementations, with seeds adjusted per partition to match Spark's behavior. String operations include Like, Rlike, and partial support for FromJson. List position generation for posexplode is implemented via \ListPositionsExpr\, and subquery scalar results are fetched from the JVM via JNI.

native/core/src/execution/expressions · high confidence

Native implementation of LPAD, RPAD, and read-side padding functions

This change introduces native Rust implementations for the Spark \lpad\, \rpad\, and \read\_side\_padding\ string functions within the \static\_invoke\ module. These functions replace the previous DataFusion defaults to provide correct Spark-compatible behavior, specifically handling Unicode padding, supporting dictionary-encoded arrays, and allowing the length argument to be a column rather than just a literal. The implementation is organized under \char\_varchar\_utils\ and exposes \spark\_lpad\, \spark\_rpad\, and \spark\_read\_side\_padding\ for use by the query engine.

_native/spark-expr/src/static\invoke · high confidence

Native implementation of Spark JSON functions

The \native/spark-expr/src/json\_funcs\ module now provides native Rust implementations for key Spark JSON functions, replacing or supplementing previous logic. This includes \from\_json\ for parsing JSON strings into structured types, \json\_array\_length\ for determining the size of JSON arrays, and \to\_json\ for converting structured data into JSON strings. These changes enable more efficient, native execution of JSON operations within the DataFusion-based engine.

_native/spark-expr/src/json\funcs · high confidence

Native implementation of Spark-compatible ISNAN and RLIKE predicate functions

The native execution engine now includes dedicated Rust implementations for the Spark \ISNAN\ and \RLIKE\ predicate functions, replacing previous fallbacks or incomplete support. The \ISNAN\ function correctly handles Float32 and Float64 arrays and scalars, ensuring that null values are treated as false in accordance with Spark semantics. The \RLIKE\ function provides a physical expression that supports Utf8, LargeUtf8, and Utf8View string types, as well as dictionary-encoded arrays, using a pre-compiled regex pattern for performance. These changes are located in the \native/spark-expr/src/predicate\_funcs\ module, which organizes these specific Spark grouping expressions.

_native/spark-expr/src/predicate\funcs · high confidence

Native implementation of Spark-compatible parse\_url and try\_parse\_url UDFs

Users can now use native, high-performance versions of the \parse\_url\ and \try\_parse\_url\ functions. This change introduces a new module in the native Spark expression layer that implements these URL parsing UDFs using RFC 3986 regex parsing to ensure exact compatibility with Spark's \java.net.URI\ behavior, addressing edge cases where upstream DataFusion implementations diverged. The implementation includes optimized query parameter extraction with regex caching and proper handling of invalid URLs via the \try\_parse\_url\ variant.

_native/spark-expr/src/url\funcs · high confidence

Native implementation of non-deterministic functions (UUID, shuffle, randn)

This change introduces the internal Rust infrastructure required to support native, Spark-compatible non-deterministic functions. It adds a bit-for-bit port of the Mersenne Twister PRNG (seeded per partition to match Spark's behavior) and a utility framework for evaluating batch random values while maintaining state across partitions. This foundation enables the native implementation of \uuid()\, \shuffle\, and \randn\ expressions, ensuring their output matches Apache Spark's determinism and distribution guarantees.

_native/spark-expr/src/nondetermenistic\funcs/internal · high confidence

Native implementation of the IF conditional expression

The native Spark expression layer now includes a dedicated implementation for the IF function, located in \native/spark-expr/src/conditional\_funcs\. This change introduces the \IfExpr\ struct, which wraps DataFusion's \CaseExpr\ to evaluate conditional logic (returning a true value if a condition is met, otherwise a false value). The implementation correctly propagates nullability based on the inputs and includes unit tests to verify evaluation against various data types.

_native/spark-expr/src/conditional\funcs · high confidence

Native implementations for Spark datetime functions

The native engine now includes native implementations for several Spark datetime expressions, including \date\_diff\, \date\_from\_unix\_date\, \date\_trunc\, \dayname\, \monthname\, \dayofweek\, \weekday\, \hour\, \minute\, \second\, \hours\ (V2 partition transform), \make\_date\, \make\_interval\, and \make\_time\. These additions replace previous fallback or less optimized paths, providing direct Spark-compatible logic for date arithmetic, truncation, name extraction, and time construction.

_native/spark-expr/src/datetime\funcs · high confidence

Native implementations for base64, JSON extraction, Levenshtein, and regex functions

This change introduces native Rust implementations for several Spark string functions in the \spark-expr\ module, replacing slower or non-native paths. Users gain support for \base64\ encoding (with optional MIME chunking) and \get\_json\_object\ (extracting values via JSONPath). The \levenshtein\ function now supports an optional distance threshold for optimized computation, while \regexp\_extract\ and \regexp\_extract\_all\ provide native regex extraction with a pattern cache to avoid recompiling regular expressions per batch. Additionally, \concat\_ws\ is implemented to handle scalar subqueries correctly, and the \contains\ function is optimized for scalar needle patterns.

_native/spark-expr/src/string\funcs · high confidence

Native implementations for new and existing aggregate functions

This change adds native Rust implementations for several Spark-compatible aggregate functions in the \agg\_funcs\ module, including \approx\_percentile\ (using a Greenwald-Khanna port), \approx\_count\_distinct\ (using a HyperLogLog++ port with Spark-compatible hashing and bias tables), and \correlation\/\covariance\ (with Spark-compatible state fields and null-on-divide-by-zero behavior). It also introduces native \GroupsAccumulator\ implementations for \collect\_list\ and \collect\_set\ to optimize grouped aggregation performance, and adds native \avg\ and \avg\_decimal\ accumulators that handle type casting and decimal precision scaling consistent with Spark.

_native/spark-expr/src/agg\funcs · high confidence

Native implementations of Spark nondeterministic functions

This change introduces native Rust implementations for several Spark nondeterministic functions within the \nondetermenistic\_funcs\ module, ensuring bit-for-bit compatibility with Spark's random number generators. The update adds support for \monotonically\_increasing\_id\ (using atomic offsets), \rand\ and \randn\ (using Spark's XORShift and Marsaglia polar methods), \randstr\ (alphanumeric string generation), \shuffle\ (list shuffling via Mersenne Twister), \uuid\ (RFC 4122 v4 generation), and \BernoulliCellSampler\ (for sampling without replacement). These expressions are designed to maintain stateful, per-partition evaluation semantics consistent with Spark, allowing users to run these nondeterministic operations natively with predictable, Spark-compatible results.

_native/spark-expr/src/nondetermenistic\funcs · high confidence

Native memory usage observability and allocator selection

The native core library now provides process-wide accounting of bytes allocated by the Rust global allocator, exposing the current balance via \current\_balance()\ for integration with executor memory usage logs and tracing. This is achieved by wrapping the selected backend allocator (jemalloc, mimalloc, or system) in an \AccountingAllocator\ that tracks allocations with minimal thread-local overhead. Additionally, the library now exposes a JNI method \isFeatureEnabled\ to check if specific features like jemalloc or hdfs-opendal are compiled in, and prints the native library version upon initialization.

native/core/src · high confidence

Native temporal kernel implementations for date and timestamp operations

The \native/spark-expr/src/kernels\ module now contains native Rust implementations for temporal operations, including optimized date truncation (year, quarter, month, week) and timezone-aware timestamp conversions. These kernels handle edge cases such as Daylight Saving Time transitions and schema consistency when converting between timezones, providing the underlying logic for date and timestamp expressions in the Spark engine.

native/spark-expr/src/kernels · high confidence

New CI validation scripts for benchmarks, configuration, and documentation

Added a suite of Python and shell scripts in dev/ci to enforce CI invariants and prevent silent failures. check-benchmark-runner.py validates that the benchmark runner discovers the expected suites and correctly handles selection filters. check-ci-config.py verifies that the change-routing policy, event policies, required-check aggregation, and artifact naming remain consistent with the workflow definitions. check-mermaid.py ensures that documentation diagrams render correctly and are present in the built site, preventing missing diagrams from going unnoticed. check-suites.py confirms that all Scala test suites are included in the Linux and macOS PR workflows. check-working-tree-clean.sh ensures the repository state is clean before proceeding. compute-changes.py replaces the external paths-filter action with an internal script to determine which jobs should run based on changed files. linux-test-profiles.py and local-ci-config.py support local CI execution and profile selection. These changes improve CI reliability by catching configuration drift and broken pipelines early.

dev/ci · high confidence

New Docker-based build environment for Comet native libraries

A new Dockerfile and build script have been added to the release tooling to standardize the compilation of Comet native libraries. The Docker image is based on Ubuntu 20.04 and includes all necessary build dependencies, such as GCC 10, Clang, Rust, and the JDK. The entrypoint script accepts a Git repository, branch, and architecture (arm64 or amd64) to clone the source and execute the build, ensuring a consistent and reproducible release environment.

dev/release/comet-rm · high confidence

New analyze\_trace binary for native memory leak detection

A new command-line tool, \analyze\_trace\, has been added to the native common module. It analyzes Comet Chrome trace event logs to detect native memory leaks by comparing process-wide allocation counters (such as \native\_allocated\ or \jemalloc\_allocated\) against the total memory reserved by Comet's memory pools. The tool identifies points where allocated bytes exceed pool reservations, reporting peak excess and sample violations to help users diagnose memory issues.

native/common/src/bin · high confidence

New developer tooling and build configuration scripts

Added several new scripts and configuration files to the \dev/\ directory to support local development and release processes. This includes \local-ci.sh\ for running Spark SQL and Iceberg CI workflows locally, \generate-release-docs.sh\ for freezing documentation content onto release branches, and \regenerate-golden-files.sh\ for updating plan stability test fixtures. Additionally, \cargo.config\ configures specific linkers for Apple Silicon and Intel Macs, \checkstyle-suppressions.xml\ provides an empty suppression list for Java style checks, \ensure-jars-have-correct-contents.sh\ validates JAR artifacts against an allowed class list, and \verify-contrib-delta-gate.sh\ ensures Delta dependencies do not leak into default builds.

dev · high confidence

New exception hierarchy for Comet native and runtime errors

The common module now introduces a dedicated exception hierarchy to better distinguish and handle errors originating from the Comet native execution layer. New classes include CometRuntimeException as the base for general runtime issues, CometNativeException for native-side errors, and specific subclasses like CometOutOfMemoryError, CometShuffleSizeLimitException, and ParquetRuntimeException for targeted error handling. Additionally, CometQueryExecutionException is added to wrap JSON-encoded error details from native execution, enabling more precise error classification and recovery in Spark applications.

common · high confidence

New native JNI bridge crate for Apache DataFusion Comet

A new \native/jni-bridge\ Rust crate has been introduced to centralize the JNI interaction layer for Apache DataFusion Comet. This crate provides the specific bindings and helper structures required for native execution to communicate with the JVM, including handles for \CometExec\ (scalar subqueries), \CometMetricNode\, \CometTaskMemoryManager\, \CometUdfBridge\, and \CometShuffleBlockIterator\. It also introduces support for pluggable S3 credentials via \CometS3CredentialDispatcher\, case-insensitive schema resolution through \CometSchemaUtils\, and task-owned remote shuffle callbacks via \JavaShufflePartitionPusher\. Additionally, it standardizes error handling with a comprehensive \CometError\ enum and ensures reliable panic backtraces are captured and propagated across the FFI boundary.

native/jni-bridge · high confidence

Support for JVM UDFs in native execution

Users can now execute Java and Scala User-Defined Functions (UDFs) using the native Comet execution engine. This change introduces a new \JvmScalarUdfExpr\ component that delegates evaluation to JVM-side \CometUDF\ implementations via JNI. It ensures correct class loading by propagating the Spark task ClassLoader and maintains access to partition-sensitive context (like \TaskContext\) by threading the task context through the FFI boundary, allowing UDFs to correctly use features such as \Rand\ or \Uuid\.

_native/spark-expr/src/jvm\udf · high confidence

Support for shuffling empty-schema record batches

The native shuffle partitioners now handle zero-column schemas (such as those produced by COUNT(\*) after column pruning) via a new EmptySchemaShufflePartitioner. This partitioner accumulates the total row count from input batches and writes a single zero-column IPC batch to partition 0, ensuring that shuffle operations on empty-schema data complete correctly without errors.

native/shuffle/src/partitioners · high confidence

Architecture

Native planner adopts modular registry-based dispatch for operators and expressions

The native query planner in \native/core/src/execution/planner\ has been refactored to use a modular registry pattern for dispatching both Spark operators and expressions. New files \operator\_registry.rs\ and \expression\_registry.rs\ introduce global registries that map protobuf operator and expression types to specific builder implementations, replacing the previous monolithic dispatch logic. This change introduces a generic extension point (\OpStruct::ContribScan\) that allows out-of-tree contrib scans, such as Delta and Lance, to be routed through dedicated shims (\delta\_scan.rs\, \lance\_scan.rs\) without modifying core planner code. Helper macros (\macros.rs\) reduce boilerplate for unary and binary expression builders. For users, this provides a more extensible and maintainable foundation for adding new native operators and expressions, while ensuring that contrib-specific scan types are handled via feature-gated modules.

native/core/src/execution/planner · high confidence

Protobuf definitions moved to a separate crate with dedicated build script

The protocol buffer definitions for Spark physical query plans (expr, metric, partitioning, operator, and config) have been moved into a dedicated native/proto crate. This change introduces a build script that automatically generates Rust code from these .proto files into a src/generated directory during compilation, streamlining the integration of query plan intermediates for the Apache DataFusion Comet project.

native/proto · high confidence

Behavioural changes

Enable Comet native execution for Iceberg tests and benchmarks

The Iceberg test suites and benchmarks now run with the Apache DataFusion Comet plugin enabled. This integrates Comet's native scan, shuffle, and write capabilities into the Iceberg test infrastructure, allowing users to validate Iceberg operations against Comet's accelerated execution engine. The changes update build dependencies to include the Comet Spark module and configure Spark sessions in test bases and benchmarks to use the Comet shuffle manager and native Iceberg features.

dev/diffs/iceberg · high confidence

Maven wrapper and network retry configuration added

The project now includes a Maven wrapper (version 3.2.0) using Maven 3.9.6, ensuring consistent builds across environments. Additionally, a new \.mvn/maven.config\ file configures HTTP retry behavior to handle transient network issues more robustly: it increases the retry count to 6, expands retryable status codes to include 408, 500, 502, and 504 (in addition to 429 and 503), sets a TCP connect timeout of 30 seconds, and sets a request timeout of 600 seconds to prevent dead connections from consuming excessive build time.

.mvn · high confidence

Native Parquet Variant projection now normalizes and unshreds data at the boundary

The native Parquet reader now handles Variant columns by normalizing storage types (such as converting unsigned integers to signed equivalents and standardizing decimal/timestamp types) and unshredding nested structures directly at the scan boundary. This ensures that Variant data read from Parquet files is consistently formatted and validated before being exposed to the query engine, resolving compatibility issues with legacy Spark storage formats and rejecting unsupported types like UUIDs.

_native/core/src/parquet/cast\column · high confidence

Native Parquet scan performance and compatibility improvements

This release introduces an eager page-index reader factory that forces loading the Parquet page index on the first metadata fetch and caching it, eliminating repeated uncached I/O for predicates that do not fully match row-group statistics. It also adds a unified field-name folding module to ensure case-insensitive Parquet reads match Spark's \toLowerCase(Locale.ROOT)\ behavior consistently across the schema adapter, nested conversions, and plan-time projection, while introducing a debug batch stream and execution plan wrapper to print batch details for troubleshooting. Additionally, the scan now supports Parquet Modular Encryption with a Spark KMS integration via a JNI-based key retriever, and improves type conversion by aligning struct casting to Spark's behavior and rejecting duplicate Parquet field names before decoding.

native/core/src/parquet · high confidence

Native cast module refactored into typed submodules

The native cast implementation in \native/spark-expr/src/conversion\_funcs\ has been reorganized from a single monolithic file into distinct modules (\boolean\, \numeric\, \string\, \temporal\, \trim\, \utils\) to improve maintainability and performance. This change introduces specialized casting logic for each data type category, including Spark-compatible boolean-to-timestamp conversion, precise Java-style float/double string formatting, and strict whitespace trimming rules that match Spark's \trimAll\ and \String.trim\ behaviors. The refactoring also adds dedicated handling for date-to-timestamp conversions with correct timezone and DST resolution, and ensures decimal-to-string output matches Spark's legacy and ANSI modes.

_native/spark-expr/src/conversion\funcs · high confidence

Native execution operators reorganized and expanded with CSV support and FFI alignment fixes

The native execution operators have been restructured into individual modules (aligned\_stream\_reader, copy, csv\_scan, expand, explode, filter, iceberg\_common, iceberg\_partition\_path) to improve modularity and maintainability. A new native CSV scan operator has been added, enabling direct reading of CSV files through DataFusion's CsvSource. The FFI import path now includes an AlignedArrowStreamReader that explicitly aligns Decimal128 buffers to prevent panics when receiving data from the JVM, addressing alignment issues in the Arrow C Data Interface. Dictionary array handling has been enhanced with dedicated copy and unpack functions that properly handle dictionary values and preserve nullability. The Expand operator now correctly derives output schemas by widening nested nullability across all projections, fixing issues with nested field nullability mismatches. The Explode operator maintains a specialized fork of DataFusion's unnesting logic with performance optimizations for list output length computation and contiguous-run fast paths. Iceberg integration improvements include a custom location generator that matches iceberg-java's partition path rendering to prevent panics on timestamptz values, and enhanced S3 credential handling with configurable providers and region defaults.

native/core/src/execution/operators · high confidence

Native math functions now support ANSI mode error handling

The native math expression implementations in \native/spark-expr/src/math\_funcs\ have been updated to respect the ANSI evaluation mode. Functions such as \abs\, \ceil\, \floor\, \log\, \modulo\, and arithmetic operations now throw specific errors (e.g., \ARITHMETIC\_OVERFLOW\, \DIVIDE\_BY\_ZERO\, \REMAINDER\_BY\_ZERO\) when \fail\_on\_error\ is enabled, matching Spark's ANSI-compliant behavior. In legacy mode, these functions continue to use wrapping arithmetic or return nulls for invalid operations, ensuring backward compatibility while providing strict error reporting for users who require ANSI SQL compliance.

_native/spark-expr/src/math\funcs · high confidence

Native query plans are now serialized via a new set of Protobuf schemas

The native execution engine now uses a dedicated set of Protobuf definitions (config.proto, expr.proto, literal.proto, metric.proto, operator.proto, partitioning.proto, and types.proto) to serialize Spark plans, expressions, and types to the native side. This replaces the previous serialization approach, enabling richer expression support (including new array, JSON, and aggregate functions), more efficient plan transmission via SQL text interning, and better error reporting with query context.

native/proto/src/proto · high confidence

Native shuffle module restructured with performance and reliability improvements

The native shuffle module has been reorganized into a separate crate, introducing a \PartitionWriter\ interface to decouple partitioning logic from storage. This refactor includes significant performance optimizations: zstd and Arrow IPC compression contexts are now reused across shuffle blocks to reduce allocation overhead, and the IPC schema is encoded once per writer instead of per block. A new \ShuffleCodecContext\ manages these reusable contexts, while a per-thread schema cache with bounded retention limits memory usage during decoding. The module also adds a \bench\_support\ seam for benchmarking partitioning without the write path, implements a \RoundRobinStrategy\ with both hash-based and positional row-group placement, and introduces \SchemaAlignExec\ to correct type and nullability drift before partitioning to ensure consistent hash-based routing. Remote shuffle support includes dictionary unpacking and schema validation to safely handle data from the JVM.

native/shuffle/src · high confidence

Native shuffle plans now support partition writer destinations

The native protocol layer has been updated to include partition writer destinations within shuffle plans. This change allows the system to specify where shuffle data should be written (e.g., local files or remote shuffle services) directly in the plan structure, replacing the previous mechanism that relied on separate index paths. This enables more flexible and efficient data shuffling strategies by explicitly defining output destinations at the protocol level.

native/proto/src · high confidence

Native struct field extraction now correctly propagates parent nullability

The native struct expression implementation has been reorganized into a dedicated \struct\_funcs\ module, introducing \CreateNamedStruct\ and \GetStructField\ expressions. A key behavioral fix ensures that extracting a field from a struct correctly reports the field as nullable if the parent struct is nullable, even when the field itself is defined as non-nullable. This aligns with Spark semantics and prevents downstream Arrow validation errors during shuffles or sorts where a non-nullable field of a nullable struct would otherwise incorrectly claim to contain no nulls.

_native/spark-expr/src/struct\funcs · high confidence

Optimized decimal overflow handling and NaN normalization in native math expressions

The native math expression engine now uses a shared fast path for decimal overflow checks, significantly improving performance when no overflow occurs by reusing input buffers instead of scanning for errors. Decimal rescale and overflow validation are fused into a single pass, and scalar float NaN/zero normalization is optimized for consistent comparison semantics. Additionally, the unscaled value extraction for decimals is vectorized for faster conversion to integers.

_native/spark-expr/src/math\funcs/internal · high confidence

Refactored hash function implementation into a modular structure

The hash function logic has been reorganized into a dedicated \hash\_funcs\ module, splitting the code into separate files for Murmur3, xxHash64, and utilities. This change introduces a new \spark\_xxhash64\ entry point that delegates to DataFusion's \SparkXxhash64\ implementation when arguments are compatible (e.g., default seed 42, supported types), while retaining a native kernel for cases where the upstream implementation is incompatible (e.g., non-default seeds, structs with null masks, dictionaries in lists). The Murmur3 implementation remains native but is now structured to handle dictionary arrays and null masks more explicitly.

_native/spark-expr/src/hash\funcs · high confidence

Refactored memory pool architecture with new fair pool and task sharing

The memory pool subsystem has been refactored to introduce a new \fair\_unified\ pool type that enforces per-consumer memory limits, replacing the previous default \greedy\_unified\ strategy. The implementation now supports task-shared memory pools via an RAII guard, ensuring pools are shared across native execution contexts within the same task attempt and cleaned up automatically. Additionally, the codebase now includes a logging pool wrapper for debugging, improved overcommit handling to prevent panics during memory growth, and stricter checks for fair pool reservations.

_native/core/src/execution/memory\pools · high confidence

A new symbolic link at .claude/skills has been added, redirecting to the ../.ai/skills directory. This change establishes a reference from the .claude location to the .ai skills, likely supporting the transition or coexistence of vendor-neutral configurations.

.claude · medium confidence

Fixes

Fixes missing PruningMetrics and Ratio values in native scan metrics

The native metrics collection now correctly exposes DataFusion 54's PruningMetrics and Ratio types to the Spark side. Previously, these metrics were surfaced as zero because DataFusion reserves their aggregation to MetricsSet; the new utils module expands them into individual integer fields (e.g., pruned/matched for pruning, part/total for ratios) so that parquet scan filtering statistics are visible in Spark metrics.

native/core/src/execution/metrics · high confidence

Fixes native linking failures caused by stale JDK library paths

The native core build script now explicitly emits the JDK's libjvm search path during compilation. This ensures that the linker can always locate libjvm, preventing build failures on CI runners or environments where the JDK installation path changes or is cached incorrectly from previous builds.

native/core · high confidence

Native shuffle now reports spill metrics correctly

Users will now see accurate native shuffle read and write metrics, including spill information, in the Spark UI and logs. Previously, these metrics were either missing or incorrect, making it difficult to monitor shuffle performance and identify spill-related bottlenecks during query execution.

spark · high confidence

Test coverage

Added benchmarks for native shuffle read and write performance; Added native benchmarks for scalar and aggregate functions; Added test utility for generating unique temporary files; Added tests for Parquet Variant projection support; Added tests for Spark-compatible expression registration and behavior; Added tests for empty native scan schema preservation; New native micro-benchmarks for allocation overhead, array conversion, and unnesting.

Dependencies

Initial release scaffolding and dependency configuration for Comet 1.1.0

This change introduces the foundational build configuration for the 1.1.0 release, establishing the Maven parent POM and the Rust workspace structure. It defines the core native crates (core, spark-expr, common, proto, jni-bridge, shuffle) and pins key dependencies, including DataFusion 55.1.0, Arrow 59.2.0, and Iceberg-rust (pinned to revision bb1e4a4). The diff also adds build-gate stubs for optional contrib integrations (Delta, Lance) and configures CI tooling, such as Gradle-based Iceberg test sharding and documentation requirements.

(dependencies) · high confidence

Written by watchdog.canine.dev from the codebase's own history, inside the signed delivery this page is composed from.

How this codebase got here

This is the PUBLIC form of this artifact. Findings are listed in full, but the details of SECURITY findings — which rule fired, in which file, on which line, and how to fix it — are deliberately withheld, and any secret-scanner results are excluded entirely. Where detail is absent here it was REMOVED FOR PUBLICATION; it is not missing from the analysis. The complete artifact is available from the repository owner.

Score

  • CAI 52 → 77 (+24.6)
  • Rubric changed (rubric-2026.08.15 → rubric-2026.09.15) — scores are not directly comparable.

Lenses

  • Code Health 100 → 84 (-15.6)
  • Architecture 96 → 100 (+4.1)
  • Maturity 58 → 77 (+19.5)
  • Readiness 34 → 80 (+46.1)
  • Security 62 → 72 (+9.7)

Resolved (58)

  • Coverage not measured — test suite did not build
  • Dimension evaluation failed
  • High IaC: DS-0002 (benchmarks/Dockerfile)
  • High IaC: DS-0002 (benchmarks/tpc/infra/docker/Dockerfile.build-comet)
  • High: security finding (details withheld)
  • High: security finding (details withheld)
  • High: security finding (details withheld)
  • High: security finding (details withheld)
  • High: security finding (details withheld)
  • High: security finding (details withheld)
  • High: security finding (details withheld)
  • High: security finding (details withheld)
  • High: security finding (details withheld)
  • High: security finding (details withheld)
  • High: security finding (details withheld)
  • High: security finding (details withheld)
  • High: security finding (details withheld)
  • High: security finding (details withheld)
  • High: security finding (details withheld)
  • High: security finding (details withheld)
  • …and 38 more

New (781)

  • AvgDecimalGroupsAccumulator::merge_batch (cognitive 17) (native/spark-expr/src/agg_funcs/avg_decimal.rs)
  • AvgDecimalGroupsAccumulator::update_batch (cognitive 18) (native/spark-expr/src/agg_funcs/avg_decimal.rs)
  • AvgGroupsAccumulator::update_batch (cognitive 18) (native/spark-expr/src/agg_funcs/avg.rs)
  • Boundary-crossing change coupling: error.rs ↔ ShimSparkErrorConverter.scala (native/common/src/error.rs)
  • Boundary-crossing change coupling: error.rs ↔ ShimSparkErrorConverter.scala (native/common/src/error.rs)
  • Boundary-crossing change coupling: error.rs ↔ ShimSparkErrorConverter.scala (native/common/src/error.rs)
  • Boundary-crossing change coupling: utils.rs ↔ CometMetricNode.scala (native/core/src/execution/metrics/utils.rs)
  • CelebornRawPartitionReader.openPartitions (cognitive 37) (spark/src/main/scala/org/apache/spark/sql/comet/execution/shuffle/CometCelebornShuffleReader.scala)
  • CelebornRawPartitionReader.openPartitions (cyclomatic 23) (spark/src/main/scala/org/apache/spark/sql/comet/execution/shuffle/CometCelebornShuffleReader.scala)
  • CelebornShuffleGenerationCoordinator.claimMapAttempt (cognitive 27) (spark/src/main/scala/org/apache/spark/sql/comet/execution/shuffle/CometCelebornShuffleManager.scala)
  • CelebornShuffleGenerationCoordinator.claimMapAttempt (cyclomatic 16) (spark/src/main/scala/org/apache/spark/sql/comet/execution/shuffle/CometCelebornShuffleManager.scala)
  • CelebornShufflePartitionPusher.<init> (cognitive 36) (spark/src/main/java/org/apache/comet/shuffle/CelebornShufflePartitionPusher.java)
  • CelebornShufflePartitionPusher.<init> (cyclomatic 32) (spark/src/main/java/org/apache/comet/shuffle/CelebornShufflePartitionPusher.java)
  • CelebornShufflePartitionPusher.finish (cognitive 16) (spark/src/main/java/org/apache/comet/shuffle/CelebornShufflePartitionPusher.java)
  • CelebornShufflePartitionPusher.finish (cyclomatic 16) (spark/src/main/java/org/apache/comet/shuffle/CelebornShufflePartitionPusher.java)
  • CelebornShufflePartitionPusher.pushClaimedPartitionData (cognitive 41) (spark/src/main/java/org/apache/comet/shuffle/CelebornShufflePartitionPusher.java)
  • CelebornShufflePartitionPusher.pushClaimedPartitionData (cyclomatic 31) (spark/src/main/java/org/apache/comet/shuffle/CelebornShufflePartitionPusher.java)
  • CelebornShufflePartitionPusher.reconcileAcceptedPushes (cognitive 31) (spark/src/main/java/org/apache/comet/shuffle/CelebornShufflePartitionPusher.java)
  • CelebornShufflePartitionPusher.reconcileAcceptedPushes (cyclomatic 22) (spark/src/main/java/org/apache/comet/shuffle/CelebornShufflePartitionPusher.java)
  • CelebornTransportCallbackTracker.installFactoryHook (cognitive 16) (spark/src/main/java/org/apache/comet/shuffle/CelebornTransportCallbackTracker.java)
  • …and 761 more

Changes since last survey

  • 300 commits — 196 feature/other, 104 fixes

By area

  • spark/src — 135 commits
  • native/spark-expr — 43 commits
  • native/core — 36 commits
  • docs/source — 25 commits
  • .github/workflows — 22 commits
  • native/shuffle — 11 commits
  • dev/diffs — 9 commits
  • .ai/skills — 4 commits
  • dev/ci — 4 commits
  • (root) — 3 commits
  • native/common — 2 commits
  • (repo) — 1 commit
  • .github/actions — 1 commit
  • contrib/delta — 1 commit
  • contrib/lance — 1 commit
  • dev/verify-contrib-delta-gate.sh — 1 commit
  • native/Cargo.lock — 1 commit

Notable commits

  • fix: fix(celeborn): reject unsafe native push completion tracking (#5665)
  • fix: fix(iceberg): don't push transform residuals as their source column, fail on residual errors (#6154)
  • fix: fix(iceberg): guard native Iceberg scan driver-metric double-post, add metrics docs and tests (#6085)
  • fix: fix: Accept explicit positive years in timestamp casts (#5858)
  • fix: fix: Bump iceberg-rust so native Iceberg writes URL-escape partition paths (#5651)
  • fix: fix: Delete completed tasks' data files when an Iceberg write job fails (#5663)
  • fix: fix: Nested floating-point IN membership does not match Spark for signed zero (#6073)
  • fix: fix: accept UTC timezone aliases in Python Arrow input (#5556)
  • fix: fix: accept dictionary encodings in remote shuffle (#5650)
  • fix: fix: align Spark 4.2 Python worker configuration (#5561)
  • fix: fix: align string to timestamp parsing with Spark's segment rules (#5682)
  • fix: fix: align time parsing and native second extraction with Spark (#5738)
  • fix: fix: apply Spark's Parquet conversion rules to nested struct/list/map fields (#5681)
  • fix: fix: apply the parent struct's null mask before hashing its fields (#5754)
  • fix: fix: attach tokio runtime threads to the JVM as daemon threads (#5748)
  • fix: fix: bound shuffle schema cache retention and preserve eviction order (#6098)
  • fix: fix: check each fair_unified reservation against its own share (#6205)
  • fix: fix: check nested TIMESTAMP_MILLIS overflow in unfiltered scans (#5740)
  • fix: fix: correct two nightly test failures on Spark 3.4 and 4.2 (#6156)
  • fix: fix: correctly rounded decimal to double/float cast matching BigDecimal.doubleValue/floatValue (#5684)
  • …and 280 more

Written by watchdog.canine.dev from the codebase's own history, inside the signed delivery this page is composed from.

Survey your own repository

apache/datafusion-comet was measured the same way every project in this corpus was: the same rubric, at a pinned commit, with the result published in full. Point a surveyor at a repository you know and see whether you agree with it.

About this page

  • The score is its most recent published measurement, taken on 27 September 2026 at a pinned commit. It is not a live figure and does not change until the project is measured again.
  • Measured at commit 605051ad239ef704f5f25d67910a446a6b6d7c70 — the exact code this score is about.
  • Scored under rubric-2026.09.15 — the same rubric and the same method as every other entry in this index.
  • Measured by watchdog.canine.dev using codehealth-analyzer preprod-d00c643c3f66.