Fichiers de processus avec des UDF

Important

Cette fonctionnalité est en version bêta. Les administrateurs d’espace de travail peuvent contrôler l’accès à cette fonctionnalité à partir de la page Aperçus . Consultez Gérer les préversions d’Azure Databricks.

Utilisez une fonction définie par l’utilisateur (UDF) pour traiter les fichiers référencés par une FILE colonne avec votre propre code et bibliothèques. La UDF reçoit chaque FILE valeur comme référence de fichier native à la langue. Il peut lire les octets du fichier ou l’ouvrir comme un chemin local, puis retourner une valeur de métadonnées, un fichier dérivé ou une sortie transformée.

Cette page présente les UDF de traitement de fichiers en Python, Scala et SQL. Pour la FILE référence typographique, voir FILE type. Pour la création générale de UDF, voir les fonctions définies par l’utilisateur (UDF) Python scalaires, les UDF Scala et Java à portée de session, et les fonctions de table définies par l’utilisateur (UDTF) en Python.

Lire les métadonnées de fichiers dans une UDF

Une valeur possède des champs FILE de métadonnées que vous pouvez lire sans ouvrir le fichier. Le tableau suivant contient les champs disponibles :

Accesseur Description
uri L’URI du fichier.
offset Un décalage dans le fichier, en octets.
size Taille du fichier, en octets.
content_type Le type MIME du fichier, quand il est connu.
checksum Une somme de contrôle utilisée pour identifier la version du fichier, sous la forme <algorithm>:<value>.

Accédez à ces champs avec la notation en point sur la FILE valeur, comme indiqué dans le code suivant :

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;

Lire le contenu du fichier dans une UDF

Une FILE valeur a deux méthodes pour lire le fichier sous-jacent :

  • as_local_file(): Retourne un chemin local que vous pouvez transmettre à n’importe quelle bibliothèque acceptant un chemin de fichier, comme une bibliothèque d’images ou de médias.
  • open(): Retourne un flux binaire qui ne lit que les octets que vous demandez, au lieu de matérialiser l’ensemble du fichier.

Les deux nécessitent un calcul Azure Databricks (un notebook ou un worker UDF) et ne sont pas disponibles sur un client Azure Databricks Connect. Vous pouvez déclarer FILE comme paramètre UDF ou type de retour dans Python, Scala et SQL. Pour l’API complète, voir FileType.

Extraire les dimensions de l’image

Vous pouvez utiliser une UDF scalaire pour renvoyer les dimensions d’une image sous forme de width x height chaîne. L’UDF appelle as_local_file() pour obtenir un chemin local, puis transmet ce chemin à une bibliothèque d’images standard (PILen Python, ImageIO dans Scala), comme montré dans le code suivant :

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()

Détecter le type d’un fichier à partir de ses octets

La UF suivante ne lit que les huit premiers octets de chaque fichier et open() détecte le type de fichier à partir de son nombre magique, sans matérialiser l’ensemble du fichier :

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()

Générer plusieurs fichiers avec un UDF de table (UDTF)

Pour transformer un fichier d’entrée en plusieurs fichiers de sortie, comme lors de la division d’une vidéo en images, utilisez un tableau UDF (UDTF). L’UDTF prend un FILE en entrée et donne une ligne par fichier de sortie, créant chaque fichier avec FileRef.from_bytes(). Déclarez la colonne du fichier comme FILE dans le schéma de returnType l’UDTF. Pour l’attribution générale des UDTF, voir Python fonction de table définie par l’utilisateur (UDTF).

