UDFs de Scala e Java com escopo de sessão

Importante

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

Esta página descreve como criar UDFs em Scala e Java com escopo de sessão no Azure Databricks. UDFs com escopo de sessão são definidos em um notebook ou trabalho e se aplicam somente ao SparkSession atual. Para obter a referência de linguagem SQL, consulte UDFs (funções escalares) definidas pelo usuário externo.

Escolha sua abordagem

Você pode definir um Scala ou Java UDF das seguintes maneiras. Para comparar todos os tipos de UDF em linguagens, governança e computação, consulte UDFs governadas pelo Catálogo do Unity vs. UDFs com escopo de sessão.

Approach Description
UDF Scala em linha Defina uma UDF em um notebook usando uma função Scala ou lambda. Escopo da sessão. Não há suporte na computação sem servidor.
Java UDF a partir de um JAR Registrar uma classe UDF pré-compilada de um JAR usando spark.udf.registerJavaFunction. Escopo da sessão. Com suporte na computação sem servidor.
UDFs em Scala ou Java governadas pelo Unity Catalog Registre uma UDF no Catálogo do Unity para governança, reutilização e capacidade de descoberta. Com suporte na computação sem servidor.

Requirements

  • As UDFs Scala na computação habilitada para o Catálogo do Unity com o modo de acesso padrão exigem o Databricks Runtime 14.2 ou posterior.
  • O suporte a instâncias do ARM em UDFs Scala em clusters habilitados para Catálogo do Unity requer o Databricks Runtime 15.2 ou posterior.
  • Registrar uma UDF Java de um JAR com spark.udf.registerJavaFunction requer o Databricks Runtime 18 LTS ou superior. Consulte Registrar uma UDF Java a partir de um JAR.

Importante

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

  • Computação clássica: combine as versões do Scala e do Spark com a sua versão do Databricks Runtime. Consulte a seção Ambiente de sistema das Notas de versão e compatibilidade do Databricks Runtime da sua versão. Por exemplo, o Databricks Runtime 18 LTS usa o Scala 2.13.16 e o Apache Spark 4.0.
  • Computação sem servidor: corresponda a versão do Scala à versão do ambiente. Consulte Versões do ambiente.

Marque a dependência do Apache Spark como provided para que ela não seja agrupada em seu JAR. Inclua apenas dependências de terceiros usadas pela sua UDF.

Registrar uma função como uma UDF

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

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

Chamar a UDF no Spark SQL

Crie uma exibição temporária e chame a UDF em uma consulta 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 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"))

Arquivos com UDF

Importante

Esse recurso está em Beta. Os administradores do workspace podem controlar o acesso a esse recurso na página Visualizações . Consulte Gerenciar visualizações do Azure Databricks.

O tipo Scala para um arquivo é FileRef. Use como parâmetro ou tipo de retorno em um UDF, tanto como tipo de nível superior quanto aninhado. Para o tipo e suas regras de aninhamento, consulte FILE tipo.

Para ler o conteúdo de um arquivo em um UDF, chame um dos seguintes em um FileRef:

  • asLocalFile(): Retorna um java.io.File que você pode passar para qualquer biblioteca que aceite um caminho.
  • open(): Retorna a java.io.InputStream para ler os bytes do arquivo. O chamador faz o fechamento.

Para produzir um novo FileRef em um UDF, chame um dos seguintes métodos estáticos:

  • FileRef.create(uri): Cria uma referência ao arquivo em uri.
  • FileRef.fromBytes(bytes, destinationPath, contentType): faz o upload bytes para o caminho do volume destinationPath como um arquivo externo e retorna uma referência.
  • FileRef.fromLocalFile(localFile, destinationPath, contentType): Faz upload de um arquivo local para o destinationPath caminho do Volume como um arquivo externo e retorna uma referência.

Retornar FileRef em uma UDF que grava em uma coluna FILE MANAGED não tem suporte.

Para exemplos em Python, Scala e SQL, incluindo processamento de imagens, detecção de tipos de arquivo e extração de quadros de vídeo, veja Arquivos de processo com UDFs.

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 da UDF com spark.udf.registerJavaFunction.

Observação

