Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
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
Prev Previous commit
Next Next commit
feat(native): add JavaScalarUdf skeleton implementing ScalarUDFImpl
Introduces JavaScalarUdf struct with its Signature, return_type, JNI
references, and a ScalarUDFImpl impl whose invoke_with_args returns
NotImplemented (placeholder for Task 6). Adds volatility_from_byte
helper. Dead-code lints suppressed with comments referencing Task 5/6.
  • Loading branch information
andygrove committed May 14, 2026
commit b4df1325ebda20a203cb109707d13917e919a049
1 change: 1 addition & 0 deletions native/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ mod csv;
mod errors;
mod proto;
mod schema;
mod udf;

pub(crate) mod proto_gen {
include!(concat!(env!("OUT_DIR"), "/datafusion_java.rs"));
Expand Down
115 changes: 115 additions & 0 deletions native/src/udf.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,115 @@
// 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.

//! Java-backed scalar UDF support.

use std::any::Any;
use std::fmt;

use datafusion::arrow::datatypes::DataType;
use datafusion::error::DataFusionError;
use datafusion::logical_expr::{
ColumnarValue, ScalarFunctionArgs, ScalarUDFImpl, Signature, Volatility,
};
use jni::objects::{GlobalRef, JStaticMethodID};

// Fields will be constructed in Task 5 (UDF registration JNI) and used in Task 6 (invoke).
#[allow(dead_code)]
pub(crate) struct JavaScalarUdf {
pub(crate) name: String,
pub(crate) signature: Signature,
pub(crate) return_type: DataType,
/// Global ref to the user's `org.apache.datafusion.ScalarUdf` instance.
pub(crate) udf_global_ref: GlobalRef,
/// Global ref to the `org.apache.datafusion.internal.JniBridge` class.
pub(crate) bridge_class: GlobalRef,
/// Method ID for `JniBridge.invokeScalarUdf`.
pub(crate) invoke_method: JStaticMethodID,
}

// SAFETY: JStaticMethodID is a JNI handle that's safe to share because the
// class it points to is held alive by `bridge_class`. We never mutate
// `invoke_method` after construction; DataFusion requires `Send + Sync` on
// `ScalarUDFImpl`.
unsafe impl Send for JavaScalarUdf {}
unsafe impl Sync for JavaScalarUdf {}

impl fmt::Debug for JavaScalarUdf {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("JavaScalarUdf")
.field("name", &self.name)
.field("return_type", &self.return_type)
.finish()
}
}

impl PartialEq for JavaScalarUdf {
fn eq(&self, other: &Self) -> bool {
// Two Java UDFs are equal iff they wrap the same registered name.
self.name == other.name
}
}

impl Eq for JavaScalarUdf {}

impl std::hash::Hash for JavaScalarUdf {
fn hash<H: std::hash::Hasher>(&self, state: &mut H) {
self.name.hash(state);
}
}

impl ScalarUDFImpl for JavaScalarUdf {
fn as_any(&self) -> &dyn Any {
self
}

fn name(&self) -> &str {
&self.name
}

fn signature(&self) -> &Signature {
&self.signature
}

fn return_type(&self, _arg_types: &[DataType]) -> datafusion::error::Result<DataType> {
Ok(self.return_type.clone())
}

fn invoke_with_args(
&self,
_args: ScalarFunctionArgs,
) -> datafusion::error::Result<ColumnarValue> {
Err(DataFusionError::NotImplemented(format!(
"JavaScalarUdf::invoke_with_args not yet implemented for '{}'",
self.name
)))
}
}

// Will be called from the UDF registration JNI entry point added in Task 5.
#[allow(dead_code)]
pub(crate) fn volatility_from_byte(byte: u8) -> datafusion::error::Result<Volatility> {
match byte {
0 => Ok(Volatility::Immutable),
1 => Ok(Volatility::Stable),
2 => Ok(Volatility::Volatile),
other => Err(DataFusionError::Execution(format!(
"unknown volatility byte: {}",
other
))),
}
}