andygrove commented on code in PR #4459:
URL: https://github.com/apache/datafusion-comet/pull/4459#discussion_r4092050234


##########
native/core/src/execution/rust_udf/imported_c.rs:
##########
@@ -0,0 +1,305 @@
+// 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.
+
+//! Adapter wrapping a C-ABI [`CometCScalarKernel`] as a DataFusion
+//! [`ScalarUDFImpl`].
+//!
+//! Lifecycle inside `invoke_with_args`:
+//!
+//! 1. Build a fresh [`CometCScalarKernelImpl`] via the kernel's `new_impl`.
+//! 2. Call `init` with the input field types (and any scalar args) to get
+//!    the return type.
+//! 3. Call `execute` once with the batch.
+//! 4. Drop the impl (its `release` callback runs).
+
+use std::ffi::CStr;
+use std::sync::Mutex;
+
+use arrow::array::ArrayRef;
+use arrow::datatypes::{DataType, Field};
+use arrow::ffi::{from_ffi_and_data_type, FFI_ArrowArray, FFI_ArrowSchema};
+use comet_udf_sdk::c_abi::{CometCScalarKernel, CometCScalarKernelImpl};
+use datafusion::common::DataFusionError;
+use datafusion::logical_expr::{
+    ColumnarValue, ScalarFunctionArgs, ScalarUDFImpl, Signature, 
TypeSignature, Volatility,
+};
+
+/// Adapter wrapping a [`CometCScalarKernel`] as a DataFusion
+/// [`ScalarUDFImpl`].
+pub struct ImportedCScalarUdf {
+    name: String,
+    /// Boxed so the kernel's address is stable; held inside a Mutex
+    /// because the FFI Drop is not Sync-safe under concurrent invocation.
+    /// The kernel itself is logically immutable post-load — the lock only
+    /// protects the FFI calls' aliasing rules. (DataFusion serializes
+    /// invocations of a given ScalarUDFImpl per-batch through
+    /// invoke_with_args anyway; the lock is defensive.)
+    kernel: Mutex<Box<CometCScalarKernel>>,

Review Comment:
   Correction to my last reply here: the second invariant was not true by 
construction. I said the kernel could not outlive the `LoadedLibrary` because 
the loader owns both, but nothing stopped an adapter from being cloned out of 
`udfs` and outliving the struct. `concurrent_invocations_share_one_adapter` did 
exactly that, and it is why CI was red from the end of August: dropping the 
`LoadedLibrary` dlclosed the library, Linux unmapped it, and the worker threads 
called into unmapped text.
   
   Fixed in f9dd01b48, though not with `Arc<CometCScalarKernel>`. The thing 
that has to be kept alive is the library rather than the kernel, so 
`ImportedCScalarUdf` now holds an `Arc<Library>`, declared after `kernel` so 
the kernel's `release` still runs while its code is mapped. 
`adapter_keeps_its_library_loaded` covers it: it loads a private copy of the 
library (the dynamic loader refcounts by file, and the cache holds the original 
open for the whole test process), drops the `LoadedLibrary`, and calls the 
adapter afterwards.
   
   One gap remains and is documented on the field: arrays a UDF returns carry 
release callbacks that also live in the library, and they hold no reference to 
it. That is safe only because the process-wide cache, which every adapter the 
planner uses comes from, never unloads a library.
   



##########
spark/src/main/scala/org/apache/comet/udf/CometNativeUDF.scala:
##########
@@ -0,0 +1,150 @@
+/*
+ * 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.
+ */
+
+package org.apache.comet.udf
+
+import scala.util.Try
+
+import org.apache.spark.sql.SparkSession
+import org.apache.spark.sql.expressions.UserDefinedFunction
+import org.apache.spark.sql.functions.udf
+import org.apache.spark.sql.types.DataType
+
+import com.fasterxml.jackson.databind.ObjectMapper
+import com.fasterxml.jackson.databind.node.ObjectNode
+
+/**
+ * Entry point for registering Rust scalar UDFs with Comet.
+ *
+ * The UDF cdylib is built against the `comet-udf-sdk` crate and exposes its 
functions through an
+ * ABI built only on the Arrow C Data Interface, so a compiled UDF is not tied 
to Comet's
+ * DataFusion version.
+ *
+ * This is an experimental API. It is deliberately not annotated
+ * `org.apache.comet.annotation.Public`, so it sits outside the enumerated 
public API in Comet's
+ * [[https://datafusion.apache.org/comet/about/versioning_policy.html 
versioning policy]] and
+ * carries no compatibility guarantee: it may change or be removed in any 
release, including a
+ * patch release, with no deprecation cycle.
+ */
+object CometRustUDF {
+
+  private val mapper: ObjectMapper = new ObjectMapper()
+
+  /**
+   * Register a single Rust UDF with an explicit signature.
+   *
+   * Validates the library on the driver (loads it, confirms a UDF named 
`name` exists). On
+   * success a stub Spark catalog UDF is installed (so SQL/DataFrame name 
resolution succeeds) and
+   * the driver-side registry is updated.
+   *
+   * Executors do not consult the driver's registry: the library path travels 
with the plan in the
+   * `RustUdfCall` proto, and each executor loads the library itself on first 
use. The path must
+   * therefore be valid on every executor, not just the driver.
+   *
+   * `deterministic` must be `true`. Comet plans every imported kernel as 
immutable, so a
+   * nondeterministic UDF cannot yet be expressed; passing `false` fails here 
rather than silently
+   * planning the function as pure.
+   */
+  def register(

Review Comment:
   Following up on `inputTypes`: until now it was only used for its length, so 
a call with other argument types went straight to the kernel's `return_field`. 
As of 6998df990 it is enforced. The serde compares each call's argument types 
with the registered ones, ignoring nullability, and refuses a mismatch at 
planning time with both signatures named, for example `registered with argument 
types (bigint) but is called with (int)`.
   
   I went with refusing rather than casting. The catalog stub is untyped, so 
Spark's analyzer inserts no casts for it (it only coerces a `ScalaUDF` that 
carries input encoders), and a conversion Comet picked on its own would be a 
semantic choice Spark never made. An explicit `cast` in the query is the fix 
the message points to. Keeping the parameter also means that adding coercion 
later, or deriving the signature from the library as in #5597, would not change 
`register`'s signature again.
   



-- 
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]

Reply via email to