Pular para o conteúdo principal

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

As UDFs Python escalares permitem execução lógica Python personalizada dentro de query SQL no Databricks. Esta página mostra como registrá-las e chamá-las, usar credenciais e segredos de serviço, e lidar com ressalvas de ordem de avaliação de subexpressões no Spark SQL.

Requisitos​

  • No Databricks Runtime 12.2 LTS e versões anteriores, as UDFs Python e Pandas não são suportadas no Unity Catalog compute que utiliza o modo de acesso padrão.

  • As UDFs escalares Python e as UDFs Pandas são suportadas no Databricks Runtime 13.3 LTS e versões superiores para todos os modos de acesso.

  • O suporte da instância Graviton para UDFs Python em clusters com Unity Catalog habilitado requer Databricks Runtime 15.2 ou superior.

Em Databricks Runtime 14.0 e abaixo, Python UDFs e Pandas UDFs não são suportados em Unity Catalog clustering que usam o modo de acesso padrão. Os UDFs escalares Python e Pandas são compatíveis com todos os modos de acesso em Databricks Runtime 14.1 e acima.

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

registrar uma função como um UDF​

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

Opcionalmente, o senhor pode definir o tipo de retorno do seu UDF. O tipo de retorno do default é StringType.

Python
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​

Python
spark.range(1, 20).createOrReplaceTempView("test")
SQL
%sql select id, squaredWithPython(id) as id_squared from test

Usar UDF com DataFrames​

Python
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, o senhor pode declarar o mesmo UDF usando a sintaxe de anotação:

Python
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 a variante é VariantType e os valores são do tipo VariantVal. Para obter informações sobre variantes, consulte Consultar dados de variantes.

Python
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}|
+---------------+
Python
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}}|
+---------------------+
Python
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}]|
+--------------------+
Python
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​

info

Beta

Este recurso está em Beta.

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

Para ler o conteúdo de um arquivo em uma UDF, chame file.as_local_file() para obter um caminho local que você possa abrir, ou file.open() para ler seus bytes como uma transmissão. Para obter exemplos em Python, Scala e SQL, incluindo processamento de imagem, detecção de tipo de arquivo e extração de quadros de vídeo, consulte Processar arquivos com UDFs. Para a referência de tipo, consulte tipoFILE.

Uma UDF registrada no Unity Catalog (CREATE FUNCTION) pode ler os metadados de um FILE, mas não o conteúdo, e não pode criar arquivos. Use uma 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 do UDF, cada valor FILE é um objeto FileRef:

Python
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​

Spark SQL (incluindo SQL e DataFrame e o conjunto de dados API) 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 avaliado da esquerda para a direita ou em qualquer outra ordem fixa. Por exemplo, lógico AND e as expressões OR não têm semântica de “curto-circuito” da esquerda para a direita.

Portanto, é perigoso confiar nos efeitos colaterais ou na ordem de avaliação das expressões Boolean 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 um UDF depender da semântica de curto-circuito no SQL para verificação de nulidade, não há garantia de que a verificação de nulidade ocorrerá antes de invocar o UDF. Por exemplo,

Python
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 o UDF strlen seja invocado após a filtragem de nulos.

Para realizar uma verificação de nulo adequada, recomendamos que você faça o seguinte:

  • Tornar o próprio UDF sensível a nulidade e fazer a verificação de nulidade dentro do próprio UDF
  • Use as expressões IF ou CASE WHEN para fazer a verificação de nulidade e chamar o UDF em uma ramificação condicional
Python
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

Access Unity Catalog secrets​

Para acessar um secret do Unity Catalog a partir de um Python UDF com escopo de sessão, consulte Use a secret in a session-scoped Python UDF. To access declared secrets from a scalar or lotes Unity Catalog Python UDF, see Use secrets in a Python UDF.

Credenciais de serviço em UDFs do Python​

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

Os requisitos variam de acordo com o tipo de UDF e o compute. Consulte Usar uma credencial de serviço em uma UDF Python.

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

Usar uma credencial de serviço em uma UDF Python escalar com escopo de sessão​

Para acessar a credencial do serviço, utilize as utilidades databricks.service_credentials.getServiceCredentialsProvider() em sua lógica UDF para inicializar os SDKs da nuvem com a credencial apropriada. Todo o código deve ser encapsulado no corpo da UDF.

Python
@udf
def use_service_credential():
from databricks.service_credentials import getServiceCredentialsProvider
import boto3

# Assuming there is a service credential named 'testcred' set up in Unity Catalog
boto3_session = boto3.Session(botocore_session=getServiceCredentialsProvider('testcred'))
# Use the S3 session 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 Python para obter os privilégios necessários.

Credenciais default em nível de compute para UDFs com escopo de sessão​

Quando usadas em UDFs Python escalares, o Databricks usa automaticamente a credencial de serviço default da variável de ambiente de compute. Esse comportamento permite que você faça referência segura a serviços externos sem gerenciar explicitamente aliases de credenciais no seu código UDF. Consulte Especificar uma credencial de serviço default para um recurso de compute

O suporte a credenciais default está disponível apenas em clusters com modo de acesso Standard e Dedicated. Ele não está disponível em SQL warehouses.

Python
@udf
def use_service_credential():
from databricks.service_credentials import getServiceCredentialsProvider
import boto3

