UDFs em Scala e Java com âmbito de sessão

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.registerJavaFunction requer 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 a java.io.File que podes passar a qualquer biblioteca que aceite um caminho.
  • open(): Devolve a java.io.InputStream para 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 em uri.
  • FileRef.fromBytes(bytes, destinationPath, contentType): Carrega bytes para o destinationPath caminho de Volume como um ficheiro externo e retorna uma referência.
  • FileRef.fromLocalFile(localFile, destinationPath, contentType): Carrega um ficheiro local para o destinationPath caminho 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:

  1. No seu espaço de trabalho do Azure Databricks, clique no ícone Dados.Catálogo para abrir o Catalog Explorer.
  2. Selecione o catálogo e depois selecione o esquema que contém o seu volume.
  3. Clique no nome do volume.
  4. Clique em Carregar para este volume e selecione o seu ficheiro JAR.
  5. Clique em Carregar.
  6. 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 IF ou CASE WHEN expressõ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