Funções escalares definidas pelo usuário (UDFs) em Python

UDFs escalares em Python permitem executar lógica Python personalizada dentro de consultas SQL no Azure Databricks. Esta página mostra como registrá-los e invocá-los, usar credenciais e segredos de serviço, e lidar com ressalvas de avaliação de subexpressões no Spark SQL.

Requirements

  • No Databricks Runtime 12.2 LTS e versões anteriores, não há suporte para UDFs do Python e do Pandas no Catálogo do Unity na computação que usa o modo de acesso padrão.

  • As UDFs escalares do Python e as UDFs do Pandas têm suporte no Databricks Runtime 13.3 LTS e superior para todos os modos de acesso.

  • O suporte a instâncias do ARM em UDFs Python em clusters habilitados para Catálogo do Unity requer o Databricks Runtime 15.2 ou posterior.

No Databricks Runtime 14.0 e abaixo, UDFs do Python e UDFs do Pandas não têm suporte em clusters do Unity Catalog que usam o modo de acesso padrão. As UDFs escalares do Python e as UDFs do Pandas têm suporte no Databricks Runtime 14.1 e superior para todos os modos de acesso.

No Databricks Runtime 14.1 e superior, você pode registrar UDFs escalares do Python no Catálogo do Unity usando a sintaxe SQL. Consulte as funções definidas pelo usuário (UDFs) de SQL e Python no Unity Catalog.

Registrar uma função como uma UDF

def squared(s):
  return s * s
spark.udf.register("squaredWithPython", squared)

Opcionalmente, você pode definir o tipo de retorno de 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 a 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 o mesmo 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 os 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}}|
+-----------------------------+

Arquivos com UDF

Importante

Esse recurso está em Beta.

O tipo PySpark para um arquivo é FileType. Use como parâmetro ou tipo de retorno em um UDF, tanto como tipo de nível superior quanto aninhado. Para saber mais sobre o tipo, suas regras de aninhamento e a API FileRef, consulte FileType.

Para ler o conteúdo de um arquivo em um UDF, ligue file.as_local_file() para obter um caminho local que você possa abrir, ou file.open() para ler seus bytes como um fluxo. 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. Para a referência tipográfica, veja FILE tipo.

Um UDF registrado no Unity Catalog (CREATE FUNCTION) pode ler os metadados de um FILE', mas não seu conteúdo, e não pode criar arquivos. Use um UDF com escopo de sessão para ler o conteúdo de um arquivo (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 de nulos

O SQL do Spark (incluindo o SQL e o DataFrame e a API do Conjunto de Dados) não garante a ordem de avaliação de subexpressões. Em especial, 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 a semântica de “curto-circuito” da esquerda para a direita.

Portanto, é perigoso confiar nos efeitos colaterais ou na ordem de avaliação de expressões booleanas, e na ordem das cláusulas WHERE e HAVING, pois essas expressões e cláusulas podem ser reordenadas durante a otimização e o planejamento da consulta. Especificamente, se uma UDF depender de semântica de curto-circuito no SQL para verificação nula, não há nenhuma garantia de que a verificação nula ocorrerá antes de invocar a 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

Essa cláusula WHERE não garante que a UDF strlen seja invocada após a filtragem de nulos.

Para executar a verificação nula adequada, é recomendável que você faça o seguinte:

  • Fazer a própria UDF reconhecer nulos e verificar nulos dentro da própria UDF
  • Usar IF expressões ou CASE WHEN 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

Acessar segredos do Unity Catalog

Para acessar um segredo do Unity Catalog de um UDF Python com escopo de sessão, veja Usar um segredo em um UDF Python com escopo de sessão. Para acessar segredos declarados de uma UDF Python escalar ou de uma UDF Python em lote do Unity Catalog, consulte Usar segredos em uma UDF Python.

Credenciais de serviço em UDFs Python

UDFs escalares em Python com escopo de sessão e UDFs escalares em Python do Unity Catalog podem usar credenciais de serviço do Unity Catalog para acessar com segurança serviços externos em nuvem. Isso é útil para integrar operações como tokens baseados em nuvem, criptografia ou gerenciamento de segredos diretamente em suas transformações de dados.

Os requisitos variam conforme o tipo de UDF e o tipo de computação. Veja Usar uma credencial de serviço em um UDF Python.

Para criar uma credencial de serviço, consulte Criar credenciais de serviço.

Use uma credencial de serviço em um UDF escalar de Python com escopo 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 da 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

UDFs com escopo de sessão usam as permissões do chamador. Consulte Usar uma credencial de serviço em uma UDF em Python para ver os privilégios necessários.

Credenciais padrão em nível de computo para UDFs com escopo de sessão

Quando usado em UDFs escalares de Python, o Databricks usa automaticamente a credencial padrão do serviç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 credencial em seu código UDF. Consulte Especificar uma credencial de serviço padrão para um recurso de computação

O suporte à credencial padrão só está disponível em clusters de modo de acesso Standard e Dedicado. Não está disponível em warehouses SQL.

Você deve instalar o azure-identity pacote para usar o DefaultAzureCredential provedor. Para instalar o pacote, consulte bibliotecas Python com escopo de notebook 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 em um UDF escalar do Unity Catalog Python

Especifique a credencial de serviço na CREDENTIALS cláusula da definição do UDF. Você pode marcar uma credencial como DEFAULT para que SDKs de nuvem patchados a usem automaticamente. Na computação clássica, essa funcionalidade requer Databricks Runtime 18.1 ou superior. Em computação serverless e em SQL warehouses Pro e Serverless, defina explicitamente o(a) environment_version da UDF como 6 ou superior. Para requisitos completos de computação, rede e permissões, veja Usar uma credencial de serviço em um UDF Python.

Obter o contexto de execução da tarefa

Use a API do TaskContext PySpark para obter informações de contexto, como identidade do usuário, marcas de cluster, ID do trabalho do Spark e muito mais. Veja Como obter o contexto de tarefa em uma UDF.

Limitações

As seguintes limitações se aplicam a UDFs do PySpark:

  • Restrições de acesso a arquivos: No Databricks Runtime 14.2 e abaixo, UDFs do PySpark em clusters compartilhados não podem acessar pastas git, arquivos de workspace ou Volumes de Catálogo do Unity.

  • Variáveis de difusão: UDFs do PySpark em clusters de modo de acesso padrão e computação sem servidor não dão suporte a variáveis de difusão.

  • Limite de memória sem servidor: UDFs do PySpark na computação sem servidor têm um limite de memória de 1 GB por UDF do PySpark. Exceder esse limite resulta em um erro de tipo UDF_PYSPARK_USER_CODE_ERROR. MEMORY_LIMIT_SERVERLESS.
  • Limite de memória no modo de acesso padrão: UDFs do PySpark no modo de acesso padrão têm um limite de memória com base na memória disponível do tipo de instância escolhido. Exceder a memória disponível resulta em um erro de tipo UDF_PYSPARK_USER_CODE_ERROR. MEMORY_LIMIT.
  • Acesso à rede em SQL warehouses sem servidor: Por padrão, as UDFs do Python em SQL warehouses sem servidor não podem fazer solicitações de rede de saída, e as consultas que tentam fazer chamadas de rede ficam travadas indefinidamente. Para habilitar o acesso de saída à rede, habilite o recurso de Visualização Pública Habilitar rede para as workloads isoladas em SQL Warehouses sem servidor na página de Visualizações do seu espaço de trabalho. Caso contrário, use computação sem servidor ou computação clássica para UDFs que exigem acesso à rede.