Lorsqu’un UDTF (ou tout UDF) écrit de nouveaux fichiers avec FileRef.from_bytes, votre code doit répondre aux exigences suivantes :

  • Créez le volume ciblé avant de lancer l’UDTF. Un travailleur Python ne peut pas créer un volume de premier niveau. Créez-le avec CREATE VOLUME IF NOT EXISTS. À l’intérieur d’un volume existant, os.makedirs() on peut créer des sous-répertoires, mais pas le volume lui-même.
  • Passer un chemin absolu dbfs: . Retourner a FileRef à une table Delta Lake nécessite un dbfs: URI, tel que dbfs:/Volumes/my_catalog/my_schema/frames/frame_00000.jpg. Un chemin nu s’élève DELTA_VIOLATE_CONSTRAINT_WITH_VALUES.
  • Les écritures de vérification sont idempotentes. Supprimez ou passez les fichiers qui existent déjà avant d’écrire. Comme FileRef.from_bytes écrit avec des flags exclusive-create, l’écriture par-dessus un fichier existant fait FileExistsErrorapparaître .

Exemple : Extraire les images vidéo

L’UDTF suivant lit une vidéo FILE, extrait chaque image avec la av bibliothèque (PyAV), l’écrit dans un volume, et donne une ligne par image :

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)

Créez la table cible avec une FILE EXTERNAL colonne, puis appelez l’UDTF pour LATERAL étendre chaque vidéo en une ligne par image :

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;

Gouvernez les colonnes FICHIER avec filtres de lignes

Gouvernez une FILE colonne avec des filtres de ligne basés sur l’identité de l’appelant ou les métadonnées du fichier.

Filtre de ligne

Un filtre de ligne est une FUD qui retourne un BOOLEANfichier . Les lignes pour lesquelles il est retourné false sont omises des résultats de requête.

Le filtre de lignes suivant ne conserve que les lignes contenant des fichiers qui font référence à un tableau Excel, en fonction des métadonnées du content_type fichier :

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)
}

Pour plus d’informations sur l’application et la gestion des filtres de lignes, y compris les étapes et limitations de l’Explorateur de catalogue, voir Appliquer manuellement des filtres de lignes et des masques de colonnes.

Enregistrer un UDF dans le catalogue Unity

Enregistrer un UDF de traitement de fichiers dans le catalogue Unity pour le gouverner avec des permissions de catalogue et le réutiliser dans les carnets, requêtes et utilisateurs. L’enregistrement et l’exécution d’une UDF nécessitent les privilèges suivants :

  • Pour créer un UDF : USAGE et CREATE sur le schéma, et USAGE sur le catalogue.
  • Pour exécuter une UDF : EXECUTE sur la UDF, ainsi USAGE que sur le schéma et le catalogue.

L’exemple suivant enregistre une UDF SQL qui retourne l’extension d’un fichier, puis appelle la UDF pour créer une nouvelle colonne :

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;

Pour enregistrer un UDF Python ou Scala dans le catalogue Unity, voir les fonctions définies par l’utilisateur (UDF) SQL et Python dans le catalogue Unity et les fonctions de table définies par l’utilisateur Python (UDTF) dans le catalogue Unity.

Sécurité : Les UDF fonctionnent avec les privilèges du propriétaire

Le code UDF s’exécute avec les privilèges du propriétaire de la fonction, et non de l’appelant de la fonction. Les privilèges du propriétaire s’appliquent à la lecture des octets d’un FILE. Un appelant disposant uniquement EXECUTE des permissions sur la UDF, et sans accès direct au volume sous-jacent, peut toujours déclencher des lectures des fichiers référencés.

Parce qu’un UDF de traitement de fichiers est un chemin d’accès régi au contenu des fichiers, considérez les effets secondaires suivants en matière de sécurité et de gouvernance :

  • Les utilisateurs peuvent accéder au contenu des fichiers via l’UDF. Accordez EXECUTE des autorisations uniquement aux utilisateurs que vous souhaitez donner un accès indirect au contenu des fichiers.
  • Les appelants héritent de l’accès aux fichiers du propriétaire. Vérifiez que le propriétaire de l’UDF dispose d’un accès au volume qui ne devrait pas être plus large que ce que devraient avoir les appelants.

Pour plus d’informations sur la manière dont Azure Databricks détermine l’utilisateur autorisé lorsque l’exécution passe dans un corps UDF, voir Utilisateur autorisé et utilisateur de session.

Étapes suivantes