Archivos de proceso con UDFs

Importante

Esta característica se encuentra en su versión beta. Los administradores del área de trabajo pueden controlar el acceso a esta característica desde la página Vistas previas . Consulte Administrar versiones preliminares de Azure Databricks.

Utiliza una función definida por el usuario (UDF) para procesar los archivos referenciados por una FILE columna con tu propio código y librerías. La UDF recibe cada FILE valor como referencia de archivo nativa del idioma. Puede leer los bytes del archivo o abrirlo como una ruta local, y luego devolver un valor de metadatos, un archivo derivado o una salida transformada.

Esta página muestra UDFs de procesamiento de archivos en Python, Scala y SQL. Para la FILE referencia tipográfica, véase FILE tipo. Para la creación general de UDF, véase funciones definidas por el usuario (UDFs) escalares en Python, UDFs de Scala y Java con alcance de sesión, y funciones de tabla definidas por el usuario (UDTFs) en Python.

Leer metadatos de archivos en una UDF

Un FILE valor tiene campos de metadatos que puedes leer sin abrir el archivo. La siguiente tabla contiene los campos disponibles:

Descriptor Descripción
uri El URI del archivo.
offset Un desplazamiento hacia el archivo, en bytes.
size Tamaño del archivo, en bytes.
content_type El tipo MIME del archivo, cuando se conoce.
checksum Una suma de comprobación utilizada para identificar la versión del archivo, como <algorithm>:<value>.

Accede a estos campos con notación de puntos sobre el FILE valor, como se muestra en el siguiente código:

Python

from pyspark.sql.functions import col, udf
from pyspark.sql.types import BooleanType

@udf(returnType=BooleanType())
def is_large_image(file):
  return file.content_type.startswith("image/") and file.size > 5_000_000

spark.read.table("documents").select(col("file").uri, is_large_image(col("file"))).display()

Scala

import org.apache.spark.sql.functions.{col, udf}

val isLargeImage = udf { (file: FileRef) =>
  file.contentType.startsWith("image/") && file.size > 5000000L
}

spark.read.table("documents").select(col("file.uri"), isLargeImage(col("file"))).display()

SQL

SELECT file.uri, file.content_type, file.size
  FROM documents
  WHERE file.content_type LIKE 'image/%'
    AND file.size > 5000000;

Leer contenido de archivos en una UDF

Un FILE valor tiene dos métodos para leer el archivo subyacente:

  • as_local_file(): Devuelve una ruta local que puedes pasar a cualquier biblioteca que acepte una ruta de archivo, como una biblioteca de imágenes o multimedia.
  • open(): Devuelve un flujo binario que solo lee los bytes que solicitas, en lugar de materializar todo el archivo.

Ambos requieren Azure Databricks compute (un notebook o un worker UDF) y no están disponibles en un cliente de Azure Databricks Connect. Puedes declarar FILE como un parámetro UDF o tipo de retorno en Python, Scala y SQL. Para la API completa, véase FileType.

Extraer dimensiones de imagen

Puedes usar una UDF escalar para devolver las dimensiones de una imagen como una width x height cadena. La UDF solicita as_local_file() obtener una ruta local, luego pasa esa ruta a una biblioteca de imágenes estándar (PILen Python, ImageIO en Scala), como se muestra en el siguiente código:

Python

from pyspark.sql.functions import col, udf
from pyspark.sql.types import StringType
from PIL import Image

@udf(returnType=StringType())
def image_resolution(file):
  # as_local_file() returns a pathlib.Path.
  with Image.open(file.as_local_file()) as img:
    return f"{img.width}x{img.height}"

spark.read.table("images").select(col("photo").uri, image_resolution(col("photo"))).display()

Scala

import org.apache.spark.sql.functions.{col, udf}
import javax.imageio.ImageIO

val imageResolution = udf { (file: FileRef) =>
  // asLocalFile() returns a java.io.File.
  val image = ImageIO.read(file.asLocalFile())
  s"${image.getWidth}x${image.getHeight}"
}

spark.read.table("images").select(col("photo.uri"), imageResolution(col("photo"))).display()

Detectar el tipo de un archivo a partir de sus bytes

La siguiente UDF lee solo los primeros ocho bytes de cada archivo con open() y detecta el tipo de archivo a partir de su número mágico, sin materializar el archivo completo:

