Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
75 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
6f6eb6f
Merge branch 'apache:main' into main
kazantsev-maksim Jul 1, 2026
c21a42e
Merge branch 'apache:main' into main
kazantsev-maksim Jul 2, 2026
618ae48
Merge branch 'apache:main' into main
kazantsev-maksim Jul 3, 2026
4d068e3
Merge branch 'apache:main' into main
kazantsev-maksim 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
593f7b6
Merge branch 'apache:main' into main
kazantsev-maksim Jul 11, 2026
b1d3a1a
Merge branch 'apache:main' into main
kazantsev-maksim 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
57d469c
work
Aug 9, 2026
3303a04
work
Aug 9, 2026
c1ec9a3
work
Aug 9, 2026
beda70b
fmt
Aug 9, 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
4 changes: 4 additions & 0 deletions native/spark-expr/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -222,4 +222,8 @@ harness = false

[[bench]]
name = "cast_int_to_decimal"
harness = false

[[bench]]
name = "contains"
harness = false
117 changes: 117 additions & 0 deletions native/spark-expr/benches/contains.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,117 @@
// 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 arrow::array::{ArrayRef, StringArray};
use arrow::datatypes::{DataType, Field};
use criterion::{criterion_group, criterion_main, Criterion};
use datafusion::common::ScalarValue;
use datafusion::config::ConfigOptions;
use datafusion::logical_expr::{ColumnarValue, ScalarFunctionArgs, ScalarUDFImpl};
use datafusion_comet_spark_expr::SparkContains;

use std::sync::Arc;

fn generate_string_array(size: usize) -> ArrayRef {
let data: Vec<Option<String>> = (0..size)
.map(|i| {
if i % 10 == 0 {
None
} else {
Some(format!(
"hello string data sample number {} with some text",
i
))
}
})
.collect();
Arc::new(StringArray::from(data))
}

fn bench_contains(c: &mut Criterion) {
let rows = 8192;
let udf = SparkContains::new();

let mut group = c.benchmark_group("string_funcs/contains");

let haystack_array = generate_string_array(rows);
let needle_scalar = ColumnarValue::Scalar(ScalarValue::Utf8(Some("sample".to_string())));
let needle_array = generate_string_array(rows);

// Общие метаданные для ScalarFunctionArgs
let arg_fields = vec![
Arc::new(Field::new("haystack", DataType::Utf8, true)),
Arc::new(Field::new("needle", DataType::Utf8, true)),
];
let return_field = Arc::new(Field::new("result", DataType::Boolean, true));
let config_options = Arc::new(ConfigOptions::new());

// 1. Array haystack vs Scalar needle (optimized path)
group.bench_function(&format!("array_vs_scalar_size_{}", rows), |b| {
b.iter(|| {
let args = ScalarFunctionArgs {
args: vec![
ColumnarValue::Array(haystack_array.clone()),
needle_scalar.clone(),
],
arg_fields: arg_fields.clone(),
number_rows: rows,
return_field: return_field.clone(),
config_options: config_options.clone(),
};
std::hint::black_box(udf.invoke_with_args(args).unwrap());
});
});

// 2. Array haystack vs Array needle
group.bench_function(&format!("array_vs_array_size_{}", rows), |b| {
b.iter(|| {
let args = ScalarFunctionArgs {
args: vec![
ColumnarValue::Array(haystack_array.clone()),
ColumnarValue::Array(needle_array.clone()),
],
arg_fields: arg_fields.clone(),
number_rows: rows,
return_field: return_field.clone(),
config_options: config_options.clone(),
};
std::hint::black_box(udf.invoke_with_args(args).unwrap());
});
});

let haystack_scalar_val = ColumnarValue::Scalar(ScalarValue::Utf8(Some("sample".to_string())));
group.bench_function(&format!("scalar_vs_array_size_{}", rows), |b| {
b.iter(|| {
let args = ScalarFunctionArgs {
args: vec![
haystack_scalar_val.clone(),
ColumnarValue::Array(needle_array.clone()),
],
arg_fields: arg_fields.clone(),
number_rows: rows,
return_field: return_field.clone(),
config_options: config_options.clone(),
};
std::hint::black_box(udf.invoke_with_args(args).unwrap());
});
});

group.finish();
}

