Aller au contenu principal

Traiter les fichiers avec des UDF

info

Bêta

Cette fonctionnalité est en bêta. Les administrateurs de Workspace peuvent contrôler l'accès à cette fonctionnalité à partir de la page Aperçus . Voir Gérer les prévisualisations Databricks.

Utilisez une fonction définie par l’utilisateur (UDF) pour traiter les fichiers référencés par une colonne FILE avec votre propre code et vos propres bibliothèques. L’UDF reçoit chaque valeur FILE en tant que référence de fichier native au langage. Il peut lire les octets du fichier ou l’ouvrir en tant que chemin local, puis renvoyer une valeur de métadonnée, 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 référence de type FILE, voir typeFILE. Pour la création générale d'UDF, consultez Fonctions scalaires définies par l'utilisateur (UDF) Python, UDF Scala et Java à portée de session et Fonctions de table définies par l'utilisateur (UDTF) Python.

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

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

Accesseur

Description

uri

L’URI du fichier.

offset

Un décalage dans le fichier, en octets.

size

La taille du fichier, en octets.

content_type

Le type MIME du fichier, lorsqu'il est connu.

checksum

Une somme de contrôle utilisée pour identifier la version du fichier, comme <algorithm>:<value>.

Accesseur

Description

uri

L’URI du fichier.

offset

Un décalage dans le fichier, en octets.

size

La taille du fichier, en octets.

content_type

Le type MIME du fichier, lorsqu'il est connu.

checksum

Une somme de contrôle utilisée pour identifier la version du fichier, comme <algorithm>:<value>.

Accédez à ces champs avec la notation par points sur la valeur FILE, comme illustré 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()

Lire le contenu d'un fichier dans une UDF

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

  • as_local_file(): renvoie un chemin local que vous pouvez transmettre à toute bibliothèque acceptant un chemin de fichier, telle qu’une bibliothèque d’images ou de médias.
  • open(): renvoie un Stream binaire qui ne lit que les octets demandés, au lieu de matérialiser le fichier entier.

Les deux nécessitent un compute Databricks (un notebook ou un worker UDF) et ne sont pas disponibles sur un client Databricks Connect. Vous pouvez déclarer FILE en tant que parameter ou type de retour UDF dans les UDF Python, Scala et SQL. Pour l’API complète, consultez FileType.

Extraire les dimensions de l'image

Vous pouvez utiliser une UDF scalaire pour renvoyer les dimensions d'une image sous forme de chaîne width x height. L’UDF appelle as_local_file() pour obtenir un chemin local, puis transmet ce chemin à une bibliothèque d’images standard (PIL en Python, ImageIO en Scala), comme illustré 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()

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

L'UDF suivante lit uniquement les huit premiers octets de chaque fichier avec open() et détecte le type de fichier à partir de son nombre magique, sans matérialiser le fichier entier :

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

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

Pour transformer un fichier d'entrée en plusieurs fichiers de sortie, par exemple lors du découpage d'une vidéo en images, utilisez une UDF de table (UDTF). L'UDTF prend un FILE en entrée et produit une ligne par fichier de sortie, en créant chaque fichier avec FileRef.from_bytes(). Déclarez la colonne de fichier comme FILE dans le schéma returnType de l'UDTF. Pour la création générale d'UDTF, consultez les fonctions de table définies par l'utilisateur (UDTF) en Python.

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

  • Créez le volume cible avant d’exécuter l’UDTF. Un Worker Python ne peut pas créer de volume de niveau supérieur. Créez-le avec CREATE VOLUME IF NOT EXISTS. Dans un volume existant, os.makedirs() peut créer des sous-répertoires, mais pas le volume lui-même.
  • Transmettez un chemin dbfs: absolu. Le renvoi d'un FileRef vers une table Delta Lake nécessite un URI dbfs:, tel que dbfs:/Volumes/my_catalog/my_schema/frames/frame_00000.jpg. Un chemin d'accès nu génère DELTA_VIOLATE_CONSTRAINT_WITH_VALUES.
  • Vérifiez que les écritures sont idempotentes. Supprimez ou ignorez les fichiers qui existent déjà avant l’écriture. Comme FileRef.from_bytes écrit avec des indicateurs de création exclusive, l’écriture sur un fichier existant déclenche FileExistsError.

Exemple : extraire des images vidéo

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

Python
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 colonne FILE EXTERNAL, puis appelez la UDTF avec LATERAL pour développer chaque vidéo en une ligne par image :

SQL
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;

Gouverner les colonnes FILE avec des filtres de ligne

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

Filtre de ligne

Un filtre de ligne est une UDF qui renvoie un BOOLEAN. Les lignes pour lesquelles elle renvoie false sont omises des résultats de la query.

Le filtre de ligne suivant ne conserve que les lignes contenant des fichiers qui font référence à une feuille de calcul Excel, en fonction des métadonnées content_type du 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);

Pour plus d’informations sur l’application et la gestion des filtres de lignes, y compris les étapes et les limitations de Catalog Explorer, consultez Appliquer manuellement des filtres de lignes et des masques de colonne.

Enregistrer une UDF dans Unity Catalog

Enregistrez une UDF de traitement de fichiers dans Unity Catalog pour la régir avec les autorisations du catalogue et la réutiliser dans des Notebooks, des query et entre les utilisateurs. [[ ## completed ##]] L'enregistrement et l'exécution d'une UDF nécessitent les privilèges suivants :

  • Pour créer une UDF : USAGE et CREATE sur le schéma, et USAGE sur le catalogue.
  • Pour exécuter une UDF : EXECUTE sur l'UDF, et USAGE sur le schéma et le catalogue.

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

SQL
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 une UDF Python ou Scala dans Unity Catalog, consultez SQL and Python user-defined functions (UDFs) in Unity Catalog et Python user-defined table functions (UDTFs) in Unity Catalog.

Sécurité : les UDF s’exécutent 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 ceux de l’appelant. Les privilèges du propriétaire s’appliquent à la lecture des octets d’un FILE. Un appelant disposant uniquement des autorisations EXECUTE sur l’UDF, et sans accès direct au volume sous-jacent, peut toujours Trigger des lectures des fichiers référencés.

Comme une UDF de traitement de fichiers est un chemin d’accès gouverné au contenu des fichiers, tenez compte des effets secondaires suivants en matière de sécurité et de gouvernance :

  • Les utilisateurs peuvent accéder au contenu des fichiers à l’aide de l’UDF. N’accordez les autorisations EXECUTE qu’aux utilisateurs auxquels 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 n'est pas plus large que celui dont les appelants devraient disposer.

Pour plus d’informations sur la façon dont Databricks détermine l’utilisateur autorisé lorsque l’exécution passe dans un corps d’UDF, consultez Utilisateur autorisé et utilisateur de session.

Étapes suivantes