UDFs Scala e Java com escopo de sessão
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 |
|---|---|
Defina uma UDF em um notebook usando uma função Scala ou lambda. Com escopo de sessão. Não suportado em compute serverless. | |
Registre uma classe UDF pré-compilada de um JAR usando | |
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.registerJavaFunctionrequer o Databricks Runtime 18.3 ou acima. Consulte Registrar um UDF Java de um JAR.
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.
- **Compute clássico**: Combine as versões Scala e Spark da sua versão do Databricks Runtime. Consulte a seção **Ambiente do sistema** das notas sobre a versão e compatibilidade do Databricks Runtime para sua versão. Por exemplo, o Databricks Runtime 18.3 usa Scala 2.13.16 e Apache Spark 4.0.
- Compute serverless : Corresponda à versão Scala da sua versão de ambiente. Consulte Versões de ambiente Serverless.
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:
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:
spark.range(1, 20).createOrReplaceTempView("test")
%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:
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.
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.
- Scala
- Java
Crie um novo projeto Scala usando sbt:
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:
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:
addSbtPlugin("com.eed3si9n" % "sbt-assembly" % "2.0.0")
Crie um novo projeto Maven usando o arquétipo quickstart:
mvn archetype:generate \
-DgroupId=com.example \
-DartifactId=my-udf \
-DarchetypeArtifactId=maven-archetype-quickstart \
-DinteractiveMode=false
Este comando cria a estrutura de projeto padrão do Maven com os diretórios src/main/java e src/test/java.
No pom.xml gerado, dentro das tags <project></project>, adicione um bloco <properties> e configure o maven-shade-plugin para gerar um fat JAR:
<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>
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:
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.
- Scala
- Java
No diretório raiz do seu projeto, execute:
sbt clean assembly
O JAR fat é criado em target/scala-2.13/ com um nome como my-udf-assembly-0.1.0-SNAPSHOT.jar.
No diretório raiz do seu projeto, execute:
mvn clean package
O JAR fat é criado em target/ com um nome como my-udf-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:
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:
- No seu workspace do Databricks, clique em
Catálogo para abrir o Catalog Explorer.
- Selecione o catálogo e, em seguida, selecione o esquema que contém seu volume.
- Clique no nome do volume.
- Clique em Fazer upload para este volume e selecione seu arquivo JAR.
- Click upload .
- 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:
# 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:
+----------+
| 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:
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
IFouCASE WHENpara fazer a verificação de nulidade e chamar o UDF em uma ramificação 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
Conjunto de dados digitados APIs
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:
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 |
| Databricks Runtime 15.4 |
| Databricks Runtime 15.4 |
(transmissão) | Databricks Runtime 15.4 |
(transmissão) | Databricks Runtime 16.1 |
(transmissão) | Databricks Runtime 16.2 |
| Databricks Runtime 18.3 |