Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
124 commits
Select commit Hold shift + click to select a range
768b3e9
impl map_from_entries
Dec 14, 2025
c68c342
Revert "impl map_from_entries"
Dec 16, 2025
d887555
Merge branch 'apache:main' into main
kazantsev-maksim Dec 16, 2025
231aa90
Merge branch 'apache:main' into main
kazantsev-maksim Dec 17, 2025
9500bbb
Merge branch 'apache:main' into main
kazantsev-maksim Dec 24, 2025
9577481
Merge branch 'apache:main' into main
kazantsev-maksim Dec 28, 2025
3791557
Merge branch 'apache:main' into main
kazantsev-maksim Jan 2, 2026
7c2f082
Merge branch 'apache:main' into main
kazantsev-maksim Jan 3, 2026
609a605
Merge branch 'apache:main' into main
kazantsev-maksim Jan 6, 2026
a151b2c
Merge branch 'apache:main' into main
kazantsev-maksim Jan 7, 2026
ad3e7f5
Merge branch 'apache:main' into main
kazantsev-maksim Jan 10, 2026
ea92e4b
Merge branch 'apache:main' into main
kazantsev-maksim Jan 14, 2026
8dfeca3
Merge branch 'apache:main' into main
kazantsev-maksim Jan 17, 2026
559741e
Merge branch 'apache:main' into main
kazantsev-maksim Jan 20, 2026
ebda14e
Merge branch 'apache:main' into main
kazantsev-maksim Jan 21, 2026
408152e
Merge branch 'apache:main' into main
kazantsev-maksim Jan 23, 2026
d7857b2
Merge branch 'apache:main' into main
kazantsev-maksim Jan 24, 2026
aef41be
Merge branch 'apache:main' into main
kazantsev-maksim Jan 29, 2026
5ac1c58
Merge branch 'apache:main' into main
kazantsev-maksim Jan 30, 2026
9ae8e23
Merge branch 'apache:main' into main
kazantsev-maksim Feb 1, 2026
5ca3888
Merge branch 'apache:main' into main
kazantsev-maksim Feb 4, 2026
160a817
Merge branch 'apache:main' into main
kazantsev-maksim Feb 5, 2026
88fc313
Merge branch 'apache:main' into main
kazantsev-maksim Feb 7, 2026
e14c180
Merge branch 'apache:main' into main
kazantsev-maksim Feb 13, 2026
610a885
Merge branch 'apache:main' into main
kazantsev-maksim Feb 20, 2026
f8acb2c
Merge branch 'apache:main' into main
kazantsev-maksim Feb 21, 2026
ec94897
Merge branch 'apache:main' into main
kazantsev-maksim Feb 26, 2026
43405e4
Merge branch 'apache:main' into main
kazantsev-maksim Feb 27, 2026
47b4915
Merge branch 'apache:main' into main
kazantsev-maksim Mar 1, 2026
26e2682
Merge branch 'apache:main' into main
kazantsev-maksim Mar 3, 2026
6cb5f07
Merge branch 'apache:main' into main
kazantsev-maksim Mar 4, 2026
ec194fb
Merge branch 'apache:main' into main
kazantsev-maksim Mar 31, 2026
256fccb
Merge branch 'apache:main' into main
kazantsev-maksim Apr 3, 2026
912c8f9
Merge branch 'apache:main' into main
kazantsev-maksim Apr 3, 2026
561a664
Merge branch 'apache:main' into main
kazantsev-maksim Apr 8, 2026
d926ef4
Merge branch 'apache:main' into main
kazantsev-maksim Apr 11, 2026
671412c
Merge branch 'apache:main' into main
kazantsev-maksim Apr 17, 2026
c9f52d1
Merge branch 'apache:main' into main
kazantsev-maksim Apr 22, 2026
67f72d9
Merge branch 'apache:main' into main
kazantsev-maksim Apr 23, 2026
314e594
Merge branch 'apache:main' into main
kazantsev-maksim Apr 24, 2026
ac8292f
Merge branch 'apache:main' into main
kazantsev-maksim May 1, 2026
c9c140e
Merge branch 'apache:main' into main
kazantsev-maksim May 7, 2026
decca58
Merge branch 'apache:main' into main
kazantsev-maksim May 13, 2026
0919b33
Merge branch 'apache:main' into main
kazantsev-maksim May 16, 2026
7495e21
Merge branch 'apache:main' into main
kazantsev-maksim May 19, 2026
0a37a60
Merge branch 'apache:main' into main
kazantsev-maksim May 21, 2026
abbba84
Merge branch 'apache:main' into main
kazantsev-maksim May 25, 2026
6020560
Merge branch 'apache:main' into main
kazantsev-maksim May 28, 2026
e2bdfb1
Merge branch 'apache:main' into main
kazantsev-maksim May 31, 2026
3edfc33
Merge branch 'apache:main' into main
kazantsev-maksim Jun 3, 2026
a39e860
Merge branch 'apache:main' into main
kazantsev-maksim Jun 4, 2026
e88dd7b
Merge branch 'apache:main' into main
kazantsev-maksim Jun 5, 2026
3e29d37
Merge branch 'apache:main' into main
kazantsev-maksim Jun 7, 2026
4068359
Merge branch 'apache:main' into main
kazantsev-maksim Jun 12, 2026
a3cb8de
Merge branch 'apache:main' into main
kazantsev-maksim Jun 13, 2026
b33726f
Merge branch 'apache:main' into main
kazantsev-maksim Jun 21, 2026
698f7a1
Merge branch 'apache:main' into main
kazantsev-maksim Jun 22, 2026
18162a6
Merge branch 'apache:main' into main
kazantsev-maksim Jun 23, 2026
6b0d500
Support native DataFusion lambda functions
Jun 28, 2026
4281483
more tests
Jun 29, 2026
d157db2
Merge branch 'main' into array_filter
kazantsev-maksim Jun 29, 2026
d056f64
fix
Jun 29, 2026
6037e7a
Merge remote-tracking branch 'origin/array_filter' into array_filter
Jun 29, 2026
30f06f4
Fix PR issues
Jul 1, 2026
47f2726
Fix PR issues
Jul 1, 2026
435365a
Fix PR issues
Jul 1, 2026
6f6eb6f
Merge branch 'apache:main' into main
kazantsev-maksim Jul 1, 2026
c21a42e
Merge branch 'apache:main' into main
kazantsev-maksim Jul 2, 2026
528d392
Merge remote-tracking branch 'origin/main' into array_filter
Jul 2, 2026
618ae48
Merge branch 'apache:main' into main
kazantsev-maksim Jul 3, 2026
b6f795d
Merge remote-tracking branch 'refs/remotes/origin/main' into array_fi…
Jul 3, 2026
fc13227
fix PR issues
Jul 3, 2026
9102c3d
Fix PR issues
Jul 3, 2026
4d068e3
Merge branch 'apache:main' into main
kazantsev-maksim Jul 3, 2026
6ef4112
Merge remote-tracking branch 'origin/main' into array_filter
Jul 3, 2026
a2f519f
Merge branch 'apache:main' into main
kazantsev-maksim Jul 4, 2026
36e13d5
Merge branch 'apache:main' into main
kazantsev-maksim Jul 7, 2026
2c8ae52
Merge branch 'apache:main' into main
kazantsev-maksim Jul 8, 2026
c90191f
fix nested lambda processing
Jul 9, 2026
9123ade
some refactoring
Jul 9, 2026
d57069d
refactoring
Jul 10, 2026
d6a1436
refactoring
Jul 11, 2026
1c6ce99
refactoring
Jul 11, 2026
593f7b6
Merge branch 'apache:main' into main
kazantsev-maksim Jul 11, 2026
ea5b38a
refactoring
Jul 11, 2026
b61a258
Merge remote-tracking branch 'origin/main' into array_filter
Jul 11, 2026
b697292
refactoring
Jul 11, 2026
27a19c6
fmt
Jul 11, 2026
b1d3a1a
Merge branch 'apache:main' into main
kazantsev-maksim Jul 14, 2026
8faa885
Merge remote-tracking branch 'origin/main' into array_filter
Jul 14, 2026
e2de8c0
Merge branch 'apache:main' into main
kazantsev-maksim Jul 26, 2026
e6fd376
Merge branch 'apache:main' into main
kazantsev-maksim Aug 1, 2026
11528e3
Merge branch 'apache:main' into main
kazantsev-maksim Aug 1, 2026
e17397f
Merge branch 'apache:main' into main
kazantsev-maksim Aug 9, 2026
b79e747
Merge branch 'apache:main' into main
kazantsev-maksim Aug 12, 2026
6673b5c
Merge remote-tracking branch 'origin/main' into array_filter
Aug 12, 2026
5824111
fix
Aug 12, 2026
d191fcf
fix
Aug 12, 2026
07153ae
fix
Aug 12, 2026
314e2cb
fix
Aug 16, 2026
b7f9582
fix
Aug 16, 2026
3988b0f
Merge branch 'apache:main' into main
kazantsev-maksim Aug 19, 2026
f705c43
Merge branch 'apache:main' into main
kazantsev-maksim Aug 20, 2026
741747d
Merge branch 'apache:main' into main
kazantsev-maksim Aug 22, 2026
9495593
Merge branch 'apache:main' into main
kazantsev-maksim Aug 27, 2026
0ac77df
Merge branch 'apache:main' into main
kazantsev-maksim Aug 28, 2026
1346dc4
Merge branch 'apache:main' into main
kazantsev-maksim Aug 31, 2026
b2e0ffd
Merge branch 'apache:main' into main
kazantsev-maksim Sep 1, 2026
0426c29
Merge branch 'apache:main' into main
kazantsev-maksim Sep 6, 2026
d6c3209
Merge remote-tracking branch 'origin/main' into array_filter
Sep 6, 2026
6c64305
work
Sep 6, 2026
14e5c2c
work
Sep 7, 2026
9bee07b
Merge branch 'apache:main' into main
kazantsev-maksim Sep 11, 2026
de5f960
Merge remote-tracking branch 'origin/main' into array_filter
Sep 11, 2026
9c9f826
resolve conflicts
Sep 11, 2026
8d0d9cd
Merge branch 'main' into array_filter
kazantsev-maksim Sep 13, 2026
50df758
Merge branch 'apache:main' into main
kazantsev-maksim Sep 13, 2026
4b49be9
Merge branch 'main' into array_filter
kazantsev-maksim Sep 15, 2026
0241d02
Merge branch 'main' into array_filter
kazantsev-maksim Sep 16, 2026
924b9fc
Merge branch 'main' into array_filter
kazantsev-maksim Sep 17, 2026
971f72a
Merge remote-tracking branch 'origin/main' into array_filter
Sep 17, 2026
c6b7599
Merge branch 'apache:main' into main
kazantsev-maksim Sep 17, 2026
65bd9f4
Merge remote-tracking branch 'origin/main' into array_filter
Sep 17, 2026
95a4302
refactoring
Sep 17, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
69 changes: 69 additions & 0 deletions native/core/src/execution/lambda.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,69 @@
// 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.

