UDF de Scala y Java con ámbito de sesión

Importante

Las UDF de Scala y Java se pueden registrar en el Catálogo de Unity para la gobernanza, la reutilización y la detectabilidad. Consulte las funciones definidas por el usuario (UDF) de Scala y Java en Unity Catalog.

En esta página se describe cómo crear UDF de Scala y Java con ámbito de sesión en Azure Databricks. Las UDF con ámbito de sesión se definen en un cuaderno o trabajo y solo se aplican a sparkSession actual. Para obtener la referencia del lenguaje SQL, consulte Funciones escalares definidas por el usuario (UDF) externas.

Elección del enfoque

Puede definir una UDF de Scala o Java de las maneras siguientes. Para comparar todos los tipos de UDF entre lenguajes, gobernanza y cómputo, consulte UDF gobernadas por Unity Catalog frente a UDF de ámbito de sesión.

Approach Description
UDF de Scala en línea Defina una UDF en un notebook usando una función de Scala o una expresión lambda. Ámbito de sesión. No es compatible con el cómputo sin servidor.
Java UDF a partir de un archivo JAR Registre una clase UDF precompilada desde un ARCHIVO JAR mediante spark.udf.registerJavaFunction. Ámbito de sesión. Compatible con la computación sin servidor.
Scala o Java UDF regulados por el catálogo de Unity Registre una UDF en el Catálogo de Unity para la gobernanza, la reutilización y la detectabilidad. Compatible con la computación sin servidor.

Requisitos

  • Las UDF de Scala en entornos de proceso habilitados para Unity Catalog con modo de acceso estándar requieren Databricks Runtime 14.2 o superior.
  • La compatibilidad con instancias ARM para UDFs de Scala en clústeres habilitados para Unity Catalog requiere Databricks Runtime 15.2 o versiones posteriores.
  • Registrar una UDF de Java desde un archivo JAR con spark.udf.registerJavaFunction requiere Databricks Runtime 18 LTS o versiones posteriores. Consulte Registro de una UDF de Java desde un archivo JAR.

Importante

Compile el archivo JAR con las mismas versiones de Scala y Apache Spark que el proceso que lo ejecuta. Una incompatibilidad puede provocar que la UDF falle en el momento del registro o de la llamada.

Marque la dependencia de Apache Spark como provided para que no se incluya en su archivo JAR. Incluya solo las dependencias de terceros que utiliza la UDF.

Registro de una función como UDF

Registrar una función de Scala como UDF mediante spark.udf.register:

val squared = (s: Long) => {
  s * s
}
spark.udf.register("square", squared)

Llamada a la UDF en Spark SQL

Cree una vista temporal y llame a la UDF en una consulta SQL:

spark.range(1, 20).createOrReplaceTempView("test")
%sql select id, square(id) as id_squared from test

Uso de UDF con dataframes

También puede llamar a una UDF mediante la API DataFrame:

import org.apache.spark.sql.functions.{col, udf}
val squared = udf((s: Long) => s * s)
display(spark.range(1, 20).select(squared(col("id")) as "id_squared"))

Archivos con UDF

Importante

Esta característica se encuentra en su versión beta. Los administradores del área de trabajo pueden controlar el acceso a esta característica desde la página Vistas previas . Consulte Administrar versiones preliminares de Azure Databricks.

El tipo Scala para un archivo es FileRef. Úselo como parámetro o tipo de retorno en una UDF, ya sea como tipo de nivel superior o anidado. Para el tipo y sus reglas de anidamiento, véase FILE tipo.

Para leer el contenido de un archivo en una UDF, llama a una de las siguientes funciones en un FileRef:

  • asLocalFile(): Devuelve un java.io.File que puedes pasar a cualquier biblioteca que acepte una ruta.
  • open(): Devuelve a java.io.InputStream para leer los bytes del archivo. La persona que llama la cierra.

Para producir un nuevo FileRef en una UDF, llama a uno de los siguientes métodos estáticos:

  • FileRef.create(uri): Crea una referencia al archivo en uri.
  • FileRef.fromBytes(bytes, destinationPath, contentType): Carga bytes en la ruta del volumen destinationPath como un archivo externo y devuelve una referencia.
  • FileRef.fromLocalFile(localFile, destinationPath, contentType): Sube un archivo local a la destinationPath ruta de Volumen como archivo externo y devuelve una referencia.

No se admite devolver un FileRef desde una UDF que escribe en una columna FILE MANAGED.

Para ejemplos en Python, Scala y SQL, incluyendo procesamiento de imágenes, detección de tipos de archivo y extracción de fotogramas de vídeo, véase Archivos de proceso con UDFs.

Registrar una UDF de Java desde un JAR

