From 47dea8e350adf711c220f18f31825ce3a77e2e0c Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Mon, 28 Sep 2026 07:53:13 -0600 Subject: [PATCH] docs: refresh stale roadmap entries Update roadmap sections whose status changed since the last pass: - Iceberg writes landed behind two experimental flags (#4658, #5361); point at the production-quality epic #5649. - Native Iceberg deletion vector reads landed (#5853); list the V3 features that still fall back, and drop the stale iceberg_scan.rs scheme-match reference. - The codegen dispatcher covers only scalar expressions, not aggregates; link the expression reference, which lists dispatched expressions. - mapInArrow and mapInPandas have experimental support (#4234). - Replace the closed TPC-DS epic #858 with #2551. - datafusion-spark function wiring (#4150) is nearly complete. - Link the hash join spill issue and the upstream DataFusion design. - Describe the delta-spark contrib scan (#5365) and the convergence proposal (#5411). --- docs/source/contributor-guide/roadmap.md | 88 ++++++++++++++---------- 1 file changed, 51 insertions(+), 37 deletions(-) diff --git a/docs/source/contributor-guide/roadmap.md b/docs/source/contributor-guide/roadmap.md index b5aa6be16fd..d318c07a539 100644 --- a/docs/source/contributor-guide/roadmap.md +++ b/docs/source/contributor-guide/roadmap.md @@ -53,30 +53,29 @@ these expressions benefit from vectorized native execution. ## Native Coverage for Codegen-Dispatched Expressions Beyond lambda bodies, a number of built-in Spark scalar expressions (some regular expression, JSON, and datetime -functions, for example) and aggregate functions route through the same JVM codegen-dispatch bridge by default, -either because their native DataFusion or `datafusion-spark` implementation has known semantic differences from -Spark, or because no native implementation exists yet. See the [compatibility guide] for the current list of -codegen-dispatched expressions. We're exploring closing these gaps so that more expressions run natively by -default, which would reduce JVM round-trips beyond what the lambda and UDF work above already covers. +functions, for example) route through the same JVM codegen-dispatch bridge by default, either because their native +DataFusion or `datafusion-spark` implementation has known semantic differences from Spark, or because no native +implementation exists yet. The [expression reference] records which expressions are codegen-dispatched today. The +bridge covers only scalar expressions, so an aggregate function with no Spark-compatible native implementation falls +back to Spark instead. We're exploring closing these gaps so that more expressions run natively by default, which +would remove JVM round-trips beyond those the lambda work above addresses. -[compatibility guide]: ../user-guide/latest/compatibility/index.md +[expression reference]: ../user-guide/latest/expressions.md ## Iceberg Table Format V3 Support -Comet landed its first Iceberg table format V3 feature (native data file decryption, [#4991]), and we want to -add more V3 features to Comet's native Iceberg scans so they don't fall back to Spark. The work is tracked -phase-by-phase in [#3376]: detecting the V3 format and supporting new V3 data types, deletion vector reads, row -lineage, and table encryption. Native deletion vector reads are prototyped in a draft PR ([#4887]), but are -blocked on upstream `iceberg-rust` support, tracked in [iceberg-rust #2792] and [iceberg-rust #2411]. Native Iceberg -scans don't support HDFS-backed tables today: the storage scheme match in `iceberg_scan.rs` only handles `file`, -`s3`/`s3a`, `gs`, and `oss`, and would need an HDFS `StorageFactory` upstream in `iceberg-storage-opendal`. We're -scoping what that work would take. +Comet's native Iceberg scans read V3 tables, including encrypted tables ([#4991]) and tables with deletion vectors +([#5853]). We want to add the remaining V3 features so that these scans don't fall back to Spark: row lineage +metadata columns, column default values, and the new V3 types (`variant`, `geometry`, `geography`, and `unknown`). +The work is tracked in [#3376], and upstream `iceberg-rust` support in [iceberg-rust #2411]. Native Iceberg scans +also don't support HDFS-backed tables today: Comet's native Iceberg storage layer handles only local files, S3 and +S3-compatible stores, GCS, and OSS, and would need an HDFS `StorageFactory` upstream in `iceberg-storage-opendal`. +We're scoping what that work would take. [#3376]: https://github.com/apache/datafusion-comet/issues/3376 -[#4887]: https://github.com/apache/datafusion-comet/pull/4887 [#4991]: https://github.com/apache/datafusion-comet/pull/4991 +[#5853]: https://github.com/apache/datafusion-comet/pull/5853 [iceberg-rust #2411]: https://github.com/apache/iceberg-rust/issues/2411 -[iceberg-rust #2792]: https://github.com/apache/iceberg-rust/issues/2792 ## TPC-H and TPC-DS Performance @@ -84,18 +83,18 @@ Comet already delivers substantial speedups over vanilla Spark on TPC-H and TPC- results for [TPC-DS] with each release. An independent [AWS Labs benchmark] comparing Comet 0.16.0 with Gluten 1.6.0 on a 3TB TPC-DS workload found that the two accelerators deliver similar overall performance. Increasing the speedup further and closing the remaining per-query gaps is an ongoing focus, tracked under [#2004] (TPC-H) and -[#858] (TPC-DS). +[#2551] (TPC-DS). [TPC-DS]: benchmark-results/tpc-ds.md [AWS Labs benchmark]: https://awslabs.github.io/data-on-eks/docs/benchmarks/spark-gluten-velox-comet-benchmark -[#858]: https://github.com/apache/datafusion-comet/issues/858 [#2004]: https://github.com/apache/datafusion-comet/issues/2004 +[#2551]: https://github.com/apache/datafusion-comet/issues/2551 ## Upstream Work in DataFusion A growing number of Spark-compatible expressions live in the `datafusion-spark` crate in the core DataFusion repository. Comet is migrating its expression implementations to that crate so that they can be shared by other -DataFusion-based projects, and is wiring up the functions that crate already provides ([#4150]). Improvements to +DataFusion-based projects, and has wired up nearly every function that crate provides ([#4150]). Improvements to core DataFusion operators (joins, aggregates, window) made in support of Comet also benefit the wider ecosystem. [#4150]: https://github.com/apache/datafusion-comet/issues/4150 @@ -103,17 +102,25 @@ core DataFusion operators (joins, aggregates, window) made in support of Comet a ## Spillable Hash Join Comet's native hash join currently requires the build side to fit entirely in memory. Adding spill-to-disk -support will allow Comet to handle larger joins without falling back to Spark, improving both reliability and -performance for memory-intensive workloads. +support ([#2545]) will allow Comet to handle larger joins without falling back to Spark, improving both reliability +and performance for memory-intensive workloads. Comet's native hash join uses DataFusion's `HashJoinExec`, and a +design for spilling in that operator, behind a flag that is off by default, is proposed upstream in +[datafusion #24768]. + +[#2545]: https://github.com/apache/datafusion-comet/issues/2545 +[datafusion #24768]: https://github.com/apache/datafusion/issues/24768 ## Java/Scala UDF Support Spark users frequently define custom UDFs in Java or Scala. Comet now dispatches scalar `ScalaUDF` expressions through a JVM codegen bridge (`CometScalaUDF`) instead of always falling back to Spark. Aggregate UDFs, table -UDFs/generators, Python/Pandas UDFs, and Hive `GenericUDF`/`SimpleUDF` still fall back to Spark entirely. +UDFs/generators, Hive `GenericUDF`/`SimpleUDF`, and Python UDFs other than `mapInArrow` and `mapInPandas` still fall +back to Spark entirely; Comet's support for those two Python APIs is experimental and disabled by default ([#4234]). Extending the codegen-dispatch approach to cover these remaining categories will reduce fallbacks further and allow more queries to run end-to-end in Comet. +[#4234]: https://github.com/apache/datafusion-comet/pull/4234 + ## Memory Management Improvements Comet coordinates memory between the JVM and native Rust execution through a custom memory pool. Improving @@ -130,30 +137,37 @@ enabled by default ([#1625]). ## Iceberg Table Writes -Comet's native scans accelerate Iceberg reads today, but writes still run entirely through Spark's Iceberg V2 write -operator. [#4322] tracks exploring native acceleration of Iceberg writes, so an ETL job could run end to end in -native code, falling back only where Comet can't match Iceberg Java's functionality. Draft proposals split the -existing V2 write operator into separate "writer" and "committer" operators so the data-file-write portion can be -planned and accelerated the same way as a native Parquet write ([#4658], [#4487]). No design here is committed yet; -we're exploring whether this approach is worth pursuing. +Comet can now write Iceberg tables natively. The feature is experimental, disabled by default, and controlled by two +settings, the second of which requires the first. The first splits Spark's Iceberg V2 write operator into separate +writer and committer operators, so the query feeding the write becomes visible to AQE and to Comet's columnar rules +([#4658]). The second delegates each task's Parquet write to `iceberg-rust` when the write passes an eligibility +check ([#5361]); writes that don't pass keep Iceberg Java's writer. Merge-on-read writes are not intercepted yet +([#6240]). The remaining work toward the original goal of [#4322], an ETL job that runs end to end in native code, is +tracked in [#5649]: correctness fixes, failure handling that matches Iceberg Java, broader coverage, and enabling +both settings by default ([#5644]). See [Iceberg Writes](iceberg-writes.md) for how the write path works. [#4322]: https://github.com/apache/datafusion-comet/issues/4322 -[#4487]: https://github.com/apache/datafusion-comet/pull/4487 [#4658]: https://github.com/apache/datafusion-comet/pull/4658 +[#5361]: https://github.com/apache/datafusion-comet/pull/5361 +[#5644]: https://github.com/apache/datafusion-comet/issues/5644 +[#5649]: https://github.com/apache/datafusion-comet/issues/5649 +[#6240]: https://github.com/apache/datafusion-comet/issues/6240 ## Delta Lake Support -Comet currently supports Spark's built-in file formats and Iceberg, but not Delta Lake ([#174]). A draft PR reuses -Comet's existing native Parquet reader to accelerate scans of plain Delta tables that use neither deletion vectors -nor column mapping ([#4669]); a separate, larger effort explores a full native Delta scan built on `delta-kernel-rs` -as an out-of-tree contrib module ([#4366]), landing in inert, gated slices ([#4700], [#4952]). Generalizing Comet's -scan-side APIs so that Delta and other non-Iceberg data sources can plug in more easily is tracked as a Table -Provider API abstraction ([#4706]). None of this is committed yet; we're exploring which approach, if any, is worth -pursuing. +Comet currently supports Spark's built-in file formats and Iceberg, but not Delta Lake ([#174]). The plugin boundary +for an out-of-tree Delta module has landed in inert, gated slices ([#4700], [#4952]), and two read paths build on +it. In the first, which is in review, delta-spark plans the scan and Comet's native Parquet reader executes it, +applying deletion vectors inside the scan; it ships as an opt-in contrib module ([#5365]). The second, still a +draft, is a full native Delta scan built on `delta-kernel-rs` ([#4366]). [#5411] proposes converging the two into +one plugin, with the Parquet path as the default and the kernel path covering what it can't express yet, such as +change data feed. Generalizing Comet's scan-side APIs so that Delta and other non-Iceberg data sources can plug in +more easily is tracked as a Table Provider API abstraction ([#4706]). [#174]: https://github.com/apache/datafusion-comet/issues/174 [#4366]: https://github.com/apache/datafusion-comet/pull/4366 -[#4669]: https://github.com/apache/datafusion-comet/pull/4669 [#4700]: https://github.com/apache/datafusion-comet/pull/4700 [#4706]: https://github.com/apache/datafusion-comet/issues/4706 [#4952]: https://github.com/apache/datafusion-comet/pull/4952 +[#5365]: https://github.com/apache/datafusion-comet/pull/5365 +[#5411]: https://github.com/apache/datafusion-comet/issues/5411