Python

from pyspark.sql.functions import col, udf
from pyspark.sql.types import StringType

@udf(returnType=StringType())
def file_signature(file):
  with file.open() as f:
    header = f.read(8)
  if header.startswith(b"%PDF"):
    return "pdf"
  if header.startswith(b"\x89PNG"):
    return "png"
  if header.startswith(b"\xff\xd8\xff"):
    return "jpeg"
  return "unknown"

spark.read.table("documents").select(col("file").uri, file_signature(col("file"))).display()

Scala

import org.apache.spark.sql.functions.{col, udf}

val fileSignature = udf { (file: FileRef) =>
  // open() returns a java.io.InputStream.
  val stream = file.open()
  try {
    val header = new Array[Byte](8)
    val n = stream.read(header)
    if (n >= 4 && header(0) == '%' && header(1) == 'P' && header(2) == 'D' && header(3) == 'F') "pdf"
    else if (n >= 4 && header(0) == 0x89.toByte && header(1) == 'P' && header(2) == 'N' && header(3) == 'G') "png"
    else if (n >= 3 && header(0) == 0xFF.toByte && header(1) == 0xD8.toByte && header(2) == 0xFF.toByte) "jpeg"
    else "unknown"
  } finally {
    stream.close()
  }
}

spark.read.table("documents").select(col("file.uri"), fileSignature(col("file"))).display()

Generar múltiples archivos con una tabla UDF (UDTF)

Para convertir un archivo de entrada en muchos archivos de salida, como al dividir un vídeo en fotogramas, utiliza un UDF de tabla (UDTF). El UDTF toma a FILE como entrada y produce una fila por archivo de salida, creando cada archivo con FileRef.from_bytes(). Declara la columna del archivo como FILE en el esquema de returnType la UDTF. Para la creación general de UDTF, véase funciones de tabla definidas por el usuario (UDTFs) en Python.

Cuando una UDTF (o cualquier UDF) escribe nuevos archivos con FileRef.from_bytes, tu código debe cumplir los siguientes requisitos:

  • Crea el volumen objetivo antes de ejecutar el UDTF. Un trabajador de Python no puede crear un volumen de primer nivel. Créalo con CREATE VOLUME IF NOT EXISTS. Dentro de un volumen existente, os.makedirs() pueden crear subdirectorios, pero no el volumen en sí.
  • Pasar un camino absoluto dbfs: . Devolver a FileRef a a una tabla Delta Lake requiere un dbfs: URI, como dbfs:/Volumes/my_catalog/my_schema/frames/frame_00000.jpg. Un camino desnudo eleva DELTA_VIOLATE_CONSTRAINT_WITH_VALUES.
  • Verificar que las escrituras son idempotentes. Elimina o omite archivos que ya existen antes de escribir. Como FileRef.from_bytes escribe con flags exclusive-create, escribir sobre un archivo existente eleva FileExistsError.

Ejemplo: Extraer fotogramas de vídeo

El siguiente UDTF lee un vídeo FILE, extrae cada fotograma con la av biblioteca (PyAV), lo escribe en un volumen y produce una fila por fotograma:

import io
import os
import av
from pyspark.sql.functions import udtf
from pyspark.sql.types import FileRef

@udtf(returnType="clip_id STRING, frame_index INT, frame FILE")
class ExtractFrames:
    def __init__(self):
        self.output_dir = "/Volumes/my_catalog/my_schema/frames/"
        os.makedirs(self.output_dir, exist_ok=True)

    def eval(self, video):
        clip_id = video.uri.split("/")[-1].split(".")[0]
        container = av.open(video.as_local_file())
        stream = container.streams.video[0]
        for i, frame in enumerate(container.decode(stream)):
            buffer = io.BytesIO()
            frame.to_image().save(buffer, format="JPEG")

            local_path = os.path.join(self.output_dir, f"{clip_id}_frame_{i:05d}.jpg")
            if os.path.exists(local_path):
                os.remove(local_path)

            yield (
                clip_id,
                i,
                FileRef.from_bytes(buffer.getvalue(), path=f"dbfs:{local_path}", content_type="image/jpeg"),
            )
        container.close()

spark.udtf.register("extract_frames", ExtractFrames)