Empaquete una UDF como UN ARCHIVO JAR, agréguela a la sesión con spark.addArtifacty registre la clase UDF con spark.udf.registerJavaFunction.

Nota:

Compatible con el modo de acceso estándar y la informática sin servidor en Databricks Runtime 18 LTS o versiones posteriores. La función registrada tiene ámbito de sesión y no está registrada en el catálogo de Unity.

En los pasos siguientes se explica cómo crear un proyecto, escribir una clase UDF, crear un archivo JAR fat y registrarlo.

Paso 1: Crear el proyecto

Configure un proyecto en Scala o Java.

Scala

Cree un nuevo proyecto de Scala mediante sbt:

sbt new scala/scala-seed.g8

Reemplace el contenido del build.sbt archivo por lo siguiente. Establezca scalaVersion y la versión de spark-sql para que coincidan con su cómputo:

scalaVersion := "2.13.16"

ThisBuild / organization := "com.example"

lazy val myUDF = (project in file("."))
  .settings(
    name := "my-udf",
    libraryDependencies += "org.apache.spark" %% "spark-sql" % "4.0.0" % "provided"
  )

Habilite el complemento sbt-assembly para compilar un archivo JAR fat. Cree o edite project/assembly.sbt y agregue:

addSbtPlugin("com.eed3si9n" % "sbt-assembly" % "2.0.0")

Java

Cree un nuevo proyecto de Maven mediante el arquetipo de inicio rápido:

mvn archetype:generate \
  -DgroupId=com.example \
  -DartifactId=my-udf \
  -DarchetypeArtifactId=maven-archetype-quickstart \
  -DinteractiveMode=false

Este comando crea la estructura de proyecto estándar de Maven con src/main/java directorios y src/test/java .

En el archivo pom.xml generado, dentro de las etiquetas <project></project>, añade un bloque <properties> y configura el maven-shade-plugin para compilar un JAR fat:

<properties>
  <maven.compiler.source>17</maven.compiler.source>
  <maven.compiler.target>17</maven.compiler.target>
  <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
</properties>

<build>
  <plugins>
    <plugin>
      <groupId>org.apache.maven.plugins</groupId>
      <artifactId>maven-shade-plugin</artifactId>
      <version>3.5.0</version>
      <executions>
        <execution>
          <phase>package</phase>
          <goals>
            <goal>shade</goal>
          </goals>
        </execution>
      </executions>
    </plugin>
  </plugins>
</build>

Paso 2: Escribir la clase UDF

La clase UDF debe implementar una de las interfaces org.apache.spark.sql.api.java.UDF (UDF1 hasta UDF22), donde el número indica cuántos argumentos de entrada acepta la UDF. Implemente el método call() con su lógica.

El controlador debe ser una clase Java. spark.udf.registerJavaFunction carga la clase mediante reflexión, por lo que debe ser una clase pública de nivel superior (o static anidada) con un compilador público sin argumentos. Scala class o object no cumple este requisito y produce un error en el momento de la llamada. Puede compilar el archivo JAR con sbt, pero la propia clase UDF debe escribirse en Java.

Creación de src/main/java/com/example/MyIntegerUDF.java:

package com.example;

import org.apache.spark.sql.api.java.UDF1;

public class MyIntegerUDF implements UDF1<Integer, Integer> {
  @Override
  public Integer call(Integer x) {
    return x + 1;
  }
}

Paso 3: Compila tu JAR fat

Empaqueta tu UDF compilada en un JAR fat.

Scala

En el directorio raíz del proyecto, ejecute:

sbt clean assembly

El JAR fat se crea en target/scala-2.13/ con un nombre similar a my-udf-assembly-0.1.0-SNAPSHOT.jar.

Java

En el directorio raíz del proyecto, ejecute:

mvn clean package

El JAR fat se crea en target/ con un nombre similar a my-udf-1.0-SNAPSHOT.jar.

Paso 4: Sube tu JAR a un volumen de Unity Catalog

Suba el archivo JAR a un volumen de Unity Catalog para que sus recursos de cómputo puedan acceder a él. Si aún no dispone de un volumen, cree uno:

CREATE VOLUME IF NOT EXISTS my_catalog.my_schema.udf_jars
COMMENT 'Storage for UDF JAR files';

Suba su archivo JAR al volumen mediante Catalog Explorer:

  1. En el área de trabajo de Azure Databricks, haga clic en Data icon.Catalog para abrir el Explorador de catálogos.
  2. Seleccione el catálogo y, a continuación, seleccione el esquema que contiene el volumen.
  3. Haga clic en el nombre del volumen.
  4. Haga clic en Cargar en este volumen y seleccione el archivo JAR.
  5. Haga clic en Cargar.
  6. Una vez finalizada la carga, haga clic en el nombre de su archivo JAR y, a continuación, haga clic en Copiar ruta para copiar la ruta del volumen. Por ejemplo: /Volumes/my_catalog/my_schema/udf_jars/my-udf-assembly-0.1.0-SNAPSHOT.jar. Necesitará esta ruta en el siguiente paso.

