Traiter les fichiers avec des UDF
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 |
|---|---|
| L’URI du fichier. |
| Un décalage dans le fichier, en octets. |
| La taille du fichier, en octets. |
| Le type MIME du fichier, lorsqu'il est connu. |
| Une somme de contrôle utilisée pour identifier la version du fichier, comme |
Accédez à ces champs avec la notation par points sur la valeur FILE, comme illustré dans le code suivant :
- Python
- Scala
- SQL
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()
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()
SELECT file.uri, file.content_type, file.size
FROM documents
WHERE file.content_type LIKE 'image/%'
AND file.size > 5000000;
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
- Scala
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()
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
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
- Scala
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()
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 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'unFileRefvers une table Delta Lake nécessite un URIdbfs:, tel quedbfs:/Volumes/my_catalog/my_schema/frames/frame_00000.jpg. Un chemin d'accès nu génèreDELTA_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éclencheFileExistsError.
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 :
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 :
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
- Python
- Scala
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);
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")
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 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 :
USAGEetCREATEsur le schéma, etUSAGEsur le catalogue. - Pour exécuter une UDF :
EXECUTEsur l'UDF, etUSAGEsur 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 :
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
EXECUTEqu’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
FILETypeai_parse_documentfonction- Type de fichier
- Fonctions définies par l’utilisateur (UDF) SQL et Python dans Unity Catalog
- Utilisateur autorisé et utilisateur de session
- Appliquer manuellement des filtres de lignes et des masques de colonne
- Fonctions scalaires définies par l’utilisateur (UDF) Python
- Fonctions de table définies par l’utilisateur (UDTF) Python
- Type FILE et données non structurées