viirya commented on issue #6499: URL: https://github.com/apache/datafusion-comet/issues/6499#issuecomment-5975454721
Thanks @andygrove for the careful review on #6500. Replying here since the main questions are about direction. **Reusing the existing operator translation.** Agreed. `local.proto` and the local planner re-describe aggregate, join, sort and limit and lower them a second time, so every newly admitted operator would need a second implementation without the Spark-compatibility handling that `PhysicalPlanner::create_plan` already has. I'd rather build the graph from the `Operator` trees that `CometExecRule` produces: - Admit a plan only when it is entirely native: native blocks connected by `CometShuffleExchangeExec`, with no JVM input. - Lower each block with `PhysicalPlanner::create_plan`, plan all file partitions of a native scan in the one graph, and replace each shuffle exchange with a `RepartitionExec`, or for range partitioning with a per-partition sort and a sort-preserving merge. - Admission then follows whatever Comet runs natively, and `local.proto` goes away. One thing I'd like to confirm: within one graph both sides of a join go through the same `RepartitionExec`, so I don't think the local exchange needs Comet's Spark-compatible Murmur3 partitioning, unless some native operator depends on it. **AQE.** Admission from the `CometRule(session, queryStagePrep = true)` instance looks like the right way to drop the session-wide requirement. Since `InsertAdaptiveSparkPlan` leaves exchange-free plans alone, those can be admitted with AQE enabled from the columnar rule, and plans with exchanges from the query stage prep rule. I'll prototype the latter to check that replacing the initial plan with a single leaf leaves AQE nothing to re-plan. **Splitting.** Happy to split it along the lines you suggested: 1. The `native/local` lifecycle, JNI bridge, result delivery and `range`. 2. Admission of exchange-free native plans (scan/filter/project) from `CometExecRule` operator trees, without requiring AQE to be disabled. 3. Exchanges through `RepartitionExec` for aggregates and joins, admitted from the query stage prep rule, plus Top-K and limit. 4. Global sort after the DataFusion 56 upgrade (#6410); see the review thread on #6500. The other review points (spill directories and size limit, native partition count, unordered limit determinism, cached plans, the Spark version gate and `memory_management.md`) would be addressed in the PR that introduces the affected code. I'll convert #6500 to a draft and keep it as a reference until the first split PR is up. -- 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]
