Pular para o conteúdo principal

UDFs Scala e Java com escopo de sessão

info

UDFs Scala e Java podem ser registradas no Unity Catalog para governança, reutilização e descoberta. Consulte funções definidas pelo usuário (UDFs) Scala e Java no Unity Catalog.

Esta página descreve como criar UDFs Scala e Java com escopo de sessão no Databricks. UDFs com escopo de sessão são definidas em um notebook ou job e se aplicam apenas à SparkSession atual. Para a referência da linguagem SQL, consulte Funções escalares definidas pelo usuário externas (UDFs).

Escolha sua abordagem

Você pode definir uma UDF Scala ou Java das seguintes maneiras. Para comparar todos os tipos de UDF em diferentes idiomas, governança e compute, consulte UDFs governadas por Unity Catalog vs. UDFs com escopo de sessão.

Abordagem

Descrição

UDF Scala Embutida

Defina uma UDF em um notebook usando uma função Scala ou lambda. Com escopo de sessão. Não suportado em compute serverless.

UDF Java de um JAR

Registre uma classe UDF pré-compilada de um JAR usando spark.udf.registerJavaFunction. Com escopo de sessão. Compatível com compute serverless.

UDF Scala ou Java governada pelo Unity Catalog

Registre uma UDF no Unity Catalog para governança, reuso e capacidade de descoberta. Compatível com compute serverless.

Abordagem

Descrição

UDF Scala Embutida

Defina uma UDF em um notebook usando uma função Scala ou lambda. Com escopo de sessão. Não suportado em compute serverless.

UDF Java de um JAR

Registre uma classe UDF pré-compilada de um JAR usando spark.udf.registerJavaFunction. Com escopo de sessão. Compatível com compute serverless.

UDF Scala ou Java governada pelo Unity Catalog

Registre uma UDF no Unity Catalog para governança, reuso e capacidade de descoberta. Compatível com compute serverless.

Requisitos

  • UDFs Scala em compute habilitado para Unity Catalog com modo de acesso padrão requerem Databricks Runtime 14.2 ou acima.
  • O suporte a instâncias ARM para UDFs Scala em clusters com Unity Catalog habilitado requer Databricks Runtime 15.2 ou superior.
  • O registro de um UDF Java de um JAR com spark.udf.registerJavaFunction requer o Databricks Runtime 18.3 ou acima. Consulte Registrar um UDF Java de um JAR.
importante

Compile seu JAR com as mesmas versões de Scala e Apache Spark do compute que o executa. Uma incompatibilidade pode fazer com que a UDF falhe no registro ou no momento da chamada.

Marque a dependência do Apache Spark como provided para que não seja empacotada em seu JAR. Inclua apenas as dependências de terceiros que seu UDF utiliza.

registrar uma função como um UDF

Registre uma função Scala como UDF usando spark.udf.register:

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

Chamar o UDF no Spark SQL

Crie uma view temporária e, em seguida, chame a UDF em uma query SQL:

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

Usar UDF com DataFrames

Você também pode chamar uma UDF usando a API do DataFrame:

Scala
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"))

Registrar uma UDF Java a partir de um JAR

Empacote uma UDF como um JAR, adicione-a à sua sessão com spark.addArtifact e registre a classe UDF com spark.udf.registerJavaFunction.

nota

Compatível no modo de acesso padrão e em compute Serverless no Databricks Runtime 18.3 ou acima. A função registrada tem escopo de sessão e não está registrada no Unity Catalog.

Os passos a seguir descrevem a criação de um projeto, a escrita de uma classe UDF, a construção de um JAR fat e o registro dele.

Passo 1: Criar seu projeto

Configure um projeto em Scala ou Java.

Crie um novo projeto Scala usando sbt:

Bash
sbt new scala/scala-seed.g8

Substitua o conteúdo do seu arquivo build.sbt pelo seguinte. Defina scalaVersion e a versão spark-sql para corresponder ao seu compute:

Scala
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 o plugin sbt-assembly para construir um JAR fat. Crie ou edite project/assembly.sbt e adicione:

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

Passo 2: Crie sua classe UDF

Sua classe UDF deve implementar uma das interfaces org.apache.spark.sql.api.java.UDF (UDF1 a UDF22), onde o número indica quantos argumentos de entrada a UDF recebe. Implemente o método call() com sua lógica.

O manipulador deve ser uma classe Java. spark.udf.registerJavaFunction carrega a classe por reflexão, portanto, deve ser uma classe pública de nível superior (ou aninhada static) com um construtor público sem argumentos. Uma Scala class ou object não satisfaz este requisito e falha no momento da chamada. Você pode criar o JAR com sbt, mas a própria classe UDF deve ser escrita em Java.