Paso 5: Registrar y llamar a la UDF

Añada el JAR a su sesión utilizando la ruta del volumen, registre la clase UDF y llámela desde Spark SQL:

# Add the JAR containing your UDF class to the session
spark.addArtifact("/Volumes/my_catalog/my_schema/udf_jars/my-udf-assembly-0.1.0-SNAPSHOT.jar")

# Register the UDF class, providing the SQL function name,
# the fully qualified class name, and the return type
from pyspark.sql.types import IntegerType

spark.udf.registerJavaFunction(
    "my_udf",
    "com.example.MyIntegerUDF",
    IntegerType(),
)

# Call the UDF from Spark SQL
spark.sql("SELECT my_udf(21)").show()

En el cómputo sin servidor y con modo de acceso estándar, debe pasar un tipo de retorno explícito. Si se omite el tipo de valor devuelto, se produce un error en UC_COMMAND_NOT_SUPPORTED_IN_SHARED_ACCESS_MODE. Las funciones de agregado definidas por el usuario (UDAFs) no se admiten con registerJavaFunction.

La consulta devuelve la salida de UDF, confirmando que la función está registrada y invocable:

+----------+
| my_udf(21)|
+----------+
|        22|
+----------+

Orden de evaluación y comprobación nula

Spark SQL (incluidos SQL y las API de DataFrame y Dataset) no garantiza el orden de evaluación de las subexpresiones. Spark no evalúa las entradas de un operador o una función de izquierda a derecha. Las expresiones lógicas AND y OR no tienen semántica de cortocircuito de izquierda a derecha.

No confíe en los efectos secundarios ni en el orden de evaluación de las expresiones booleanas, ni en el orden de las cláusulas WHERE y HAVING. El optimizador de consultas puede reordenar estas expresiones y cláusulas. Si una UDF se basa en la semántica de cortocircuito para la comprobación de valores nulos, Spark no garantiza que la comprobación de valores nulos se ejecute antes de la UDF. Por ejemplo:

spark.udf.register("strlen", (s: String) => s.length)
spark.sql("select s from test1 where s is not null and strlen(s) > 1") // no guarantee

Esta WHERE cláusula no garantiza que Spark invoque la strlen UDF después de filtrar los valores NULL.

Para controlar la comprobación de valores NULL, Databricks recomienda cualquiera de las siguientes opciones:

  • Hacer que la propia UDF sea compatible con valores NULL y realice la comprobación de valores NULL dentro de la UDF
  • Uso de expresiones IF o CASE WHEN para realizar la comprobación de valores NULL e invocar la UDF en una rama condicional
spark.udf.register("strlen_nullsafe", (s: String) => if (s != null) s.length else -1)
spark.sql("select s from test1 where s is not null and strlen_nullsafe(s) > 1") // ok
spark.sql("select s from test1 where if(s is not null, strlen(s), null) > 1")   // ok

API de conjuntos de datos tipados

Nota:

Esta característica se admite en clústeres habilitados para catálogos de Unity con el modo de acceso estándar en Databricks Runtime 15.4 y versiones posteriores.

Utiliza las API de Dataset tipadas para ejecutar transformaciones como map, filter y agregaciones en Datasets con una función definida por el usuario.

En el ejemplo siguiente se usa la map() API para modificar un número de una columna de resultado en una cadena con prefijo:

spark.range(3).map(f => s"row-$f").show()

En este ejemplo se usa map(), pero el mismo patrón se aplica a otras API de conjunto de datos con tipo, como filter(), mapPartitions(), foreach(), foreachPartition(), reduce()y flatMap().

Características de UDF de Scala y compatibilidad de Databricks Runtime

Las siguientes características requieren versiones mínimas de Databricks Runtime en clústeres habilitados para catálogos de Unity en modo de acceso estándar (compartido).

Característica Versión mínima de Databricks Runtime
UDF escalares Databricks Runtime 14.2
Dataset.map, Dataset.mapPartitions, Dataset.filter, , Dataset.reduce, Dataset.flatMap Databricks Runtime 15.4
KeyValueGroupedDataset.flatMapGroups, KeyValueGroupedDataset.mapGroups Databricks Runtime 15.4
(Transmisión en directo) foreachWriter Sink Databricks Runtime 15.4
(Transmisión en directo) foreachBatch Databricks Runtime 16.1
(Transmisión en directo) KeyValueGroupedDataset.flatMapGroupsWithState Databricks Runtime 16.2
spark.udf.registerJavaFunction (Java UDF desde un archivo JAR) Databricks Runtime 18 LTS