//! Helpers for planning DataFusion higher-order functions (HOFs) coming
//! from Spark.
//!
//! The planner needs three things that don't belong in `planner.rs`:
//! 1. A stack of *lambda scopes* so nested `NamedLambdaVariable`s resolve
//! by Spark `exprId` (immune to name shadowing / column collisions).
//! 2. A drop-guard that pops a scope on any exit path (`?`, panic-safe).
//! 3. A tiny `PhysicalExpr` wrapper that keeps *unused* lambda parameters
//! visible in `children()` so `LambdaExpr::new`'s projection compaction
//! stays consistent with the runtime batch layout.

use std::cell::RefCell;
use std::collections::HashMap;

use arrow::datatypes::FieldRef;
use datafusion::common::Result;

/// Maps Spark `exprId` -> (column index in the extended body schema, field).
pub(crate) type LambdaScope = HashMap<i64, (usize, FieldRef)>;

/// A stack of lambda variable scopes, innermost last.
/// Planning is single-threaded per planner, so `RefCell` is sufficient to manage
/// the stack of scopes during the recursive planning process.
#[derive(Default)]
pub(crate) struct LambdaScopes {
stack: RefCell<Vec<LambdaScope>>,
}

impl LambdaScopes {
/// Resolve a lambda variable by Spark `exprId`, searching innermost
/// scope first.
pub(crate) fn resolve_variable(&self, expr_id: i64) -> Option<(usize, FieldRef)> {
self.stack
.borrow()
.iter()
.rev()
.find_map(|s| s.get(&expr_id).cloned())
}

/// Push `scope`, run `f`, pop unconditionally. The pop happens on both
/// the `Ok` and `Err` paths — this replaces the earlier RAII guard.
pub(crate) fn with_scope<T, E>(
&self,
scope: LambdaScope,
f: impl FnOnce() -> Result<T, E>,
) -> Result<T, E> {
self.stack.borrow_mut().push(scope);
let out = f();
self.stack.borrow_mut().pop();
out
}
}
1 change: 1 addition & 0 deletions native/core/src/execution/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ pub mod operators;
pub(crate) mod planner;
pub mod serde;
pub use datafusion_comet_shuffle as shuffle;
pub mod lambda;
mod memory_pools;
pub(crate) mod sort;
pub(crate) mod spark_config;
Expand Down
163 changes: 155 additions & 8 deletions native/core/src/execution/planner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -91,9 +91,10 @@ use datafusion::{
prelude::SessionContext,
};
use datafusion_comet_spark_expr::{
create_comet_physical_fun, create_comet_physical_fun_with_eval_mode, BinaryOutputStyle,
BloomFilterAgg, BloomFilterMightContain, CometCollectList, CometCollectSet, CsvWriteOptions,
EvalMode, SparkArraysZipFunc, SparkBloomFilterVersion, SparkPercentile, SumInteger, ToCsv,
create_comet_hof_func, create_comet_physical_fun, create_comet_physical_fun_with_eval_mode,
BinaryOutputStyle, BloomFilterAgg, BloomFilterMightContain, CometCollectList, CometCollectSet,
CsvWriteOptions, EvalMode, SparkArraysZipFunc, SparkBloomFilterVersion, SparkPercentile,
SumInteger, ToCsv,
};
use iceberg::expr::Bind;