# The default service credential for the compute is automatically used
boto3_session = boto3.Session()
# Use the S3 client to perform operations

Exemplo de função AWS Lambda de UDF em Python com escopo de sessão​

O exemplo a seguir usa uma credencial de serviço para chamar uma função do AWS Lambda a partir de uma UDF do Python. Ele executa as seguintes ações:

  1. Recupere a credencial default utilizando o provedor de credenciais do serviço Databricks.
  2. Configura uma sessão boto3.
  3. Invoca uma função Lambda para processar uma string de entrada.
Python
from pyspark.sql.functions import udf
from pyspark.sql.types import StringType

@udf(StringType())
def call_lambda_udf(input_str):
import boto3
import json
import base64
from databricks.service_credentials import getServiceCredentialsProvider
from pyspark.taskcontext import TaskContext

# Create a session using the default Unity Catalog service credential
session = boto3.Session()
client = session.client("lambda", region_name="us-west-2")

# Optionally attach Spark TaskContext metadata to the Lambda request
user_ctx = {"custom": {"user": TaskContext.get().getLocalProperty("user")}}

# Build the Lambda payload
payload = json.dumps({
"values": [input_str],
"is_debug": False
})

# Encode context for Lambda's client context
encoded_ctx = base64.b64encode(json.dumps(user_ctx).encode("utf-8")).decode("utf-8")

# Call the Lambda function
response = client.invoke(
FunctionName="HashValuesFunction",
InvocationType="RequestResponse",
ClientContext=encoded_ctx,
Payload=payload,
)

response_payload = json.loads(response["Payload"].read().decode("utf-8"))

if "errorMessage" in response_payload:
raise Exception(response_payload["errorMessage"])

return response_payload["values"][0]

Usar uma credencial de serviço em uma UDF Python scalar do Unity Catalog​

Especifique a credencial de serviço na cláusula CREDENTIALS da definição da UDF. Você pode marcar uma credencial como DEFAULT para que os SDKs de cloud com patch a usem automaticamente. No compute clássico, este recurso requer o Databricks Runtime 18.1 ou superior. Em compute serverless e em SQL warehouses Pro e serverless, defina explicitamente o environment_version da UDF como 6 ou acima. Para obter requisitos completos de compute, rede e permissões, consulte Usar uma credencial de serviço em uma UDF Python.

Exemplo de credencial de serviço: função AWS Lambda​

The following example uses a credencial de serviço to call an AWS Lambda function from a scalar Unity Catalog Python UDF. The example uses serverless environment version 6.

SQL
CREATE OR REPLACE FUNCTION main.test.call_lambda_func(data STRING, debug BOOLEAN)
RETURNS STRING
LANGUAGE PYTHON
CREDENTIALS (
`scalar-uc-udf-service-creds-example-cred` DEFAULT
)
ENVIRONMENT (
dependencies = '["boto3"]',
environment_version = '6'
)
AS $$
import base64
import boto3
import json
from pyspark.taskcontext import TaskContext

session = boto3.Session()
client = session.client("lambda", region_name="us-west-2")

user_context = {"custom": {"user": TaskContext.get().getLocalProperty("user")}}
payload = json.dumps({"values": [data], "is_debug": debug})

response = client.invoke(
FunctionName="HashValuesFunction",
InvocationType="RequestResponse",
ClientContext=base64.b64encode(json.dumps(user_context).encode("utf-8")).decode(
"utf-8"
),
Payload=payload,
)

response_payload = json.loads(response["Payload"].read().decode("utf-8"))
if "errorMessage" in response_payload:
raise Exception(str(response_payload))

return response_payload["values"][0]
$$;

Chame a UDF após ela ser registrada:

SQL
SELECT main.test.call_lambda_func(data, false)
FROM VALUES
('abc'),
('def')
AS t(data)

Obter o contexto de execução da tarefa​

Use o TaskContext PySpark API para obter informações de contexto, como a identidade do usuário, a tag do cluster, o ID do spark job e muito mais. Consulte Obter contexto de tarefa em um UDF.

Limitações​

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

  • Restrições de acesso a arquivos: Em Databricks Runtime 14.2 e abaixo, os UDFs PySpark em clustering compartilhado não podem acessar pastas Git, arquivos workspace ou volumes Unity Catalog.

  • Variáveis de difusão: PySpark UDFs em clustering de modo de acesso padrão e serverless compute não oferecem suporte a variáveis de broadcast.

  • perfil de instância: PySpark UDFs em clustering de modo de acesso padrão e serverless compute não são compatíveis com o perfil de instância.

  • Limite de memória em serverless : PySpark Os UDFs em serverless compute têm um limite de memória de 1 GB 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 : as 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 do tipo UDF_PYSPARK_USER_CODE_ERROR.MEMORY_LIMIT.

  • Acesso à rede em um data warehouse SQL serverless : Por default, as UDFs Python em um data warehouse SQL serverless 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 rede de saída, habilite o recurso de Pré-visualização Pública "Habilitar rede para cargas de trabalho isoladas no SQL Warehouse sem servidor" na página de Pré-visualizações do seu workspace. Caso contrário, utilize compute serverless ou compute clássica para UDFs que exigem acesso à rede.