criterion_group!(benches, bench_contains);
criterion_main!(benches);
118 changes: 76 additions & 42 deletions native/spark-expr/src/string_funcs/contains.rs
Original file line number Diff line number Diff line change
Expand Up @@ -83,15 +83,14 @@ fn spark_contains(haystack: &ColumnarValue, needle: &ColumnarValue) -> Result<Co

// Array haystack, scalar needle - OPTIMIZED PATH
(ColumnarValue::Array(haystack_array), ColumnarValue::Scalar(needle_scalar)) => {
let result = contains_with_arrow_scalar(haystack_array, needle_scalar)?;
let result = contains_array_scalar(haystack_array, needle_scalar)?;
Ok(ColumnarValue::Array(result))
}

// Scalar haystack, array needle - less common
(ColumnarValue::Scalar(haystack_scalar), ColumnarValue::Array(needle_array)) => {
let haystack_array = haystack_scalar.to_array_of_size(needle_array.len())?;
let result = arrow_contains(&haystack_array, needle_array)?;
Ok(ColumnarValue::Array(Arc::new(result)))
let result = contains_scalar_array(haystack_scalar, needle_array)?;
Ok(ColumnarValue::Array(result))
}

// Both scalars - compute single result
Expand All @@ -102,9 +101,24 @@ fn spark_contains(haystack: &ColumnarValue, needle: &ColumnarValue) -> Result<Co
}
}

/// Helper to safely extract string reference from ScalarValue
#[inline]
fn get_string_scalar_value<'a>(scalar: &'a ScalarValue, arg_name: &str) -> Result<&'a str> {
match scalar {
ScalarValue::Utf8(Some(s))
| ScalarValue::LargeUtf8(Some(s))
| ScalarValue::Utf8View(Some(s)) => Ok(s.as_str()),
_ => exec_err!(
"contains function requires string type for {}, got {:?}",
arg_name,
scalar.data_type()
),
}
}

