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.
As funções definidas pelo utilizador (UDFs) permitem-lhe reutilizar e partilhar código que estende as capacidades incorporadas no Azure Databricks. Use UDFs para executar tarefas específicas, como cálculos complexos, transformações ou manipulações de dados personalizadas.
Quando usar uma função UDF vs. Apache Spark?
Use UDFs para lógica difícil de expressar com funções integradas do Apache Spark. As funções integradas do Apache Spark são otimizadas para processamento distribuído e oferecem melhor desempenho em escala. Para obter mais informações, consulte Functions.
A Databricks recomenda UDFs para consultas ad hoc, limpeza manual de dados, análise exploratória de dados e operações em conjuntos de dados de pequeno a médio porte. Os casos de uso comuns para UDFs incluem criptografia de dados, descriptografia, hashing, análise JSON e validação.
Use métodos Apache Spark para operações em grandes conjuntos de dados e quaisquer cargas de trabalho que corram regularmente ou continuamente, incluindo trabalhos ETL e operações de streaming.
Compreender os tipos de UDF
Selecione um tipo UDF nas guias a seguir para ver uma descrição, um exemplo e um link para saber mais.
UDF escalar
UDFs escalares operam em uma única linha e retornam um único valor de resultado para cada linha. Eles podem ser governados pelo Unity Catalog ou com escopo de sessão.
O exemplo a seguir usa um UDF escalar para calcular o comprimento de cada nome em uma name coluna e adicionar o valor em uma nova coluna name_length.
+-------+-------+
| name | score |
+-------+-------+
| alice | 10.0 |
| bob | 20.0 |
| carol | 30.0 |
| dave | 40.0 |
| eve | 50.0 |
+-------+-------+
-- Create a SQL UDF for name length
CREATE OR REPLACE FUNCTION main.test.get_name_length(name STRING)
RETURNS INT
RETURN LENGTH(name);
-- Use the UDF in a SQL query
SELECT name, main.test.get_name_length(name) AS name_length
FROM your_table;
+-------+-------+-------------+
| name | score | name_length |
+-------+-------+-------------+
| alice | 10.0 | 5 |
| bob | 20.0 | 3 |
| carol | 30.0 | 5 |
| dave | 40.0 | 4 |
| eve | 50.0 | 3 |
+-------+-------+-------------+
Para implementar isto num notebook Azure Databricks usando PySpark:
from pyspark.sql.functions import udf
from pyspark.sql.types import IntegerType
@udf(returnType=IntegerType())
def get_name_length(name):
return len(name)
df = df.withColumn("name_length", get_name_length(df.name))
# Show the result
display(df)
Veja funções definidas pelo utilizador (UDFs) SQL e Python no Unity Catalog e funções definidas pelo utilizador (UDFs) escalares em Python.
UDFs escalares em lote
Processe dados em lotes, mantendo a paridade de linha de entrada/saída 1:1. Isso reduz a sobrecarga de operações linha a linha para processamento de dados em grande escala. As UDFs de lote também mantêm o estado entre lotes para serem executadas com mais eficiência, reutilizar recursos e lidar com cálculos complexos que precisam de contexto em blocos de dados.
Eles podem ser governados pelo Unity Catalog ou com escopo de sessão.
O seguinte Batch Unity Catalog Python UDF calcula o IMC durante o processamento de lotes de linhas:
+-------------+-------------+
| weight_kg | height_m |
+-------------+-------------+
| 90 | 1.8 |
| 77 | 1.6 |
| 50 | 1.5 |
+-------------+-------------+
%sql
CREATE OR REPLACE FUNCTION main.test.calculate_bmi_pandas(weight_kg DOUBLE, height_m DOUBLE)
RETURNS DOUBLE
LANGUAGE PYTHON
PARAMETER STYLE PANDAS
HANDLER 'handler_function'
AS $$
import pandas as pd
from typing import Iterator, Tuple
def handler_function(batch_iter: Iterator[Tuple[pd.Series, pd.Series]]) -> Iterator[pd.Series]:
for weight_series, height_series in batch_iter:
yield weight_series / (height_series ** 2)
$$;
select main.test.calculate_bmi_pandas(cast(70 as double), cast(1.8 as double));
+--------+
| BMI |
+--------+
| 27.8 |
| 30.1 |
| 22.2 |
+--------+
Veja as funções SQL e Python definidas pelo utilizador (UDFs) no Catálogo Unity e as funções Python definidas pelo utilizador (UDFs) por lotes no Catálogo Unity.
UDFs não escalares
UDFs não escalares operam em conjuntos de dados/colunas inteiros com relações de entrada/saída flexíveis (1:N ou muitos:muitos).
Os UDFs de pandas em lote com escopo de sessão podem ser dos seguintes tipos:
- De Série Para Série
- Iterador de Série para Iterador de Série
- Iterador de múltiplas séries para iterador de uma série
- Conversão de série para escalar
Segue-se um exemplo de um pandas UDF de Série para Série.
from pyspark.sql.functions import pandas_udf
import pandas as pd
df = spark.createDataFrame([(70, 1.75), (80, 1.80), (60, 1.65)], ["Weight", "Height"])
@pandas_udf("double")
def calculate_bmi_pandas(weight: pd.Series, height: pd.Series) -> pd.Series:
return weight / (height ** 2)
df.withColumn("BMI", calculate_bmi_pandas(df["Weight"], df["Height"])).display()
Consulte as funções definidas pelo utilizador do pandas.
UDAF
As UDAFs operam em várias linhas e retornam um único resultado agregado. As UDAFs estão limitadas ao âmbito da sessão.
O exemplo UDAF a seguir agrega pontuações por comprimento de nome.
from pyspark.sql.functions import pandas_udf
from pyspark.sql import SparkSession
import pandas as pd
# Define a pandas UDF for aggregating scores
@pandas_udf("int")
def total_score_udf(scores: pd.Series) -> int:
return scores.sum()
# Group by name length and aggregate
result_df = (df.groupBy("name_length")
.agg(total_score_udf(df["score"]).alias("total_score")))
display(result_df)
+-------------+-------------+
| name_length | total_score |
+-------------+-------------+
| 3 | 70.0 |
| 4 | 40.0 |
| 5 | 40.0 |
+-------------+-------------+
Consulte funções definidas pelo utilizador do pandas para Python e funções agregadas definidas pelo utilizador em Scala (UDAFs).
Funções de Tabela Definidas pelo Utilizador (UDTFs)
Um UDTF usa um ou mais argumentos de entrada e retorna várias linhas (e possivelmente várias colunas) para cada linha de entrada. Eles podem ser governados pelo Unity Catalog ou com escopo de sessão.
O UDTF a seguir cria uma tabela usando uma lista fixa de dois argumentos inteiros:
CREATE OR REPLACE FUNCTION get_sum_diff(x INT, y INT)
RETURNS TABLE (sum INT, diff INT)
LANGUAGE PYTHON
HANDLER 'GetSumDiff'
AS $$
class GetSumDiff:
def eval(self, x: int, y: int):
yield x + y, x - y
$$;
SELECT * FROM get_sum_diff(10, 3);
+-----+------+
| sum | diff |
+-----+------+
| 13 | 7 |
+-----+------+
Para implementar isto num notebook Azure Databricks usando PySpark:
from pyspark.sql.functions import lit, udtf
@udtf(returnType="sum: int, diff: int")
class GetSumDiff:
def eval(self, x: int, y: int):
yield x + y, x - y
GetSumDiff(lit(1), lit(2)).show()
Consulte UDTFs do Unity Catalog e UDTFs com escopo específico de sessão.
UDFs governados pelo Unity Catalog vs. UDFs com âmbito de sessão
O Unity Catalog armazena UDFs governados pelo Unity Catalog para melhorar a governação, a reutilização e a facilidade de descoberta. Define UDFs com âmbito da sessão num notebook ou tarefa, limitadas à SparkSession atual. Pode definir e aceder a UDFs com âmbito de sessão usando SQL, Python ou Scala.
Utilize a tabela seguinte para decidir entre as duas categorias e, em seguida, consulte os guias de consulta rápida apresentados em seguida para os tipos específicos de UDF de cada uma.
| Consideração | O Catálogo Unity governava os UDFs | UDFs com âmbito de sessão |
|---|---|---|
| Melhor para | Partilhar funções com segurança entre equipas, notebooks, tarefas e armazéns de dados SQL. | Desenvolvimento rápido e iterativo dentro de um único caderno ou tarefa. |
| Languages | SQL, Python, Scala e Java. | SQL, Python e Scala. |
| Governação e partilha | Governado pelas permissões do Catálogo Unity e detectável no Explorador de Catálogos. | Limitado à SparkSession atual. Não governado ou partilhado. |
| Persistência | Persistiu no Unity Catalog e foi reutilizável ao longo das sessões. | Existe apenas para a sessão atual. |
Folha de dicas de UDFs geridas pelo Unity Catalog
As UDFs governadas pelo Unity Catalog permitem que funções personalizadas sejam definidas, usadas, compartilhadas com segurança e controladas em ambientes de computação. Consulte funções definidas pelo utilizador (UDFs) em SQL e Python no Unity Catalog.
| Tipo UDF | Computação suportada | Descrição |
|---|---|---|
| Unity Catálogo Python UDF |
|
Defina um UDF em Python e registre-o no Unity Catalog para governança. UDFs escalares operam em uma única linha e retornam um único valor de resultado para cada linha. |
| Batch Unity Catálogo Python UDF |
|
Defina um UDF em Python e registre-o no Unity Catalog para governança. Realizar operações em lote com múltiplos valores e devolver múltiplos resultados. Reduz a sobrecarga das operações realizadas linha por linha para o processamento de dados em grande escala. |
| Catálogo Unity Python UDTF |
|
Defina um UDTF em Python e registre-o no Unity Catalog para governança. Um UDTF usa um ou mais argumentos de entrada e retorna várias linhas (e possivelmente várias colunas) para cada linha de entrada. |
| Unity Catalog Scala ou Java UDF |
|
Defina um UDF em Scala ou Java e registe-o no Unity Catalog para governação. UDFs escalares operam em uma única linha e retornam um único valor de resultado para cada linha. Requer Scala 2.13.16, JDK 17 e Environment Version 4. |
Folha de truques UDFs com escopo de sessão para computação isolada do usuário
Define UDFs com âmbito da sessão num notebook ou trabalho, limitadas à SparkSession atual. Pode definir e aceder a UDFs com âmbito de sessão usando SQL, Python ou Scala.
| Tipo UDF | Computação suportada | Descrição |
|---|---|---|
| Python escalar |
|
UDFs escalares operam em uma única linha e retornam um único valor de resultado para cada linha. |
| Python não-escalar |
|
UDFs não escalares incluem pandas_udf, mapInPandas, mapInArrow, applyInPandas. Os Pandas UDFs usam a Seta Apache para transferir dados e os pandas para trabalhar com os dados. Os Pandas UDFs suportam operações vetorizadas que podem aumentar consideravelmente o desempenho em UDFs escalares linha a linha. |
| UDTFs de Python |
|
Um UDTF usa um ou mais argumentos de entrada e retorna várias linhas (e possivelmente várias colunas) para cada linha de entrada. |
| UDFs escalares Scala |
|
UDFs escalares operam em uma única linha e retornam um único valor de resultado para cada linha. |
| Scala ou Java UDF a partir de um JAR |
|
Registar uma classe UDF pré-compilada a partir de um JAR usando spark.udf.registerJavaFunction. Veja Registar um UDF Java a partir de um JAR. |
| Scala UDAFs |
|
As UDAFs operam em várias linhas e retornam um único resultado agregado. |
Considerações sobre desempenho
Funções internas e UDFs SQL são as opções mais eficientes.
UDFs Scala são geralmente mais rápidos do que UDFs Python.
- UDFs Scala não isoladas são executadas na Java Virtual Machine (JVM), evitando a sobrecarga de mover dados para dentro e para fora da JVM.
- Os UDFs Scala isolados têm de mover dados para dentro e fora da JVM, mas ainda assim podem ser mais rápidos do que os UDFs Python porque gerem a memória de forma mais eficiente.
UDFs Python e UDFs pandas tendem a ser mais lentos do que os UDFs Scala porque têm de serializar os dados e movê-los da JVM para o interpretador Python.
- As UDFs Pandas são até 100x mais rápidas que as UDFs Python porque usam a Seta Apache para reduzir os custos de serialização.