Criar src/main/java/com/example/MyIntegerUDF.java:

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;
}
}

Passo 3: Construa seu JAR fat

Empacote sua UDF compilada em um fat JAR.

No diretório raiz do seu projeto, execute:

Bash
sbt clean assembly

O JAR fat é criado em target/scala-2.13/ com um nome como my-udf-assembly-0.1.0-SNAPSHOT.jar.

O Passo 4: upload seu JAR para um volume do Unity Catalog

Faça upload do JAR para um volume do Unity Catalog para que o compute possa acessá-lo. Caso ainda não haja um volume, crie um:

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

Faça upload do seu arquivo JAR para o volume usando o Catalog Explorer:

  1. No seu workspace do Databricks, clique em Ícone de dados. Catálogo para abrir o Catalog Explorer.
  2. Selecione o catálogo e, em seguida, selecione o esquema que contém seu volume.
  3. Clique no nome do volume.
  4. Clique em Fazer upload para este volume e selecione seu arquivo JAR.
  5. Click upload .
  6. Após a conclusão do upload, clique no nome do seu arquivo JAR e, em seguida, clique em Copiar caminho para copiar o caminho do volume. Por exemplo, /Volumes/my_catalog/my_schema/udf_jars/my-udf-assembly-0.1.0-SNAPSHOT.jar. Você precisará deste caminho no próximo passo.

Passo 5: Registre e chame a UDF

Adicione o JAR à sua sessão usando seu caminho de volume, registre a classe UDF e chame-a do Spark SQL:

Python
# 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()

Em compute serverless e com modo de acesso padrão, é necessário informar um tipo de retorno explícito. Omitir o tipo de retorno falha com UC_COMMAND_NOT_SUPPORTED_IN_SHARED_ACCESS_MODE. Funções agregadas definidas pelo usuário (UDAFs) não têm suporte com registerJavaFunction.

A query retorna a saída da UDF, confirmando que a função está registrada e pode ser chamada:

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

Ordem de avaliação e verificação de nulos

O Spark SQL (incluindo SQL e as APIs DataFrame e Dataset) não garante a ordem de avaliação de subexpressões. O Spark não avalia as entradas de um operador ou função da esquerda para a direita. As expressões lógicas AND e OR não têm semântica de curto-circuito da esquerda para a direita.

Não dependa dos efeitos colaterais ou da ordem de avaliação das expressões booleanas, ou da ordem das cláusulas WHERE e HAVING. O otimizador de query pode reordenar essas expressões e cláusulas. Se uma UDF depender da semântica de curto-circuito para verificação de nulos, o Spark não garante que a verificação de nulos seja executada antes da UDF. Por exemplo:

Scala
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 cláusula WHERE não garante que o Spark invoca o UDF strlen depois de filtrar os nulos.

Para a verificação de valores nulos, a Databricks recomenda uma das seguintes opções:

  • Torne o UDF sensível a nulos e faça a verificação de nulos dentro do UDF.
  • Use as expressões IF ou CASE WHEN para fazer a verificação de nulidade e chamar o UDF em uma ramificação condicional
Scala
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

Conjunto de dados digitados APIs

nota

Esse recurso é suportado no clustering habilitado para o Unity Catalog com modo de acesso padrão em Databricks Runtime 15.4 e acima.

Use APIs de Dataset tipadas para executar transformações como mapear, filtrar e agregações em Datasets com uma função definida pelo usuário.

O exemplo a seguir usa a API map() para modificar um número em uma coluna de resultado para uma string prefixada:

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

Este exemplo usa map(), mas o mesmo padrão se aplica a outras APIs de dataset tipadas como filter(), mapPartitions(), foreach(), foreachPartition(), reduce() e flatMap().

Scala UDF Recurso e compatibilidade Databricks Runtime

Os seguintes recursos exigem versões mínimas do Databricks Runtime em clusters habilitados para Unity Catalog no modo de acesso padrão (compartilhado).

Recurso

Versão Mínima do Databricks Runtime

UDFs 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

(transmissão) foreachWriter Sink

Databricks Runtime 15.4

(transmissão) foreachBatch

Databricks Runtime 16.1

(transmissão) KeyValueGroupedDataset.flatMapGroupsWithState

Databricks Runtime 16.2

spark.udf.registerJavaFunction (UDF Java de um JAR)

Databricks Runtime 18.3

Recurso

Versão Mínima do Databricks Runtime

UDFs 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

(transmissão) foreachWriter Sink

Databricks Runtime 15.4

(transmissão) foreachBatch

Databricks Runtime 16.1

(transmissão) KeyValueGroupedDataset.flatMapGroupsWithState

Databricks Runtime 16.2

spark.udf.registerJavaFunction (UDF Java de um JAR)

Databricks Runtime 18.3