Crea la tabla objetivo con una FILE EXTERNAL columna, luego llama al UDTF con LATERAL para expandir cada vídeo en una fila por fotograma:

CREATE TABLE my_catalog.my_schema.drive_frames (
  clip_id STRING,
  frame_index INT,
  frame FILE EXTERNAL
);

INSERT INTO my_catalog.my_schema.drive_frames
  SELECT *
  FROM my_catalog.my_schema.drive_clips AS c
  JOIN LATERAL extract_frames(c.video) AS f;

Columnas de FILE de gobierno con filtros de fila

Goberna una FILE columna con filtros de fila basados en la identidad del llamante o en los metadatos del archivo.

Filtro de fila

Un filtro de fila es una UDF que devuelve un BOOLEANarchivo . Las filas para las que se devuelve false se omiten en los resultados de la consulta.

El siguiente filtro de filas mantiene solo filas con archivos que hacen referencia a una hoja de cálculo de Excel, basándose en los metadatos del content_type archivo:

SQL

CREATE FUNCTION excel_only(file FILE)
  RETURN file.content_type IN (
    'application/vnd.openxmlformats-officedocument.spreadsheetml.sheet',
    'application/vnd.ms-excel');

ALTER TABLE documents SET ROW FILTER excel_only ON (file);

Python

from pyspark.sql.functions import udf
from pyspark.sql.types import BooleanType

@udf(returnType=BooleanType())
def excel_only(file):
  return file.content_type in (
    "application/vnd.openxmlformats-officedocument.spreadsheetml.sheet",
    "application/vnd.ms-excel")

Scala

import org.apache.spark.sql.functions.udf

val excelOnly = udf { (file: FileRef) =>
  Set(
    "application/vnd.openxmlformats-officedocument.spreadsheetml.sheet",
    "application/vnd.ms-excel").contains(file.contentType)
}

Para más información sobre cómo aplicar y gestionar filtros de fila, incluyendo los pasos y limitaciones del Explorador de Catálogos, consulte Aplicar manualmente filtros de fila y máscaras de columna.

Registrar una UDF en el Catálogo de Unity

Registrar un UDF de procesamiento de archivos en el Catálogo de Unity para gobernarlo con permisos de catálogo y reutilizarlo entre cuadernos, consultas y usuarios. Registrar y ejecutar una UDF requiere los siguientes privilegios:

  • Para crear una UDF: USAGE y CREATE en el esquema, y USAGE en el catálogo.
  • Para ejecutar una UDF: EXECUTE en la UDF, y USAGE en el esquema y catálogo.

El siguiente ejemplo registra un UDF SQL que devuelve la extensión de un archivo, y luego llama al UDF para crear una nueva columna:

CREATE FUNCTION my_catalog.my_schema.file_extension(file FILE)
  RETURNS STRING
  RETURN lower(element_at(split(file.uri, '\\.'), -1));

SELECT file.uri, my_catalog.my_schema.file_extension(file) AS extension
  FROM documents;

Para registrar un UDF de Python o Scala en el Catálogo de Unity, consulte funciones definidas por el usuario (UDFs) de SQL y Python en el Catálogo de Unity y funciones de tabla definidas por el usuario (UDTFs) de Python en el Catálogo de Unity.

Seguridad: Los UDF funcionan con los privilegios del propietario

El código UDF se ejecuta con los privilegios del propietario de la función, no del llamador de la función. Los privilegios del propietario se aplican a la lectura de los bytes de un FILE. Un llamador con solo EXECUTE permisos sobre la UDF y sin acceso directo al volumen subyacente puede seguir activando lecturas de los archivos referenciados.

Dado que un UDF de procesamiento de archivos es una ruta de acceso gobernada al contenido de los archivos, considera los siguientes efectos secundarios de seguridad y gobernanza:

  • Los usuarios pueden acceder al contenido de los archivos usando la UDF. Concede EXECUTE permisos solo a los usuarios que pretendas dar acceso indirecto al contenido del archivo.
  • Los llamantes heredan el acceso al archivo del propietario. Verifica que el propietario de la UDF tenga acceso a volumen no más amplio del que deberían tener los llamantes.

Para más información sobre cómo Azure Databricks determina el usuario autorizado cuando la ejecución cruza a un cuerpo UDF, consulte Usuario autorizado y usuario de sesión.

Pasos siguientes