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.
Os UDFs escalares em Python permitem executar lógica Python personalizada dentro de consultas SQL no Azure Databricks. Esta página mostra como os registar e invocar, usar credenciais e segredos de serviço, e lidar com as ressalvas da ordem de avaliação das subexpressões no Spark SQL.
Requerimentos
No Databricks Runtime 12.2 LTS e versões anteriores, UDFs Python e UDFs Pandas não são suportados no Unity Catalog que utiliza o modo de acesso padrão na computação.
Os UDFs Python escalares e os UDFs Pandas são suportados no Databricks Runtime 13.3 LTS e versões superiores para todos os modos de acesso.
O suporte a instâncias ARM para UDFs Python em clusters ativados com Unity Catalog requer Databricks Runtime 15.2 ou superior.
No Databricks Runtime 14.0 e inferiores, os UDFs Python e os UDFs Pandas não têm suporte em clusters do Unity Catalog que usam o modo de acesso padrão. As UDFs Scalar Python e as UDFs Pandas são suportadas para todos os modos de acesso no Databricks Runtime 14.1 e superiores.
No Databricks Runtime 14.1 e superiores, pode registar UDFs escalares em Python no Unity Catalog usando sintaxe SQL. Consulte funções definidas pelo utilizador (UDFs) em SQL e Python no Unity Catalog.
Registrar uma função como UDF
def squared(s):
return s * s
spark.udf.register("squaredWithPython", squared)
Opcionalmente, você pode definir o tipo de retorno do seu UDF. O tipo de retorno padrão é StringType.
from pyspark.sql.types import LongType
def squared_typed(s):
return s * s
spark.udf.register("squaredWithPython", squared_typed, LongType())
Chamar o UDF no Spark SQL
spark.range(1, 20).createOrReplaceTempView("test")
%sql select id, squaredWithPython(id) as id_squared from test
Usar UDF com DataFrames
from pyspark.sql.functions import udf
from pyspark.sql.types import LongType
squared_udf = udf(squared, LongType())
df = spark.table("test")
display(df.select("id", squared_udf("id").alias("id_squared")))
Como alternativa, você pode declarar a mesma UDF usando a sintaxe de anotação:
from pyspark.sql.functions import udf
@udf("long")
def squared_udf(s):
return s * s
df = spark.table("test")
display(df.select("id", squared_udf("id").alias("id_squared")))
Variantes com UDF
O tipo PySpark para variante é VariantType e os valores são do tipo VariantVal. Para obter informações sobre variantes, consulte Consultar dados de variantes.
from pyspark.sql.functions import col, lit, udf
from pyspark.sql.types import VariantType, VariantVal
# Return Variant
@udf(returnType = VariantType())
def toVariant(jsonString):
return VariantVal.parseJson(jsonString)
spark.range(1).select(lit('{"a" : 1}').alias("json")).select(toVariant(col("json"))).display()
+---------------+
|toVariant(json)|
+---------------+
| {"a":1}|
+---------------+
from pyspark.sql.functions import col, lit, udf
from pyspark.sql.types import StructField, StructType, VariantType, VariantVal
# Return Struct<Variant>
@udf(returnType = StructType([StructField("v", VariantType(), True)]))
def toStructVariant(jsonString):
return {"v": VariantVal.parseJson(jsonString)}
spark.range(1).select(lit('{"a" : 1}').alias("json")).select(toStructVariant(col("json"))).display()
+---------------------+
|toStructVariant(json)|
+---------------------+
| {"v":{"a":1}}|
+---------------------+
from pyspark.sql.functions import col, lit, udf
from pyspark.sql.types import ArrayType, VariantType, VariantVal
# Return Array<Variant>
@udf(returnType = ArrayType(VariantType()))
def toArrayVariant(jsonString):
return [VariantVal.parseJson(jsonString)]
spark.range(1).select(lit('{"a" : 1}').alias("json")).select(toArrayVariant(col("json"))).display()
+--------------------+
|toArrayVariant(json)|
+--------------------+
| [{"a":1}]|
+--------------------+
from pyspark.sql.functions import col, lit, udf
from pyspark.sql.types import MapType, StringType, VariantType, VariantVal
# Return Map<String, Variant>
@udf(returnType = MapType(StringType(), VariantType(), True))
def toMapVariant(jsonString):
return {"v1": VariantVal.parseJson(jsonString), "v2": VariantVal.parseJson("[" + jsonString + "]")}
spark.range(1).select(lit('{"a" : 1}').alias("json")).select(toMapVariant(col("json"))).display()
+-----------------------------+
| toMapVariant(json)|
+-----------------------------+
|{"v2":[{"a":1}],"v1":{"a":1}}|
+-----------------------------+
Ficheiros com a UDF
Importante
Este recurso está em versão Beta.
O tipo PySpark para um ficheiro é FileType. Utilize-o como parâmetro ou tipo de retorno numa UDF, quer como tipo de nível superior quer aninhado. Para o tipo, as suas regras de aninhamento e a API FileRef, consulte FileType.
Para ler o conteúdo de um ficheiro numa UDF, ligue file.as_local_file() para obter um caminho local que possa abrir ou file.open() para ler os seus bytes como um fluxo. 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. Para a referência tipográfica, veja FILE tipo.
Um UDF registado no Unity Catalog (CREATE FUNCTION) pode ler os metadados de um FILE', mas não o seu conteúdo, e não pode criar ficheiros. Use um UDF com âmbito de sessão para ler o conteúdo de um ficheiro (open, as_local_file) ou criar um (from_bytes, from_local_file).
No corpo da UDF, cada FILE valor é um FileRef objeto:
from pyspark.sql.functions import col, udf
from pyspark.sql.types import FileRef, StringType
@udf(returnType=StringType())
def file_content_type(file: FileRef) -> str:
return file.content_type
df = spark.table("documents")
display(df.select(col("file").uri, file_content_type(col("file"))))
Ordem de avaliação e verificação nula
O Spark SQL (incluindo SQL e a API DataFrame e Dataset) não garante a ordem de avaliação das subexpressões. Em particular, as entradas de um operador ou função não são necessariamente avaliadas da esquerda para a direita ou em qualquer outra ordem fixa. Por exemplo, as expressões lógicas AND e OR não têm semântica de "curto-circuito" ao serem avaliadas da esquerda para a direita.
Portanto, é perigoso confiar nos efeitos colaterais ou na ordem de avaliação das expressões booleanas, e na ordem das cláusulas WHERE e HAVING, uma vez que tais expressões e cláusulas podem ser reordenadas durante a otimização e o planejamento da consulta. Especificamente, se um UDF depende de semântica de curto-circuito no SQL para verificação nula, não há garantia de que a verificação nula acontecerá antes de invocar o UDF. Por exemplo,
spark.udf.register("strlen", lambda s: len(s), "int")
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 strlen UDF seja invocado após a filtragem dos nulos.
Para executar a verificação nula adequada, recomendamos que você siga um destes procedimentos:
- Tornar a UDF capaz de lidar com valores nulos e fazer a verificação de nulos dentro da 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", lambda s: len(s) if not s is None else -1, "int")
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
Aceder aos segredos do Unity Catalog
Para aceder a um segredo do Unity Catalog de um UDF Python com âmbito de sessão, veja Usar um segredo num UDF Python com âmbito de sessão. Para aceder a segredos declarados de um UDF Python escalar ou Batch Unity Catalog, veja Usar segredos num UDF Python.
Credenciais de serviço em UDFs Python
As UDFs escalares Python com âmbito de sessão e as UDFs escalares Python do Unity Catalog podem usar credenciais de serviço do Unity Catalog para aceder de forma segura a serviços externos na cloud. Isso é útil para integrar operações como tokenização baseada em nuvem, criptografia ou gerenciamento de segredos diretamente em suas transformações de dados.
Os requisitos variam consoante o tipo de UDF e o tipo de computação. Veja Usar uma credencial de serviço num UDF Python.
Para criar uma credencial de serviço, consulte Criar credenciais de serviço.
Use uma credencial de serviço num UDF escalar Python com âmbito de sessão
Para acessar a credencial de serviço, use o databricks.service_credentials.getServiceCredentialsProvider() utilitário em sua lógica UDF para inicializar SDKs de nuvem com a credencial apropriada. Todo o código deve ser encapsulado no corpo UDF.
@udf
def use_service_credential():
from azure.mgmt.web import WebSiteManagementClient
# Assuming there is a service credential named 'testcred' set up in Unity Catalog
web_client = WebSiteManagementClient(subscription_id, credential = getServiceCredentialsProvider('testcred'))
# Use web_client to perform operations
Permissões de credenciais de serviço
Os UDFs com âmbito de sessão utilizam as permissões do chamador. Veja Usar uma credencial de serviço num UDF Python para obter os privilégios necessários.
Credenciais predefinidas ao nível de computação para UDFs com âmbito de sessão
Quando usado em UDFs escalares de Python, o Databricks utiliza automaticamente a credencial de serviço padrão da variável do ambiente de computação. Esse comportamento permite que você faça referência segura a serviços externos sem gerenciar explicitamente aliases de credenciais em seu código UDF. Consulte Especificar uma credencial de serviço padrão para um recurso de computação
O suporte a credenciais padrão só está disponível em clusters de modo de acesso Padrão e Dedicado. Não está disponível em armazéns SQL.
Você deve instalar o azure-identity pacote para usar o DefaultAzureCredential provedor. Para instalar o pacote, consulte Bibliotecas Python com escopo de bloco de anotações ou Bibliotecas com escopo de computação.
@udf
def use_service_credential():
from azure.identity import DefaultAzureCredential
from azure.mgmt.web import WebSiteManagementClient
# DefaultAzureCredential is automatically using the default service credential for the compute
web_client_default = WebSiteManagementClient(DefaultAzureCredential(), subscription_id)
# Use web_client to perform operations
Use uma credencial de serviço num UDF escalar do Unity Catalog Python
Especifique a credencial de serviço na CREDENTIALS cláusula da definição do UDF. Podes marcar uma credencial como DEFAULT para que os SDKs cloud atualizados a usem automaticamente. Na computação clássica, esta funcionalidade requer Databricks Runtime 18.1 ou superior. Em processamento sem servidor e em SQL warehouses Pro e Serverless, defina explicitamente o environment_version da UDF como 6 ou superior. Para requisitos completos de computação, redes e permissões, consulte Usar uma credencial de serviço num UDF Python.
Obter contexto de execução da tarefa
Utilize a API PySpark do TaskContext para obter informações de contexto, como identidade do utilizador, tags do cluster, ID do trabalho Spark e mais. Veja Obter o contexto da tarefa num UDF.
Limitações
As seguintes limitações se aplicam aos UDFs do PySpark:
Restrições de acesso a ficheiros: No Databricks Runtime 14.2 e inferior, as UDFs do PySpark em clusters compartilhados não podem acessar pastas Git, arquivos de espaço de trabalho ou volumes de catálogo Unity.
Variáveis de transmissão: UDFs do PySpark em clusters de modo de acesso padrão e computação sem servidor não suportam variáveis de difusão.
- Limite de memória em serverless: PySpark UDFs em computação sem servidor têm um limite de memória de 1GB por PySpark UDF. Exceder esse limite resulta em um erro do tipo UDF_PYSPARK_USER_CODE_ERROR. MEMORY_LIMIT_SERVERLESS.
- Limite de memória no modo de acesso padrão: UDFs PySpark no modo de acesso padrão têm um limite de memória baseado na memória disponível do tipo de instância escolhido. Exceder a memória disponível resulta em um erro do tipo UDF_PYSPARK_USER_CODE_ERROR. MEMORY_LIMIT.
- Acesso à rede em armazéns de SQL sem servidor: Por defeito, as UDFs de Python em armazéns de SQL sem servidor não podem fazer pedidos de saída para a rede, e as consultas que tentam efetuar chamadas de rede ficam suspensas indefinidamente. Para ativar o acesso de saída à rede, ative a funcionalidade de Pré-visualização Pública Ativar a rede para cargas de trabalho isoladas em Serverless SQL Warehouses na página Pré-visualizações do seu espaço de trabalho. Caso contrário, use computação serverless ou computação clássica para UDFs que exijam acesso à rede.