/// Optimized contains for array haystack with scalar needle.
/// Uses Arrow's native scalar handling for better performance.
fn contains_with_arrow_scalar(
fn contains_array_scalar(
haystack_array: &ArrayRef,
needle_scalar: &ScalarValue,
) -> Result<ArrayRef> {
Expand All @@ -114,17 +128,7 @@ fn contains_with_arrow_scalar(
}

// Extract the needle string
let needle_str = match needle_scalar {
ScalarValue::Utf8(Some(s))
| ScalarValue::LargeUtf8(Some(s))
| ScalarValue::Utf8View(Some(s)) => s.clone(),
_ => {
return exec_err!(
"contains function requires string type for needle, got {:?}",
needle_scalar.data_type()
)
}
};
let needle_str = get_string_scalar_value(needle_scalar, "needle")?;

// Create scalar array for needle - tells Arrow to use optimized paths
let needle_scalar_array = StringArray::new_scalar(needle_str);
Expand All @@ -134,6 +138,21 @@ fn contains_with_arrow_scalar(
Ok(Arc::new(result))
}

fn contains_scalar_array(
haystack_scalar: &ScalarValue,
needle_array: &ArrayRef,
) -> Result<ArrayRef> {
if haystack_scalar.is_null() {
return Ok(Arc::new(BooleanArray::new_null(needle_array.len())));
}

let haystack_str = get_string_scalar_value(haystack_scalar, "haystack")?;
let haystack_scalar_array = StringArray::new_scalar(haystack_str.to_string());

let result = arrow_contains(&haystack_scalar_array, needle_array)?;
Ok(Arc::new(result))
}

/// Contains for two scalar values.
fn contains_scalar_scalar(
haystack_scalar: &ScalarValue,
Expand All @@ -144,29 +163,8 @@ fn contains_scalar_scalar(
return Ok(ScalarValue::Boolean(None));
}

let haystack_str = match haystack_scalar {
ScalarValue::Utf8(Some(s))
| ScalarValue::LargeUtf8(Some(s))
| ScalarValue::Utf8View(Some(s)) => s.as_str(),
_ => {
return exec_err!(
"contains function requires string type for haystack, got {:?}",
haystack_scalar.data_type()
)
}
};

let needle_str = match needle_scalar {
ScalarValue::Utf8(Some(s))
| ScalarValue::LargeUtf8(Some(s))
| ScalarValue::Utf8View(Some(s)) => s.as_str(),
_ => {
return exec_err!(
"contains function requires string type for needle, got {:?}",
needle_scalar.data_type()
)
}
};
let haystack_str = get_string_scalar_value(haystack_scalar, "haystack")?;
let needle_str = get_string_scalar_value(needle_scalar, "needle")?;

Ok(ScalarValue::Boolean(Some(
haystack_str.contains(needle_str),
Expand All @@ -188,7 +186,7 @@ mod tests {
])) as ArrayRef;
let needle = ScalarValue::Utf8(Some("world".to_string()));

let result = contains_with_arrow_scalar(&haystack, &needle).unwrap();
let result = contains_array_scalar(&haystack, &needle).unwrap();
let bool_array = result.as_any().downcast_ref::<BooleanArray>().unwrap();

assert!(bool_array.value(0)); // "hello world" contains "world"
Expand Down Expand Up @@ -218,7 +216,7 @@ mod tests {
])) as ArrayRef;
let needle = ScalarValue::Utf8(None);

let result = contains_with_arrow_scalar(&haystack, &needle).unwrap();
let result = contains_array_scalar(&haystack, &needle).unwrap();
let bool_array = result.as_any().downcast_ref::<BooleanArray>().unwrap();

// Null needle should produce null results
Expand All @@ -231,11 +229,47 @@ mod tests {
let haystack = Arc::new(StringArray::from(vec![Some("hello world"), Some("")])) as ArrayRef;
let needle = ScalarValue::Utf8(Some("".to_string()));

let result = contains_with_arrow_scalar(&haystack, &needle).unwrap();
let result = contains_array_scalar(&haystack, &needle).unwrap();
let bool_array = result.as_any().downcast_ref::<BooleanArray>().unwrap();

// Empty string is contained in any string
assert!(bool_array.value(0));
assert!(bool_array.value(1));
}

#[test]
fn test_contains_scalar_array_null_haystack() {
let haystack = ScalarValue::Utf8(None);
let needle = Arc::new(StringArray::from(vec![
Some("hello world"),
Some("foo bar"),
])) as ArrayRef;

let result = contains_scalar_array(&haystack, &needle).unwrap();
let bool_array = result.as_any().downcast_ref::<BooleanArray>().unwrap();

// Null haystack should produce null results for all array elements
assert!(bool_array.is_null(0));
assert!(bool_array.is_null(1));
}

#[test]
fn test_spark_contains_dispatcher_scalar_array() {
let haystack = ColumnarValue::Scalar(ScalarValue::Utf8(Some("abc".to_string())));
let needle =
ColumnarValue::Array(
Arc::new(StringArray::from(vec![Some("a"), Some("bc"), Some("d")])) as ArrayRef,
);

let result = spark_contains(&haystack, &needle).unwrap();
let array = match result {
ColumnarValue::Array(arr) => arr,
_ => panic!("Expected ColumnarValue::Array"),
};
let bool_array = array.as_any().downcast_ref::<BooleanArray>().unwrap();

assert!(bool_array.value(0));
assert!(bool_array.value(1));
assert!(!bool_array.value(2));
}
}
Loading