Nota:
El acceso a esta página requiere autorización. Puede intentar iniciar sesión o cambiar directorios.
El acceso a esta página requiere autorización. Puede intentar cambiar los directorios.
Los UDF escalares de Python te permiten ejecutar lógica Python personalizada dentro de consultas SQL en Azure Databricks. Esta página muestra cómo registrarlos e invocarlos, usar credenciales y secretos de servicio, y gestionar las advertencias sobre el orden de evaluación de subexpresiones en Spark SQL.
Requisitos
En Databricks Runtime 12.2 LTS y versiones posteriores, las UDF de Python y las UDF de Pandas no se admiten en el proceso de Catálogo de Unity que usa el modo de acceso estándar.
Las UDF escalares de Python y las UDF de Pandas se admiten en Databricks Runtime 13.3 LTS y versiones posteriores para todos los modos de acceso.
La compatibilidad de instancias ARM para UDFs de Python en clústeres habilitados con Unity Catalog requiere Databricks Runtime 15.2 o superior.
En Databricks Runtime 14.0 y versiones anteriores, las UDF de Python y UDF de Pandas no se admiten en clústeres de Catálogo de Unity que usan el modo de acceso estándar. Las UDF escalares de Python y las UDF de Pandas se admiten para todos los modos de acceso en Databricks Runtime 14.1 y versiones posteriores.
En Databricks Runtime 14.1 y versiones posteriores, puede registrar UDF escalares de Python en el catálogo de Unity mediante la sintaxis SQL. Consulte las funciones definidas por el usuario (UDFs) de SQL y Python en Unity Catalog.
Registro de una función como UDF
def squared(s):
return s * s
spark.udf.register("squaredWithPython", squared)
Puede configurar opcionalmente el tipo de retorno de su UDF. El tipo de valor devuelto predeterminado es StringType.
from pyspark.sql.types import LongType
def squared_typed(s):
return s * s
spark.udf.register("squaredWithPython", squared_typed, LongType())
Llamada a la UDF en Spark SQL
spark.range(1, 20).createOrReplaceTempView("test")
%sql select id, squaredWithPython(id) as id_squared from test
Uso de UDF con 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, puede declarar la misma UDF mediante la sintaxis de anotación:
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 con UDF
El tipo PySpark para variant es VariantType y los valores son de tipo VariantVal. Para obtener información sobre las variantes, consulte Consultar datos 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}}|
+-----------------------------+
Archivos con UDF
Importante
Esta característica se encuentra en su versión beta.
El tipo PySpark para un archivo es FileType. Úselo como parámetro o tipo de retorno en una UDF, ya sea como tipo de nivel superior o anidado. Para el tipo, sus reglas de anidamiento y la FileRef API, véase FileType.
Para leer el contenido de un archivo en una UDF, llama file.as_local_file() para obtener una ruta local que puedas abrir o file.open() para leer sus bytes como un flujo. Para ejemplos en Python, Scala y SQL, incluyendo procesamiento de imágenes, detección de tipos de archivo y extracción de fotogramas de vídeo, véase Archivos de proceso con UDFs. Para la referencia tipográfica, véase FILE tipo.
Una UDF registrada en el Catálogo de Unity (CREATE FUNCTION) puede leer los metadatos de un FILE', pero no su contenido, y no puede crear archivos. Utiliza un UDF con alcance de sesión para leer el contenido de un archivo (open, as_local_file) o crear uno (from_bytes, from_local_file).
En el cuerpo UDF, cada FILE valor es un 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"))))
Orden de evaluación y comprobación de valores nulos
Spark SQL (incluidas SQL y las DataFrame API y Dataset API) no garantiza el orden de evaluación de las subexpresiones. En concreto, las entradas de un operador o función no se evalúan necesariamente de izquierda a derecha ni en ningún otro orden fijo. Por ejemplo, las expresiones lógicas AND y OR no tienen semántica de "cortocircuito" de izquierda a derecha.
Por lo tanto, es peligroso basarse en los efectos secundarios o el orden de evaluación de las expresiones booleanas y el orden de las cláusulas WHERE y HAVING, ya que estas expresiones y cláusulas se pueden reordenar durante la optimización y el planeamiento de consultas. En concreto, si una UDF se basa en la semántica de cortocircuito en SQL para comprobar valores NULL, no hay ninguna garantía de que se produzca la comprobación nula antes de invocar la UDF. Por ejemplo,
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 cláusula WHERE no garantiza que se invoque la UDF strlen después de filtrar los valores NULL.
Para realizar una comprobación correcta de los valores NULL, se recomienda realizar una de las siguientes acciones:
- Hacer que la UDF tenga en cuenta los valores NULL y realizar la comprobación de valores NULL dentro de la propia UDF
- Uso de expresiones
IFoCASE WHENpara realizar la comprobación de valores NULL e invocar la UDF en una rama 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
Acceder a los secretos de Unity Catalog
Para acceder a un secreto de Unity Catalog desde una UDF de Python con ámbito de sesión, consulte Usar un secreto en una UDF de Python con ámbito de sesión. Para acceder a secretos declarados desde una UDF escalar de Python de Unity Catalog o una UDF por lotes de Python de Unity Catalog, consulte Usar secretos en una UDF de Python.
Credenciales de servicio en UDFs de Python
Los UDFs escalares de Python con alcance de sesión y los UDFs escalares de Unity Catalog Python pueden utilizar las credenciales de servicio del Catálogo de Unity para acceder de forma segura a servicios en la nube externa. Esto es útil para integrar operaciones como la tokenización basada en la nube, el cifrado o la administración de secretos directamente en las transformaciones de datos.
Los requisitos varían según el tipo de UDF y el cálculo. Ver Usar una credencial de servicio en un UDF de Python.
Para crear una credencial de servicio, consulte Creación de credenciales de servicio.
Utiliza una credencial de servicio en un UDF escalar de Python con alcance de sesión
Para acceder a la credencial de servicio, use la utilidad en la databricks.service_credentials.getServiceCredentialsProvider() lógica de UDF para inicializar los SDKs de la nube con la credencial apropiada. Todo el código debe estar encapsulado en el cuerpo de la 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
Permisos de credenciales de servicio
Las UDF de ámbito de sesión utilizan los permisos del usuario que realiza la llamada. Consulta Usar una credencial de servicio en un UDF de Python para obtener los privilegios requeridos.
Credenciales predeterminadas a nivel de cómputo para UDFs con alcance de sesión
Cuando se utiliza en UDFs escalares de Python, Databricks utiliza automáticamente la credencial de servicio predeterminada de la variable de entorno de cálculo. Este comportamiento permite hacer referencia de forma segura a servicios externos sin administrar explícitamente alias de credenciales en el código UDF. Consulte Especificación de una credencial de servicio predeterminada para un recurso de proceso.
La compatibilidad con credenciales predeterminadas solo está disponible en clústeres de modo de acceso estándar y dedicado. No está disponible en almacenes SQL.
Debe instalar el azure-identity paquete para usar el DefaultAzureCredential proveedor. Para instalar el paquete, consulte Bibliotecas de Python con ámbito de cuaderno o Bibliotecas con ámbito de proceso.
@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
Utiliza una credencial de servicio en un UDF escalar de Unity Catalog Python
Especifica la credencial de servicio en la CREDENTIALS cláusula de la definición de la UDF. Puedes marcar una credencial como DEFAULT para que los SDKs de la nube parcheados la usen automáticamente. En computación clásica, esta característica requiere Databricks Runtime 18.1 o superior. En el procesamiento sin servidor y en los almacenes SQL Pro y serverless, establece explícitamente el environment_version de la UDF en 6 o superior. Para los requisitos completos de computación, redes y permisos, véase Usar una credencial de servicio en un UDF de Python.
Obtención del contexto de ejecución de tareas
Use taskContext PySpark API para obtener información de contexto, como la identidad del usuario, las etiquetas de clúster, el identificador de trabajo de Spark, etc. Consulte Obtener contexto de tarea en una UDF.
Limitaciones
Las limitaciones siguientes se aplican a las UDF de PySpark:
Restricciones de acceso a archivos: En Databricks Runtime 14.2 y versiones posteriores, las UDF de PySpark en clústeres compartidos no pueden acceder a carpetas de Git, archivos de área de trabajo ni volúmenes de catálogo de Unity.
Variables de difusión: UDF de PySpark en clústeres de modo de acceso estándar y el proceso sin servidor no admiten variables de difusión.
- Límite de memoria en sin servidor: las UDF de PySpark en el proceso sin servidor tienen un límite de memoria de 1 GB por UDF de PySpark. Si se supera este límite, se produce un error de tipo UDF_PYSPARK_USER_CODE_ERROR. MEMORY_LIMIT_SERVERLESS.
- Límite de memoria en modo de acceso estándar: las UDF de PySpark en el modo de acceso estándar tienen un límite de memoria en función de la memoria disponible del tipo de instancia elegido. Si se supera la memoria disponible, se produce un error de tipo UDF_PYSPARK_USER_CODE_ERROR. MEMORY_LIMIT.
- Acceso de red en almacenes SQL sin servidor: De forma predeterminada, las UDF de Python en almacenes SQL sin servidor no pueden realizar solicitudes de red salientes, y las consultas que intentan realizar llamadas de red se bloquean indefinidamente. Para habilitar el acceso de red saliente, active la característica de Vista previa pública Habilitar la red para cargas de trabajo aisladas en almacenes de SQL sin servidor en la página Versiones preliminares de su área de trabajo. De lo contrario, use el proceso sin servidor o el proceso clásico para las UDF que requieren acceso a la red.