From 9601e81ba8fa2e8b1e7b3d92fbdc976e91f7a380 Mon Sep 17 00:00:00 2001 From: rich7420 Date: Mon, 21 Sep 2026 17:56:22 +0800 Subject: [PATCH] feat: support map inputs for explode --- .../expression-audits/generator_funcs.md | 8 +- docs/source/user-guide/latest/expressions.md | 3 +- docs/source/user-guide/latest/operators.md | 12 +- .../src/execution/expressions/map_entries.rs | 119 ++++++++++++++++ native/core/src/execution/expressions/mod.rs | 1 + native/core/src/execution/planner.rs | 131 ++++++++++++++++-- native/proto/src/proto/operator.proto | 4 +- .../apache/spark/sql/comet/operators.scala | 6 +- .../sql-tests/expressions/array/explode.sql | 4 +- .../expressions/array/posexplode.sql | 4 +- .../sql-tests/expressions/map/explode.sql | 93 +++++++++++++ .../comet/exec/CometGenerateExecSuite.scala | 41 ++++-- .../sql/benchmark/CometExplodeBenchmark.scala | 11 +- 13 files changed, 381 insertions(+), 56 deletions(-) create mode 100644 native/core/src/execution/expressions/map_entries.rs create mode 100644 spark/src/test/resources/sql-tests/expressions/map/explode.sql diff --git a/docs/source/contributor-guide/expression-audits/generator_funcs.md b/docs/source/contributor-guide/expression-audits/generator_funcs.md index 3023cfaa774..d34a0b61071 100644 --- a/docs/source/contributor-guide/expression-audits/generator_funcs.md +++ b/docs/source/contributor-guide/expression-audits/generator_funcs.md @@ -23,18 +23,18 @@ ## explode -- Handled at the operator level as a `GenerateExec` (`CometExplodeExec`), not via the expression serde maps, so it is not auto-detected by the function-registry checkbox logic. Compatible for array inputs; map inputs fall back ([#2837](https://github.com/apache/datafusion-comet/issues/2837)). +- Handled at the operator level as a `GenerateExec` (`CometExplodeExec`), not via the expression serde maps, so it is not auto-detected by the function-registry checkbox logic. Compatible for array and map inputs. Maps emit key/value columns, with a position column for `posexplode`. ## explode_outer -- Same `CometExplodeExec` path as `explode`. Compatible for array inputs; empty and NULL arrays both emit one null-valued row per Spark's `outer` semantics, which the planner requests as DataFusion's `NullHandling::PreserveAndExpandEmpty`. Map inputs fall back. +- Same `CometExplodeExec` path as `explode`. Compatible for array and map inputs; empty and NULL collections both emit one null-valued row per Spark's `outer` semantics, which the planner requests as DataFusion's `NullHandling::PreserveAndExpandEmpty`. Map entries reuse the list unnest path and expand into key/value columns. ## posexplode -- Handled at the operator level as a `GenerateExec` (`CometExplodeExec`), like `explode`. Compatible for array inputs; map inputs fall back ([#2837](https://github.com/apache/datafusion-comet/issues/2837)). +- Handled at the operator level as a `GenerateExec` (`CometExplodeExec`), like `explode`. Compatible for array and map inputs. Maps emit key/value columns, with a position column for `posexplode`. ## posexplode_outer -- Same `CometExplodeExec` path as `posexplode`. Compatible for array inputs; empty and NULL arrays both emit one row with null `pos` and null `value` per Spark's `outer` semantics, which the planner requests as DataFusion's `NullHandling::PreserveAndExpandEmpty`. +- Same `CometExplodeExec` path as `posexplode`. Compatible for array and map inputs; empty and NULL collections both emit one row with null `pos` and null generated columns per Spark's `outer` semantics, which the planner requests as DataFusion's `NullHandling::PreserveAndExpandEmpty`. [Spark Expression Support]: ../../user-guide/latest/expressions.md diff --git a/docs/source/user-guide/latest/expressions.md b/docs/source/user-guide/latest/expressions.md index cbc0b6e6406..efda85c89db 100644 --- a/docs/source/user-guide/latest/expressions.md +++ b/docs/source/user-guide/latest/expressions.md @@ -332,8 +332,7 @@ The type-name conversion functions (`bigint`, `binary`, `boolean`, `date`, `deci ## generator_funcs `explode`, `explode_outer`, `posexplode`, and `posexplode_outer` are supported via -`CometExplodeExec` (operator-level, not expression-level) for array input; map input falls back -to Spark ([#2837](https://github.com/apache/datafusion-comet/issues/2837)). Enabled by default via +`CometExplodeExec` for array and map inputs (operator-level, not expression-level). Enabled by default via `spark.comet.exec.explode.enabled`. | Function | Status | Implementation | Notes | diff --git a/docs/source/user-guide/latest/operators.md b/docs/source/user-guide/latest/operators.md index c9ba5e491e1..0acb2b616ac 100644 --- a/docs/source/user-guide/latest/operators.md +++ b/docs/source/user-guide/latest/operators.md @@ -108,12 +108,12 @@ omitted from the tables below and may be reconsidered based on demand: ## Generators and set operations -| Operator | Status | Notes | -| -------------- | ------ | ---------------------------------------------------------------------------------------------------------------- | -| `GenerateExec` | ✅ | Supports `explode`, `explode_outer`, `posexplode`, `posexplode_outer` over arrays. `inline` / `stack` fall back. | -| `ExpandExec` | ✅ | | -| `UnionExec` | ✅ | | -| `CoalesceExec` | ✅ | | +| Operator | Status | Notes | +| -------------- | ------ | ------------------------------------------------------------------------------------------------------------------------- | +| `GenerateExec` | ✅ | Supports `explode`, `explode_outer`, `posexplode`, `posexplode_outer` over arrays and maps. `inline` / `stack` fall back. | +| `ExpandExec` | ✅ | | +| `UnionExec` | ✅ | | +| `CoalesceExec` | ✅ | | ## Writes diff --git a/native/core/src/execution/expressions/map_entries.rs b/native/core/src/execution/expressions/map_entries.rs new file mode 100644 index 00000000000..8188a80730d --- /dev/null +++ b/native/core/src/execution/expressions/map_entries.rs @@ -0,0 +1,119 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +use std::fmt::{Display, Formatter}; +use std::hash::{Hash, Hasher}; +use std::sync::Arc; + +use arrow::array::{Array, ListArray, MapArray, RecordBatch}; +use arrow::datatypes::{DataType, FieldRef, Schema}; +use datafusion::common::{exec_err, Result as DataFusionResult}; +use datafusion::physical_expr::PhysicalExpr; +use datafusion::physical_plan::ColumnarValue; + +/// Exposes a map's entries as a list for `ExplodeExec`, sharing the input buffers. +/// Preserve the original entry fields, including nullability and metadata; the +/// SQL `map_entries` function rebuilds those fields and loses their metadata. +#[derive(Debug, Clone)] +pub struct MapEntriesExpr { + child: Arc, +} + +impl MapEntriesExpr { + pub fn new(child: Arc) -> Self { + Self { child } + } +} + +impl Display for MapEntriesExpr { + fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { + write!(f, "map_entries({})", self.child) + } +} + +impl PartialEq for MapEntriesExpr { + fn eq(&self, other: &Self) -> bool { + self.child.eq(&other.child) + } +} + +impl Eq for MapEntriesExpr {} + +impl Hash for MapEntriesExpr { + fn hash(&self, state: &mut H) { + self.child.hash(state); + } +} + +impl PhysicalExpr for MapEntriesExpr { + fn fmt_sql(&self, f: &mut Formatter<'_>) -> std::fmt::Result { + Display::fmt(self, f) + } + + fn return_field(&self, input_schema: &Schema) -> DataFusionResult { + let field = self.child.return_field(input_schema)?; + let DataType::Map(entries, _) = field.data_type() else { + return exec_err!( + "MapEntriesExpr expected Map input, got {}", + field.data_type() + ); + }; + Ok(Arc::new( + field + .as_ref() + .clone() + .with_data_type(DataType::List(Arc::clone(entries))), + )) + } + + fn evaluate(&self, batch: &RecordBatch) -> DataFusionResult { + let array = self.child.evaluate(batch)?.into_array(batch.num_rows())?; + let Some(map) = array.as_any().downcast_ref::() else { + return exec_err!( + "MapEntriesExpr expected Map input, got {}", + array.data_type() + ); + }; + let DataType::Map(entries, _) = map.data_type() else { + unreachable!("MapArray downcast guarantees DataType::Map"); + }; + let list = ListArray::try_new( + Arc::clone(entries), + map.offsets().clone(), + Arc::new(map.entries().clone()), + map.nulls().cloned(), + )?; + Ok(ColumnarValue::Array(Arc::new(list))) + } + + fn children(&self) -> Vec<&Arc> { + vec![&self.child] + } + + fn with_new_children( + self: Arc, + children: Vec>, + ) -> DataFusionResult> { + if children.len() != 1 { + return exec_err!( + "MapEntriesExpr expects exactly 1 child, got {}", + children.len() + ); + } + Ok(Arc::new(Self::new(Arc::clone(&children[0])))) + } +} diff --git a/native/core/src/execution/expressions/mod.rs b/native/core/src/execution/expressions/mod.rs index e174bd37475..b1da89aacc2 100644 --- a/native/core/src/execution/expressions/mod.rs +++ b/native/core/src/execution/expressions/mod.rs @@ -22,6 +22,7 @@ pub mod bitwise; pub mod comparison; pub mod list_positions; pub mod logical; +pub mod map_entries; pub mod nullcheck; pub mod partition; pub mod random; diff --git a/native/core/src/execution/planner.rs b/native/core/src/execution/planner.rs index 80f8d890a34..20b383fe2e5 100644 --- a/native/core/src/execution/planner.rs +++ b/native/core/src/execution/planner.rs @@ -40,6 +40,7 @@ use crate::execution::operators::IcebergWriteExec; use crate::execution::operators::{PartitionedRankLimitExec, WindowFnKind}; use crate::execution::{ expressions::list_positions::ListPositionsExpr, + expressions::map_entries::MapEntriesExpr, expressions::subquery::Subquery, operators::{ CometFilterExec, ExecutionError, ExpandExec, ExplodeExec, ParquetCompression, @@ -2076,7 +2077,7 @@ impl PhysicalPlanner { let (scans, shuffle_scans, child) = self.create_plan(&children[0], inputs, partition_count)?; - // Create the expression for the array to explode + // Create the expression for the collection to explode let child_expr = if let Some(child_expr) = &explode.child { self.create_expr(child_expr, child.schema())? } else { @@ -2092,6 +2093,22 @@ impl PhysicalPlanner { .name() .to_string(); + // Expose maps as List>, sharing their entries buffers. + // Reuse the list path for outer rows, positions and bounded output batches, + // then flatten the entry struct into Spark's two output columns. + let map_fields = match child_expr.data_type(&child_schema)? { + DataType::Map(entries, _) => match entries.data_type() { + DataType::Struct(fields) => Some(fields.clone()), + _ => unreachable!("Map entries must be a struct"), + }, + _ => None, + }; + let child_expr: Arc = if map_fields.is_some() { + Arc::new(MapEntriesExpr::new(child_expr)) + } else { + child_expr + }; + // Both posexplode variants reference the array twice: once for positions // and once for values. Materialize computed arrays so both references // share one evaluation. A plain Column is already materialized. @@ -2186,11 +2203,22 @@ impl PhysicalPlanner { } }; - output_fields.push(Field::new( - array_field.name(), - element_type, - true, // Element is nullable after unnesting - )); + let struct_unnests = if let Some(fields) = map_fields { + output_fields.extend(fields.iter().map(|field| { + field + .as_ref() + .clone() + .with_nullable(explode.outer || field.is_nullable()) + })); + vec![array_input_index] + } else { + output_fields.push(Field::new( + array_field.name(), + element_type, + true, // Element is nullable after unnesting + )); + vec![] + }; let output_schema = Arc::new(Schema::new(output_fields)); @@ -2220,7 +2248,7 @@ impl PhysicalPlanner { let unnest_exec = Arc::new(ExplodeExec::new( project_exec, list_unnests, - vec![], // No struct columns to unnest + struct_unnests, output_schema, unnest_options, )?); @@ -6342,11 +6370,23 @@ mod tests { #[tokio::test] async fn explode_evaluates_array_once_per_batch() { + check_explode_evaluates_collection_once_per_batch(false).await; + } + + #[tokio::test] + async fn explode_evaluates_map_once_per_batch() { + check_explode_evaluates_collection_once_per_batch(true).await; + } + + async fn check_explode_evaluates_collection_once_per_batch(map: bool) { + use arrow::array::{AsArray, MapArray, StructArray}; use arrow::datatypes::Int32Type; use datafusion::common::tree_node::{Transformed, TreeNode}; use datafusion::logical_expr::{create_udf, Volatility}; use datafusion::physical_plan::projection::ProjectionExec; - use spark_expression::data_type::{data_type_info::DatatypeStruct, DataTypeInfo, ListInfo}; + use spark_expression::data_type::{ + data_type_info::DatatypeStruct, DataTypeInfo, ListInfo, MapInfo, + }; let array_type = spark_expression::DataType { type_id: 14, @@ -6365,6 +6405,55 @@ mod tests { Some(vec![Some(20)]), ])) as ArrayRef; + let (array_type, arrays) = if map { + // Slice away a leading entry to exercise non-zero map offsets. + let values = ListArray::from_iter_primitive::(vec![ + Some(vec![Some(99)]), + Some(vec![Some(10), None]), + Some(vec![]), + None, + Some(vec![Some(20)]), + ]) + .slice(1, 4); + let fields: Fields = vec![ + Field::new("key", DataType::Int32, false) + .with_metadata([("PARQUET:field_id".to_string(), "10".to_string())].into()), + Field::new("value", DataType::Int32, true) + .with_metadata([("PARQUET:field_id".to_string(), "11".to_string())].into()), + ] + .into(); + let entries = StructArray::new( + fields.clone(), + vec![ + Arc::new(Int32Array::from(vec![99, 1, 2, 3])), + Arc::clone(values.values()), + ], + None, + ); + let maps = MapArray::new( + Arc::new(Field::new("entries", DataType::Struct(fields), false)), + values.offsets().clone(), + entries, + values.nulls().cloned(), + false, + ); + let map_type = spark_expression::DataType { + type_id: 15, + type_info: Some(Box::new(DataTypeInfo { + datatype_struct: Some(DatatypeStruct::Map(Box::new(MapInfo { + key_type: Some(Box::new(create_proto_datatype())), + value_type: Some(Box::new(create_proto_datatype())), + value_contains_null: true, + key_field_id: Some(10), + value_field_id: Some(11), + }))), + })), + }; + (map_type, Arc::new(maps) as ArrayRef) + } else { + (array_type, arrays) + }; + for outer in [false, true] { for position in [false, true] { for computed in [false, true] { @@ -6443,20 +6532,21 @@ mod tests { let results = collect(native_plan.execute(0, task_ctx).unwrap()) .await .unwrap(); - let context = - format!("outer={outer}, position={position}, computed={computed}"); + let context = format!( + "outer={outer}, position={position}, computed={computed}, map={map}" + ); assert_eq!( calls.load(Ordering::Relaxed), if computed { 2 } else { 0 }, "{context}" ); // The array is pre-projected only to share one evaluation between the - // `pos` and `value` references, so only a computed child needs it. `outer` + // `pos` and `value` references, so a computed child or map wrapper needs it. `outer` // does not, since it is now a `UnnestOptions` mode rather than a wrapper // expression around the child. assert_eq!( projections, - 1 + usize::from(position && computed), + 1 + usize::from(position && (computed || map)), "{context}" ); let expected_values = if outer { @@ -6476,6 +6566,23 @@ mod tests { }) .collect(); assert_eq!(values, expected_values.repeat(2), "{context}"); + if map { + let expected_keys = if outer { + vec![Some(1), Some(2), None, None, Some(3)] + } else { + vec![Some(1), Some(2), Some(3)] + }; + let keys: Vec<_> = results + .iter() + .flat_map(|batch| { + batch + .column(batch.num_columns() - 2) + .as_primitive::() + .iter() + }) + .collect(); + assert_eq!(keys, expected_keys.repeat(2), "{context}"); + } if position { let expected_positions = if outer { vec![Some(0), Some(1), None, None, Some(0)] diff --git a/native/proto/src/proto/operator.proto b/native/proto/src/proto/operator.proto index 7b3a7aec0d2..2619cf0be71 100644 --- a/native/proto/src/proto/operator.proto +++ b/native/proto/src/proto/operator.proto @@ -945,9 +945,9 @@ message Expand { } message Explode { - // The array expression to explode into multiple rows + // The array or map expression to explode into multiple rows spark.spark_expression.Expr child = 1; - // Whether this is explode_outer (produces null row for empty/null arrays) + // Whether this is explode_outer (produces null row for empty/null collections) bool outer = 2; // Expressions for other columns to project alongside the exploded values repeated spark.spark_expression.Expr project_list = 3; diff --git a/spark/src/main/scala/org/apache/spark/sql/comet/operators.scala b/spark/src/main/scala/org/apache/spark/sql/comet/operators.scala index 952b53def8e..a465cebba16 100644 --- a/spark/src/main/scala/org/apache/spark/sql/comet/operators.scala +++ b/spark/src/main/scala/org/apache/spark/sql/comet/operators.scala @@ -1711,12 +1711,8 @@ object CometExplodeExec extends CometOperatorSerde[GenerateExec] { return Unsupported(Some(s"Unsupported generator: ${op.generator.nodeName}")) } op.generator.children.head.dataType match { - case _: ArrayType => + case _: ArrayType | _: MapType => Compatible() - case _: MapType => - // TODO add support for map types - // https://github.com/apache/datafusion-comet/issues/2837 - Unsupported(Some("Comet only supports explode/explode_outer for arrays, not maps")) case other => Unsupported(Some(s"Unsupported data type: $other")) } diff --git a/spark/src/test/resources/sql-tests/expressions/array/explode.sql b/spark/src/test/resources/sql-tests/expressions/array/explode.sql index b1bdb8b7493..881186e264e 100644 --- a/spark/src/test/resources/sql-tests/expressions/array/explode.sql +++ b/spark/src/test/resources/sql-tests/expressions/array/explode.sql @@ -328,7 +328,7 @@ CREATE TABLE test_explode_empty(id int, arr array) USING parquet query SELECT id, explode_outer(arr) FROM test_explode_empty --- ===== Map falls back to Spark: outer form must also fall back ===== +-- ===== Map inputs: outer form preserves empty and null maps ===== statement CREATE TABLE test_explode_map(id int, m map) USING parquet @@ -339,7 +339,7 @@ INSERT INTO test_explode_map VALUES (2, map()), (3, NULL) -query expect_fallback(Comet only supports explode/explode_outer for arrays, not maps) +query SELECT id, explode_outer(m) FROM test_explode_map -- ===== TIMESTAMP_NTZ (Spark 3.4+): distinct Arrow layout (no tz) from LTZ ===== diff --git a/spark/src/test/resources/sql-tests/expressions/array/posexplode.sql b/spark/src/test/resources/sql-tests/expressions/array/posexplode.sql index eb537586f0e..d31c00619c4 100644 --- a/spark/src/test/resources/sql-tests/expressions/array/posexplode.sql +++ b/spark/src/test/resources/sql-tests/expressions/array/posexplode.sql @@ -93,8 +93,8 @@ INSERT INTO test_posexplode_map VALUES (1, map('a', 1, 'b', 2)), (2, map('c', 3)) --- posexplode over a map falls back to Spark (Comet only supports array inputs, not maps) -query expect_fallback(Comet only supports explode/explode_outer for arrays, not maps) +-- Map entries expand to position, key and value columns. +query SELECT id, posexplode(m) FROM test_posexplode_map -- ===== posexplode_outer across non-int element types ===== diff --git a/spark/src/test/resources/sql-tests/expressions/map/explode.sql b/spark/src/test/resources/sql-tests/expressions/map/explode.sql new file mode 100644 index 00000000000..5e21151cfd2 --- /dev/null +++ b/spark/src/test/resources/sql-tests/expressions/map/explode.sql @@ -0,0 +1,93 @@ +-- Licensed to the Apache Software Foundation (ASF) under one +-- or more contributor license agreements. See the NOTICE file +-- distributed with this work for additional information +-- regarding copyright ownership. The ASF licenses this file +-- to you under the Apache License, Version 2.0 (the +-- "License"); you may not use this file except in compliance +-- with the License. You may obtain a copy of the License at +-- +-- http://www.apache.org/licenses/LICENSE-2.0 +-- +-- Unless required by applicable law or agreed to in writing, +-- software distributed under the License is distributed on an +-- "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +-- KIND, either express or implied. See the License for the +-- specific language governing permissions and limitations +-- under the License. + +-- Config: spark.comet.exec.explode.enabled=true + +statement +CREATE TABLE explode_maps(id int, m map) USING parquet + +statement +INSERT INTO explode_maps VALUES + (1, map('b', 20, 'a', 10)), (2, map()), (3, NULL), (4, map('null', NULL)) + +query +SELECT id, explode(m) FROM explode_maps + +query +SELECT id, m, explode_outer(m) FROM explode_maps + +query +SELECT id, posexplode(m) FROM explode_maps + +query +SELECT id, m, posexplode_outer(m) FROM explode_maps + +query +SELECT posexplode_outer(m) FROM explode_maps + +query +SELECT id, k, v FROM explode_maps LATERAL VIEW OUTER explode(m) e AS k, v +WHERE k IS NULL OR v IS NULL + +query +SELECT id, p, k, v FROM explode_maps LATERAL VIEW OUTER posexplode(m) e AS p, k, v +WHERE p = 0 OR p IS NULL + +-- Computed and scalar maps, including non-nullable values and nested fields. +query +SELECT id, posexplode(map(id, named_struct('v', coalesce(id, 0)))) FROM explode_maps + +query +SELECT id, explode(map(1, 10, 2, 20)) FROM explode_maps + +query +SELECT id, posexplode_outer(cast(NULL AS map)) FROM explode_maps + +-- The untyped map() constructor still falls back before the cast. +query expect_fallback(MapType(NullType,NullType,false)) +SELECT id, explode_outer(cast(map() AS map)) FROM explode_maps + +statement +CREATE TABLE explode_nested_maps( + id int, m map, struct, m: map>>) USING parquet + +statement +INSERT INTO explode_nested_maps VALUES + (1, map(array(2, 1), named_struct('a', array(10, NULL), 'm', map('x', unhex('FF00'))), + array(3), NULL)), + (2, map()), (3, NULL) + +query +SELECT id, explode(m) FROM explode_nested_maps + +query +SELECT id, explode_outer(m) FROM explode_nested_maps + +query +SELECT id, posexplode(m) FROM explode_nested_maps + +query +SELECT id, posexplode_outer(m) FROM explode_nested_maps + +-- Project the output fields and compose map expansion with array expansion. +query +SELECT id, k, v.a, v.m FROM explode_nested_maps LATERAL VIEW OUTER explode(m) e AS k, v + +query +SELECT id, p, k, item FROM explode_nested_maps +LATERAL VIEW OUTER posexplode(m) e AS p, k, v +LATERAL VIEW OUTER explode(v.a) a AS item diff --git a/spark/src/test/scala/org/apache/comet/exec/CometGenerateExecSuite.scala b/spark/src/test/scala/org/apache/comet/exec/CometGenerateExecSuite.scala index 4cb6a7f7efd..a81e73dbecd 100644 --- a/spark/src/test/scala/org/apache/comet/exec/CometGenerateExecSuite.scala +++ b/spark/src/test/scala/org/apache/comet/exec/CometGenerateExecSuite.scala @@ -224,16 +224,31 @@ class CometGenerateExecSuite extends CometTestBase { } } - test("explode with map input falls back") { - withSQLConf( - CometConf.COMET_EXEC_LOCAL_TABLE_SCAN_ENABLED.key -> "true", - CometConf.COMET_EXEC_EXPLODE_ENABLED.key -> "true") { - val df = Seq((1, Map("a" -> 1, "b" -> 2)), (2, Map("c" -> 3))) - .toDF("id", "map") - .selectExpr("id", "explode(map) as (key, value)") - checkSparkAnswerAndFallbackReason( - df, - "Comet only supports explode/explode_outer for arrays, not maps") + for (generator <- Seq("explode", "explode_outer", "posexplode", "posexplode_outer")) { + test(s"$generator with map input across batch boundaries") { + withSQLConf( + CometConf.COMET_EXEC_EXPLODE_ENABLED.key -> "true", + CometConf.COMET_BATCH_SIZE.key -> "4") { + val rows = (0 until 16).map { i => + val m = i % 4 match { + case 0 => null.asInstanceOf[Map[Int, java.lang.Integer]] + case 1 => Map.empty[Int, java.lang.Integer] + case _ => + (0 until 13).map { j => + j -> (if (j % 3 == 0) null else java.lang.Integer.valueOf(i * 100 + j)) + }.toMap + } + (i, m) + } + withParquetDataFrame(rows) { input => + // One map exceeds the output batch size. Carry the map through too, so + // outer padding cannot accidentally replace the original empty map. + val query = input.toDF("id", "m").selectExpr("id", "m", s"$generator(m)") + val (_, plan) = checkSparkAnswerAndOperator(query) + assert(collect(plan) { case e: CometExplodeExec => e }.nonEmpty) + checkSparkSchema(query) + } + } } } @@ -390,16 +405,16 @@ class CometGenerateExecSuite extends CometTestBase { } } - test("posexplode with map input falls back") { + test("posexplode with map input falls back when disabled") { withSQLConf( CometConf.COMET_EXEC_LOCAL_TABLE_SCAN_ENABLED.key -> "true", - CometConf.COMET_EXEC_EXPLODE_ENABLED.key -> "true") { + CometConf.COMET_EXEC_EXPLODE_ENABLED.key -> "false") { val df = Seq((1, Map("a" -> 1, "b" -> 2)), (2, Map("c" -> 3))) .toDF("id", "map") .selectExpr("id", "posexplode(map) as (pos, key, value)") checkSparkAnswerAndFallbackReason( df, - "Comet only supports explode/explode_outer for arrays, not maps") + "Native support for operator GenerateExec is disabled") } } diff --git a/spark/src/test/scala/org/apache/spark/sql/benchmark/CometExplodeBenchmark.scala b/spark/src/test/scala/org/apache/spark/sql/benchmark/CometExplodeBenchmark.scala index aa16417ca2c..06d5008c0b1 100644 --- a/spark/src/test/scala/org/apache/spark/sql/benchmark/CometExplodeBenchmark.scala +++ b/spark/src/test/scala/org/apache/spark/sql/benchmark/CometExplodeBenchmark.scala @@ -59,11 +59,8 @@ import org.apache.comet.CometConf * processes. Worth knowing when reading these numbers: real queries do get that filter, so a * plain `explode` in production usually sees an array column with no nulls and no empty rows. * - * Only array inputs are covered. Comet declines to convert a generator over a map - * (https://github.com/apache/datafusion-comet/issues/2837), so the Comet arm of such a case would - * be Spark's `GenerateExec` behind a columnar-to-row transition and its timing would say nothing - * about `CometExplodeExec`. The nesting group below therefore reaches its event list through an - * `array>` where a map would be the more natural modeling choice. + * Only array inputs are covered. The nesting group below reaches its event list through an + * `array>`. */ object CometExplodeBenchmark extends CometBenchmarkBase { @@ -382,9 +379,7 @@ object CometExplodeBenchmark extends CometBenchmarkBase { * * The shape is the one asked for in review: a customer profile whose event list sits eight * struct accessors down, holding a second array of four-field structs inside each element. The - * requested outer container was a map keyed by platform, which is where a real schema would put - * it; that is an `array>` here because Comet has no native generator - * over maps yet (#2837) and the Comet arm would silently be Spark. + * outer container is an `array>`. * * The whole event struct is counted rather than one of its fields, so nested schema pruning * cannot narrow the exploded element and leave the case measuring a two-column gather. It does