viirya opened a new issue, #6499: URL: https://github.com/apache/datafusion-comet/issues/6499
### What is the problem the feature request solves? When Spark runs with a local master, the driver and all executors share one JVM, but Comet still executes a query the distributed way: each Spark task builds and runs its own native plan, stages exchange data through Spark shuffle files, and each result partition is encoded by its task and decoded on the driver. In a single process these boundaries are pure overhead. Partial aggregates, join inputs and sort runs are written to and read back from shuffle files that DataFusion could exchange in memory, and the native plans of one stage cannot share operator state. DataFusion can already execute a whole physical plan, including hash repartitioning, on its own Tokio runtime. For queries that Comet fully supports, the single-process case could hand the whole query to DataFusion and keep Spark only for parsing, analysis, optimization and result delivery. ### Describe the potential solution Add an experimental, opt-in local execution mode (an internal `spark.comet.exec.local.enabled` option, off by default): - Admit a query only as a whole, before ordinary Comet conversion: Spark 4.1, an in-process local master, AQE disabled, batch queries without subqueries, and a fully supported plan. Any other query keeps the existing Comet/Spark path, and there is no fallback once native output has started. - Plan the whole query as one DataFusion graph in the driver, reusing Comet's native Parquet scan and Spark-compatible expressions. Replace Spark exchanges with DataFusion repartitioning inside that graph, so the query has no Spark shuffle and runs as a single Spark result task. Each execution owns a fresh graph and a query-scoped memory pool with spilling. - Start with Parquet scan/filter/project, grouped and global `COUNT`/`MIN`/`MAX`, shuffled hash joins, and terminal global sort, Top-K and limit, and leave everything else to the existing path. - Keep Spark's cancellation, job groups and SQL metrics by still running the result task as a Spark job, while delivering `collect`/`take` results to the driver within the same JVM instead of encoding them in the task. ### Additional context Related to #1204, which avoids repeating native plan construction across tasks within an executor. This mode instead targets the single-process case, where the whole query can be one native graph. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