Com suporte no modo de acesso padrão e computação sem servidor no Databricks Runtime 18 LTS ou superior. A função registrada tem escopo de sessão e não está registrada no Catálogo do Unity.

As etapas a seguir explicam como criar um projeto, escrever uma classe UDF, criar um JAR gordo e registrá-lo.

Etapa 1: Criar seu projeto

Configure um projeto no Scala ou Java.

Scala

Crie um novo projeto Scala usando sbt:

sbt new scala/scala-seed.g8

Substitua o conteúdo do build.sbt arquivo pelo seguinte. Defina scalaVersion e a spark-sql versão para corresponder à sua computação:

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 plug-in sbt-assembly para criar um JAR gordo. Crie ou edite project/assembly.sbt e adicione:

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

Java

Crie um 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

Esse 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>

Etapa 2: escrever sua classe UDF

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

O manipulador deve ser uma classe Java. spark.udf.registerJavaFunction carrega a classe por reflexão, portanto, ela deve ser uma classe pública de nível superior (ou static aninhada) com um construtor público sem argumentos. Um Scala class ou object não atende a esse 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;
  }
}

Etapa 3: Criar seu fat JAR

Empacote o UDF compilado em um JAR gordo.

Scala

No diretório raiz do projeto, execute:

sbt clean assembly

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

Java

No diretório raiz do projeto, execute:

mvn clean package

O fat JAR é criado em target/ com um nome como my-udf-1.0-SNAPSHOT.jar.

Etapa 4: Faça upload do seu JAR para um volume do Unity Catalog

Carregue o JAR em um volume do Catálogo do Unity para que sua computação possa acessá-lo. Se você ainda não tiver um volume, crie um:

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

Carregue seu arquivo JAR no volume usando o Gerenciador de Catálogos:

  1. No workspace Azure Databricks, clique em Data icon.Catalog para abrir o Catalog Explorer.
  2. Selecione o catálogo e, em seguida, selecione o esquema que contém o volume.
  3. Clique no nome do volume.
  4. Clique em Fazer upload para este volume e selecione seu arquivo JAR.
  5. Clique em Carregar.
  6. Depois que o upload for concluído, clique no nome do arquivo JAR e 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ê precisa desse caminho na próxima etapa.

Etapa 5: Registrar e chamar a UDF

Adicione o JAR à sessão usando seu caminho de volume, registre a classe da 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 computação sem servidor e no modo de acesso padrão, você deve especificar um tipo de retorno explícito. A omissão do tipo de retorno causa falha com UC_COMMAND_NOT_SUPPORTED_IN_SHARED_ACCESS_MODE. Não há suporte para UDAFs (funções de agregação definidas pelo usuário) com registerJavaFunction.

A consulta retorna a saída UDF, confirmando que a função está registrada e que 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 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 dependa dos efeitos colaterais ou da ordem de avaliação das expressões boolianas ou da ordem das WHEREHAVING cláusulas. O otimizador de consulta pode reordenar essas expressões e cláusulas. Se uma UDF depender de semântica de curto-circuito para verificação nula, o Spark não garantirá que a verificação nula 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

Essa WHERE cláusula não garante que o Spark invoque a strlen UDF depois de filtrar os valores nulos.

Para lidar com a verificação nula, o Databricks recomenda um dos seguintes:

  • Fazer a própria UDF reconhecer nulos e verificar nulos dentro da UDF
  • Use expressões IF ou CASE WHEN para fazer a verificação de nulidade 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 do conjunto de dados tipado

Observação

Esse recurso tem suporte em clusters habilitados para Catálogo do Unity, no modo de acesso padrão, no Databricks Runtime 15.4 e superior.

Use APIs de conjunto de dados tipadas para executar transformações como mapa, filtro e agregações em conjuntos de dados com uma função definida pelo usuário.

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

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

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

Recursos do Scala UDF e compatibilidade do Databricks Runtime

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

Característica Versão mínima do Databricks Runtime
UDFs escalares Databricks Runtime 14.2
Dataset.map, Dataset.mapPartitions, Dataset.filter, , Dataset.reduceDataset.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 arquivo JAR) Databricks Runtime 18 LTS