Nota
O acesso a esta página requer autorização. Pode tentar iniciar sessão ou alterar os diretórios.
O acesso a esta página requer autorização. Pode tentar alterar os diretórios.
Importante
Os UDFs Scala e Java podem ser registados no Catálogo Unity para governação, reutilização e descoberta. Veja Scala e funções definidas pelo utilizador (UDFs) em Java no Catálogo Unity.
Esta página descreve como criar UDFs Scala e Java com âmbito de sessão no Azure Databricks. Os UDFs com âmbito de sessão são definidos num caderno ou trabalho e aplicam-se apenas ao SparkSession atual. Para a referência à linguagem SQL, veja Funções escalares definidas pelo utilizador externo (UDFs).
Escolha a sua abordagem
Pode definir um UDF Scala ou Java das seguintes formas. Para comparar todos os tipos de UDF ao nível das linguagens, da governação e da computação, veja UDFs governadas pelo Unity Catalog vs. UDFs com âmbito de sessão.
| Approach | Descrição |
|---|---|
| Inline Scala UDF | Defina um UDF num caderno usando uma função Scala ou lambda. Com âmbito de sessão. Não é suportado na computação sem servidor. |
| Java UDF a partir de um JAR | Registar uma classe UDF pré-compilada a partir de um JAR usando spark.udf.registerJavaFunction. Com âmbito de sessão. Suportado em computação sem servidor. |
| Scala governado pelo Unity Catalog ou Java UDF | Registe um UDF no Catálogo Unity para governação, reutilização e descoberta. Suportado em computação sem servidor. |
Requerimentos
- As UDFs de Scala em computação com o Unity Catalog ativado e com o modo de acesso padrão requerem o Databricks Runtime 14.2 ou posterior.
- O suporte a instâncias ARM para UDFs Scala em clusters com Unity Catalog requer Databricks Runtime 15.2 ou superior.
- Registar um UDF Java a partir de um JAR
spark.udf.registerJavaFunctionrequer Databricks Runtime 18 LTS ou superior. Veja Registar um UDF Java a partir de um JAR.
Importante
Constrói o teu JAR contra as mesmas versões Scala e Apache Spark que o computador que o executa. Uma incompatibilidade pode fazer com que a UDF falhe no registo ou na hora da chamada.
- Computação clássica: Compare as versões Scala e Spark da sua versão de Runtime Databricks. Consulte a secção de ambiente do sistema nas notas de versão do Databricks Runtime, versões e compatibilidade para a sua versão. Por exemplo, o Databricks Runtime 18 LTS utiliza o Scala 2.13.16 e o Apache Spark 4.0.
- Computação sem servidor: Faça corresponder a versão do Scala à versão do seu Ambiente. Consulte as versões do ambiente sem servidor .
Marca a dependência do Apache Spark como provided para que não seja incluída no teu JAR. Inclua apenas dependências de terceiros que o seu UDF utilize.
Registrar uma função como UDF
Registar uma função de Scala como um UDF usando spark.udf.register:
val squared = (s: Long) => {
s * s
}
spark.udf.register("square", squared)
Chamar o UDF no Spark SQL
Crie uma vista temporária e depois chame o UDF numa consulta SQL:
spark.range(1, 20).createOrReplaceTempView("test")
%sql select id, square(id) as id_squared from test
Usar UDF com DataFrames
Também pode chamar um UDF usando a 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"))
Ficheiros com a UDF
Importante
Este recurso está em versão Beta. Os administradores do espaço de trabalho podem controlar o acesso a esse recurso na página Visualizações . Ver Gerir as pré-visualizações de Azure Databricks.
O tipo de Scala para um ficheiro é FileRef. Use-o como parâmetro ou tipo de retorno numa UDF, seja como tipo de topo ou aninhado. Para o tipo e as suas regras de aninhamento, veja FILE tipo.
Para ler o conteúdo de um ficheiro numa UDF, chame um dos seguintes num FileRef:
-
asLocalFile(): Devolve ajava.io.Fileque podes passar a qualquer biblioteca que aceite um caminho. -
open(): Devolve ajava.io.InputStreampara ler os bytes do ficheiro. O interlocutor encerra-o.
Para produzir um novo FileRef numa UDF, chame um dos seguintes métodos estáticos:
-
FileRef.create(uri): Cria uma referência ao ficheiro emuri. -
FileRef.fromBytes(bytes, destinationPath, contentType): Carregabytespara odestinationPathcaminho de Volume como um ficheiro externo e retorna uma referência. -
FileRef.fromLocalFile(localFile, destinationPath, contentType): Carrega um ficheiro local para odestinationPathcaminho do Volume como ficheiro externo e devolve uma referência.
Devolver a FileRef de um UDF que escreve numa FILE MANAGED coluna não é suportado.
Para exemplos em Python, Scala e SQL, incluindo processamento de imagem, deteção de tipos de ficheiro e extração de fotogramas de vídeo, veja Processos de ficheiros com UDFs.
Registar um UDF Java a partir de um JAR
Empacota um UDF como JAR, adiciona-o à tua sessão com spark.addArtifact, e regista a classe UDF com spark.udf.registerJavaFunction.
Nota
Suportado em modo de acesso padrão e computação serverless no Databricks Runtime 18 LTS ou superior. A função registada tem âmbito de sessão e não está registada no Unity Catalog.
Os passos seguintes passam pela criação de um projeto, pela escrita de uma aula UDF, pela construção de um JAR gordo e pelo registo.
Passo 1: Crie o seu projeto
Configura um projeto em Scala ou Java.
Scala
Crie um novo projeto Scala usando sbt:
sbt new scala/scala-seed.g8
Substitua o conteúdo do arquivo build.sbt pelo seguinte. Defina scalaVersion e a spark-sql versão para corresponder ao seu cálculo:
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"
)
Ativa o plugin sbt-assembly para construir um JAR grosso. Criar ou editar project/assembly.sbt e adicionar:
addSbtPlugin("com.eed3si9n" % "sbt-assembly" % "2.0.0")
Java
Crie um novo projeto Maven usando o arquétipo de início rápido:
mvn archetype:generate \
-DgroupId=com.example \
-DartifactId=my-udf \
-DarchetypeArtifactId=maven-archetype-quickstart \
-DinteractiveMode=false
Este comando cria a estrutura padrão do projeto Maven com src/main/java diretórios e src/test/java diretórios.
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: Escreva a sua classe UDF
A sua classe UDF deve implementar uma das org.apache.spark.sql.api.java.UDF interfaces (UDF1 até UDF22), onde o número indica quantos argumentos de entrada a UDF recebe. Implementa o método call() com a tua lógica.
O handler deve ser uma classe Java.
spark.udf.registerJavaFunction carrega a classe por reflexão, pelo que deve ser uma classe pública de topo (ou static aninhada) com um construtor público no-arg. Um Scala class ou object não satisfaz este requisito e falha no momento da chamada. Podes construir o JAR com sbt, mas a própria classe UDF tem de 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: Constrói o teu JAR de gordura
Empacote o UDF compilado num fat JAR.
Scala
A partir do diretório raiz do seu projeto, execute:
sbt clean assembly
O fat JAR é criado em target/scala-2.13/ com um nome semelhante a my-udf-assembly-0.1.0-SNAPSHOT.jar.
Java
A partir do diretório raiz do seu projeto, execute:
mvn clean package
O fat JAR é criado em target/ com um nome semelhante a my-udf-1.0-SNAPSHOT.jar.
Passo 4: Carregue o seu JAR para um volume do Unity Catalog
Carrega o JAR para um volume do Unity Catalog para que o teu computador possa aceder a ele. Se ainda não tem um volume, crie um:
CREATE VOLUME IF NOT EXISTS my_catalog.my_schema.udf_jars
COMMENT 'Storage for UDF JAR files';
Carregue o seu ficheiro JAR para o volume usando o Explorador de Catálogos:
- No seu espaço de trabalho do Azure Databricks, clique no
Catálogo para abrir o Catalog Explorer.
- Selecione o catálogo e depois selecione o esquema que contém o seu volume.
- Clique no nome do volume.
- Clique em Carregar para este volume e selecione o seu ficheiro JAR.
- Clique em Carregar.
- Depois de o carregamento terminar, clique no nome do ficheiro JAR e depois 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. Precisas deste caminho no próximo passo.
Passo 5: Registe-se e ligue para a UDF
Adiciona o JAR à tua sessão usando o seu caminho de volume, regista a classe UDF e chama-a a partir 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()
Na computação sem servidor e na computação com modo de acesso padrão, deve passar um tipo de retorno explícito. A omissão do tipo de retorno falha com UC_COMMAND_NOT_SUPPORTED_IN_SHARED_ACCESS_MODE. Funções agregadas definidas pelo utilizador (UDAFs) não são suportadas com registerJavaFunction.
A consulta devolve a saída do UDF, confirmando que a função está registada e pode ser chamada:
+----------+
| my_udf(21)|
+----------+
| 22|
+----------+
Ordem de avaliação e verificação de nulidade
O SPARK SQL (incluindo SQL e as APIs DataFrame e Dataset) não garante a ordem de avaliação das 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 se baseie nos efeitos secundários nem na ordem de avaliação das expressões booleanas, nem na ordem das cláusulas WHERE e HAVING. O otimizador de consultas pode reordenar estas expressões e cláusulas. Se um UDF depende de semântica de curto-circuito para verificação de nulos, a Spark não garante que a verificação de nulo 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 WHERE cláusula não garante que o Spark invoque a strlen UDF depois de filtrar os nulos.
Para lidar com a verificação de valores nulos, a Databricks recomenda uma das seguintes opções:
- Fazer com que a própria UDF trate valores nulos e faça a verificação desses valores na própria UDF
- Use
IFouCASE WHENexpressões para fazer a verificação nula e invocar a 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
APIs de conjunto de dados tipadas
Nota
Esse recurso é suportado em clusters habilitados para Unity Catalog com modo de acesso padrão no Databricks Runtime 15.4 e superior.
Use APIs de Conjunto de Dados tipadas para executar transformações como mapeamento, filtro e agregações em Conjuntos de Dados com uma função definida pelo utilizador.
O exemplo seguinte usa a map() API para modificar um número numa coluna de resultados para uma cadeia prefixada:
spark.range(3).map(f => s"row-$f").show()
Este exemplo usa map(), mas o mesmo padrão aplica-se a outras APIs de Conjuntos de Dados tipadas, como filter(), mapPartitions(), foreach()foreachPartition(), reduce(), , e flatMap().
Recursos do Scala UDF e compatibilidade do Databricks Runtime
As seguintes funcionalidades requerem versões mínimas de Databricks Runtime em clusters habilitados pelo Unity Catalog em modo de acesso padrão (partilhado).
| Característica | Versão mínima do Databricks Runtime |
|---|---|
| Funções Definidas pelo Usuário (UDFs) escalares | Tempo de execução do Databricks 14.2 |
Dataset.map, Dataset.mapPartitions, Dataset.filter, Dataset.reduce, Dataset.flatMap |
Tempo de execução do Databricks 15.4 |
KeyValueGroupedDataset.flatMapGroups, KeyValueGroupedDataset.mapGroups |
Tempo de execução do Databricks 15.4 |
(Transmissão em fluxo) foreachWriter Sink |
Tempo de execução do Databricks 15.4 |
(Transmissão em fluxo) foreachBatch |
Tempo de execução do Databricks 16.1 |
(Transmissão em fluxo) KeyValueGroupedDataset.flatMapGroupsWithState |
Tempo de execução do Databricks 16.2 |
spark.udf.registerJavaFunction (Java UDF a partir de um JAR) |
Databricks Runtime 18 LTS |