Expand All @@ -111,13 +112,15 @@ use datafusion::datasource::listing::PartitionedFile;
use datafusion::logical_expr::type_coercion::functions::fields_with_udf;
use datafusion::logical_expr::type_coercion::other::get_coerce_type_for_case_expression;
use datafusion::logical_expr::{
AggregateUDF, ReturnFieldArgs, ScalarUDF, TypeSignature, WindowFrame, WindowFrameBound,
WindowFrameUnits, WindowFunctionDefinition,
AggregateUDF, HigherOrderUDF, LambdaParametersProgress, ReturnFieldArgs, ScalarUDF,
TypeSignature, ValueOrLambda, WindowFrame, WindowFrameBound, WindowFrameUnits,
WindowFunctionDefinition,
};
use datafusion::physical_expr::expressions::{Literal, StatsType};
use datafusion::physical_expr::expressions::{LambdaExpr, LambdaVariable, Literal, StatsType};
use datafusion::physical_expr::window::WindowExpr;
use datafusion::physical_expr::LexOrdering;
use datafusion::physical_expr::{HigherOrderFunctionExpr, LexOrdering};

use crate::execution::lambda::{LambdaScope, LambdaScopes};
use crate::parquet::parquet_exec::init_datasource_exec;
use arrow::array::{
new_empty_array, Array, ArrayRef, BinaryBuilder, BooleanArray, Date32Array, Decimal128Array,
Expand All @@ -132,7 +135,7 @@ use datafusion::physical_plan::filter::FilterExec;
use datafusion::physical_plan::joins::NestedLoopJoinExec;
use datafusion::physical_plan::limit::GlobalLimitExec;
use datafusion::physical_plan::unnest::ListUnnest;
use datafusion_comet_proto::spark_expression::ListLiteral;
use datafusion_comet_proto::spark_expression::{HigherOrderFunc, LambdaFunction, ListLiteral};
use datafusion_comet_proto::spark_operator::SparkFilePartition;
use datafusion_comet_proto::{
spark_expression::{
Expand Down Expand Up @@ -307,6 +310,9 @@ pub struct PhysicalPlanner {
/// Task-owned destination for remote shuffle blocks, registered on the driving Spark task
/// thread before native planning. Only explicit RSS destinations may use it.
shuffle_partition_pusher: Option<Arc<dyn ShufflePartitionPusher>>,
/// Stack of lambda variable scopes, innermost last. Planning is
/// single-threaded per planner, so `RefCell` is sufficient.
lambda_scopes: LambdaScopes,
}

impl Default for PhysicalPlanner {
Expand All @@ -326,6 +332,7 @@ impl PhysicalPlanner {
task_context: None,
class_loader: None,
shuffle_partition_pusher: None,
lambda_scopes: LambdaScopes::default(),
}
}

Expand Down Expand Up @@ -725,6 +732,18 @@ impl PhysicalPlanner {
_ => func,
}
}
ExprStruct::HighOrderFunc(hof) => {
self.create_high_order_function_expr(hof, input_schema)
}
ExprStruct::NamedLambdaVariable(nlv) => {
let (idx, field) = self.lambda_scopes.resolve_variable(nlv.expr_id).ok_or_else(|| {
GeneralError(format!(
"Lambda variable '{}' (exprId={}) is not bound in any enclosing lambda scope",
nlv.name, nlv.expr_id
))
})?;
Ok(Arc::new(LambdaVariable::new(idx, field)))
}
ExprStruct::CaseWhen(case_when) => {
let when_then_pairs = case_when
.when
Expand Down Expand Up @@ -3670,6 +3689,134 @@ impl PhysicalPlanner {
}
}

fn create_high_order_function_expr(
&self,
expr: &HigherOrderFunc,
input_schema: SchemaRef,
) -> Result<Arc<dyn PhysicalExpr>, ExecutionError> {
let udf = create_comet_hof_func(&expr.func_name, &self.session_ctx.state())?;

// 1. Plan value args.
let value_args: Vec<Arc<dyn PhysicalExpr>> = expr
.value_args
.iter()
.map(|e| self.create_expr(e, Arc::clone(&input_schema)))
.collect::<Result<_, _>>()?;

// 2. Resolve lambda param field types via the UDF (mirrors runtime).
let param_fields = Self::resolve_lambda_param_fields(
&udf,
&expr.func_name,
&value_args,
expr.lambdas.len(),
input_schema.as_ref(),
)?;

// 3. Plan lambdas with resolved param fields.
let lambdas: Vec<Arc<dyn PhysicalExpr>> = expr
.lambdas
.iter()
.zip(&param_fields)
.map(|(l, fields)| self.create_lambda_expr(l, &input_schema, fields))
.collect::<Result<_, _>>()?;

// 4. NOTE: assumes value args precede lambdas (holds for array_filter).
let mut args = value_args;
args.extend(lambdas);

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

A short comment that this assumes all value args precede all lambdas would help. That holds for array_filter and the current single-lambda functions, but would not generalize to a future HOF with interleaved value and lambda args.


Ok(Arc::new(HigherOrderFunctionExpr::try_new_with_schema(
udf,
args,
&input_schema,
Arc::new(ConfigOptions::default()),
)?))
}

fn resolve_lambda_param_fields(
udf: &HigherOrderUDF,
func_name: &str,
value_args: &[Arc<dyn PhysicalExpr>],
lambda_count: usize,
schema: &Schema,
) -> Result<Vec<Vec<FieldRef>>, ExecutionError> {
let mut planning_fields: Vec<ValueOrLambda<FieldRef, Option<FieldRef>>> = value_args
.iter()
.map(|e| Ok(ValueOrLambda::Value(e.return_field(schema)?)))
.collect::<Result<_, DataFusionError>>()?;
planning_fields.extend(std::iter::repeat_n(
ValueOrLambda::Lambda(None),
lambda_count,
));

match udf.lambda_parameters(0, &planning_fields)? {
LambdaParametersProgress::Complete(items) if items.len() >= lambda_count => Ok(items),
LambdaParametersProgress::Complete(items) => Err(GeneralError(format!(
"{func_name}: expected parameter fields for {lambda_count} lambdas, got {}",
items.len()
))),
LambdaParametersProgress::Partial(_) => Err(GeneralError(format!(
"{func_name}: multi-step lambda resolution is not supported yet"
))),
}
}

fn create_lambda_expr(
&self,
lambda: &LambdaFunction,
input_schema: &SchemaRef,
param_fields: &[FieldRef],
) -> Result<Arc<dyn PhysicalExpr>, ExecutionError> {
if param_fields.len() < lambda.args.len() {
return Err(GeneralError(format!(
"lambda declares {} params but the function resolved only {}",
lambda.args.len(),
param_fields.len()
)));
}

// Build extended schema = input schema ++ lambda params, and the scope.
let mut body_fields: Vec<FieldRef> = input_schema.fields().iter().map(Arc::clone).collect();
let mut scope = LambdaScope::with_capacity(lambda.args.len());
let mut scope_entries: Vec<(usize, FieldRef)> = Vec::with_capacity(lambda.args.len());
let mut param_names: Vec<String> = Vec::with_capacity(lambda.args.len());

for (arg, resolved) in lambda.args.iter().zip(param_fields) {
// Runtime uses `param.renamed(name)` — do the same here.
let field: FieldRef = Arc::new(resolved.as_ref().clone().with_name(&arg.name));
let idx = body_fields.len();
if scope
.insert(arg.expr_id, (idx, Arc::clone(&field)))
.is_some()
{
return Err(GeneralError(format!(
"duplicate lambda variable exprId {} ('{}')",
arg.expr_id, arg.name
)));
}
scope_entries.push((idx, Arc::clone(&field)));
param_names.push(arg.name.clone());
body_fields.push(field);
}

let body_schema = Arc::new(Schema::new(
body_fields
.iter()
.map(|f| f.as_ref().clone())
.collect::<Vec<_>>(),
));
let lambda_body = lambda
.body
.as_ref()
.ok_or_else(|| GeneralError("lambda has no body".to_string()))?;

// Plan the body under this scope; the guard pops on any `?` / drop.
let body_expr = self
.lambda_scopes
.with_scope(scope, || self.create_expr(lambda_body, body_schema))?;

Ok(Arc::new(LambdaExpr::try_new(param_names, body_expr)?))
}

fn create_scalar_function_expr(
&self,
expr: &ScalarFunc,
Expand Down
20 changes: 20 additions & 0 deletions native/proto/src/proto/expr.proto
Original file line number Diff line number Diff line change
Expand Up @@ -94,6 +94,8 @@ message Expr {
Shuffle shuffle = 72;
RandStr rand_str = 73;
Uuid uuid = 74;
HigherOrderFunc high_order_func = 75;
NamedLambdaVariable named_lambda_variable = 76;
}

reserved 20;
Expand Down Expand Up @@ -675,3 +677,21 @@ message JvmScalarUdf {
// Whether the result column may contain nulls.
bool return_nullable = 4;
}

message HigherOrderFunc {
string func_name = 1;
repeated Expr value_args = 2;
repeated LambdaFunction lambdas = 3;
}

message NamedLambdaVariable {
string name = 1;
DataType data_type = 2;
bool nullable = 3;
int64 expr_id = 4;
}

message LambdaFunction {
Expr body = 1;
repeated NamedLambdaVariable args = 2;
}
30 changes: 30 additions & 0 deletions native/spark-expr/src/comet_high_order_funcs.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,30 @@
// 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 datafusion::common::DataFusionError;
use datafusion::execution::FunctionRegistry;
use datafusion::logical_expr::HigherOrderUDF;
use std::sync::Arc;

pub fn create_comet_hof_func(
func_name: &str,
registry: &dyn FunctionRegistry,
) -> Result<Arc<HigherOrderUDF>, DataFusionError> {
registry.higher_order_function(func_name).map_err(|e| {
DataFusionError::Execution(format!("HOF {func_name} not found in the registry: {e}"))
})
}
2 changes: 2 additions & 0 deletions native/spark-expr/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@ pub use predicate_funcs::{spark_isnan, RLike};

mod agg_funcs;
mod array_funcs;
mod comet_high_order_funcs;
mod comet_scalar_funcs;
pub mod hash_funcs;

Expand Down Expand Up @@ -71,6 +72,7 @@ pub use conditional_funcs::*;
pub use conversion_funcs::*;
pub use nondetermenistic_funcs::*;

pub use comet_high_order_funcs::create_comet_hof_func;
pub use comet_scalar_funcs::{
create_comet_physical_fun, create_comet_physical_fun_with_eval_mode,
register_all_comet_functions,
Expand Down
10 changes: 10 additions & 0 deletions spark/src/main/scala/org/apache/comet/CometConf.scala
Original file line number Diff line number Diff line change
Expand Up @@ -409,6 +409,16 @@ object CometConf extends ShimCometConf {
.booleanConf
.createWithDefault(true)

val COMET_EXEC_HIGHER_ORDER_FUNCTION_NATIVE_ENABLED: ConfigEntry[Boolean] =
conf("spark.comet.exec.higherOrderFunction.native.enabled")
.category(CATEGORY_EXEC)
.doc(
"When enabled, supported higher-order functions (e.g. filter) are executed by the " +
"native DataFusion engine. Shapes the native path cannot handle fall back to the " +
"codegen dispatcher, and finally to Spark.")
.booleanConf
.createWithDefault(true)

val COMET_SHUFFLE_NATIVE_HASH_PARTITIONING_ENABLED: ConfigEntry[Boolean] =
conf("spark.comet.shuffle.native.partitioning.hash.enabled")
.withAlternative("spark.comet.native.shuffle.partitioning.hash.enabled")
Expand Down
Loading