Dateien mit UDFs bearbeiten

Important

Dieses Feature befindet sich in der Betaversion. Arbeitsbereichsadministratoren können den Zugriff auf dieses Feature über die Vorschauseite steuern. Siehe Manage Azure Databricks Previews.

Verwenden Sie eine benutzerdefinierte Funktion (UDF), um die Dateien zu verarbeiten, die in einer Spalte FILE mit eigenem Code und Bibliotheken referenziert werden. Der UDF erhält jeden FILE-Wert als sprachspezifische Dateireferenz. In Python ist diese Referenz ein FileRef-Objekt, das Sie aus pyspark.sql.types importieren können. Der UDF kann die Bytes der Datei lesen oder als lokalen Pfad öffnen und dann einen Metadatenwert, eine abgeleitete Datei oder eine transformierte Ausgabe zurückgeben.

Diese Seite zeigt Dateiverarbeitungs-UDFs in Python, Scala und SQL. Für die FILE Typreferenz siehe FILE Typ. Für allgemeine UDF-Erstellung siehe Python skalare benutzerdefinierte Funktionen (UDFs),Session-scoped Scala und Java UDFs sowie Python user-defined table functions (UDTFs).

Dateimetadaten in einem UDF lesen

Ein FILE Wert enthält Metadatenfelder, die du lesen kannst, ohne die Datei zu öffnen. Die folgende Tabelle enthält die verfügbaren Felder:

Accessor Description
uri Die URI der Datei.
offset Ein Offset innerhalb der Datei, in Byte.
size Die Größe der Datei in Byte.
content_type Der MIME-Typ der Datei, wenn bekannt.
checksum Eine Prüfsumme, die zum Identifizieren der Dateiversion verwendet wird, z. B. <algorithm>:<value>.

Greifen Sie auf diese Felder mit Punktzeichen auf dem FILE Wert zu, wie im folgenden Code gezeigt:

Python

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

@udf(returnType=BooleanType())
def is_large_image(file: FileRef) -> bool:
  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;

Dateiinhalte in einem UDF lesen

Ein FILE Wert hat zwei Methoden, um die zugrundeliegende Datei zu lesen:

  • as_local_file(): Gibt einen lokalen Pfad zurück, den Sie an jede Bibliothek weitergeben können, die einen Dateipfad akzeptiert, wie z. B. eine Bild- oder Medienbibliothek.
  • open(): Gibt einen Binärstrom zurück, der nur die von dir angeforderten Bytes liest, anstatt die gesamte Datei zu materialisieren.

Beide erfordern Azure Databricks Compute (einen Notebook- oder UDF-Worker) und sind auf einem Azure Databricks Connect-Client nicht verfügbar. Du kannst als UDF-Parameter oder Rückgabetyp in Python-, Scala- und SQL-UDFs deklarierenFILE. Für die vollständige API siehe FileType.

Bildmaße extrahieren

Du kannst ein skalares UDF verwenden, um die Abmessungen eines Bildes als Zeichenkette width x height zurückzugeben. Die UDF ruft as_local_file() auf, um einen lokalen Pfad zu erhalten, und übergibt diesen Pfad dann an eine Standard-Bildbibliothek (PIL in Python, ImageIO in Scala), wie im folgenden Code dargestellt:

Python

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

@udf(returnType=StringType())
def image_resolution(file: FileRef) -> str:
  # 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()

Erkennen Sie den Dateityp anhand ihrer Bytes

Das folgende UDF liest nur die ersten acht Bytes jeder Datei mit open() und erkennt den Dateityp anhand seiner magischen Zahl, ohne die gesamte Datei zu materialisieren:

Python

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

@udf(returnType=StringType())
def file_signature(file: FileRef) -> str:
  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()

Mehrere Dateien mit einer Tabelle UDF (UDTF) generieren

Um eine Eingabedatei in viele Ausgabedateien umzuwandeln, zum Beispiel beim Aufteilen eines Videos in Frames, verwenden Sie eine Tabellen-UDF (UDTF). Das UDTF nimmt eine FILE als Eingabe und gibt für jede Ausgabedatei eine Zeile zurück, wobei jede Datei mit FileRef.from_bytes() erstellt wird. Deklarieren Sie die Dateispalte als FILE im returnType-Schema der UDTF. Für allgemeine UDTF-Erstellung siehe Python user-defined table functions (UDTFs).

Wenn ein UDTF (oder ein beliebiger UDF) neue Dateien mit FileRef.from_bytesschreibt, muss Ihr Code die folgenden Anforderungen erfüllen:

  • Erstellen Sie das Zielvolumen, bevor Sie das UDTF ausführen. Ein Python-Worker kann kein Top-Level-Volume erstellen. Erstellen Sie es mit CREATE VOLUME IF NOT EXISTS. Innerhalb eines bestehenden Volumes os.makedirs() können Unterverzeichnisse erstellt werden, aber nicht das Volume selbst.
  • Geben Sie einen absoluten dbfs: Pfad an. Um a FileRef an eine Delta Lake-Tabelle zurückzugeben, ist eine dbfs: URI erforderlich, wie zum Beispiel dbfs:/Volumes/my_catalog/my_schema/frames/frame_00000.jpg. Ein leerer Pfad löst DELTA_VIOLATE_CONSTRAINT_WITH_VALUES aus.
  • Überprüfen Sie, ob Schreibvorgänge idempotent sind. Lösche oder überspringe Dateien, die bereits existieren, bevor du schreibst. Da FileRef.from_bytes mit Flags für exklusives Erstellen schreibt, löst das Schreiben über eine vorhandene Datei FileExistsError aus.

Beispiel: Videoframes extrahieren

Das folgende UDTF liest ein Video FILE ein, extrahiert jedes Einzelbild mithilfe der av-Bibliothek (PyAV), schreibt es in ein Volume und gibt pro Einzelbild eine Zeile zurück:

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: FileRef):
        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)

Erstellen Sie die Zieltabelle mit einer Spalte FILE EXTERNAL, und rufen Sie dann die UDTF mit LATERAL auf, um jedes Video zu einer Zeile pro Frame zu erweitern:

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;

Regele FILE-Spalten mit Zeilenfiltern

Steuern Sie eine FILESpalte mit Zeilenfiltern anhand der Identität des Aufrufers oder der Metadaten der Datei.

Zeilenfilter

Ein Zeilenfilter ist eine UDF, die ein BOOLEAN zurückgibt. Zeilen, für die sie zurückgibt false , werden in den Abfrageergebnissen weggelassen.

Der folgende Zeilenfilter behält nur Zeilen mit Dateien, die auf eine Excel-Tabelle verweisen, basierend auf den content_type Metadaten der Datei:

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, FileRef

@udf(returnType=BooleanType())
def excel_only(file: FileRef) -> bool:
  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)
}

Weitere Informationen zum Anwenden und Management von Zeilenfiltern, einschließlich der Schritte und Einschränkungen im Catalog Explorer, finden Sie unter Manuelle Anwendung von Zeilenfiltern und Spaltenmasken.

Registrieren Sie ein UDF im Unity-Katalog

Registrieren Sie ein Dateiverarbeitungs-UDF im Unity-Katalog, um es mit Katalogberechtigungen zu verwalten und es in Notizbüchern, Abfragen und Benutzern wiederzuverwenden. Die Registrierung und Ausführung eines UDF erfordert folgende Rechte:

  • Um ein UDF zu erstellen: USAGE und CREATE auf dem Schema und USAGE im Katalog.
  • Ausführen eines UDF: EXECUTE auf dem UDF und USAGE auf dem Schema und Katalog.

Das folgende Beispiel registriert ein SQL-UDF, das die Dateierweiterung zurückgibt und dann das UDF aufruft, um eine neue Spalte zu erstellen:

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;

Um eine Python- oder Scala-UDF im Unity Catalog zu registrieren, siehe SQL und Python user-defined functions (UDFs) im Unity Catalog sowie Python user-defined table functions (UDTFs) im Unity Catalog.

Sicherheit: UDFs nutzen die Rechte des Besitzers

UDF-Code läuft mit den Rechten des Funktionsbesitzers, nicht des Funktionsaufrufers. Die Rechte des Besitzers gelten für das Lesen der Bytes eines FILE. Ein Aufrufer mit nur EXECUTE Berechtigungen auf der UDF und ohne direkten Zugriff auf das zugrundeliegende Volume kann dennoch Lesungen der referenzierten Dateien auslösen.

Da ein Dateiverarbeitungs-UDF ein gesteuerter Zugriffspfad zu Dateiinhalten ist, sollten Sie folgende Sicherheits- und Governance-Nebenwirkungen berücksichtigen:

  • Benutzer können mit dem UDF auf Dateiinhalte zugreifen. Gewähren Sie EXECUTE Berechtigungen nur Nutzern, denen Sie indirekten Zugriff auf Dateiinhalte geben möchten.
  • Anrufer erben den Dateizugriff des Eigentümers. Stellen Sie sicher, dass der UDF-Besitzer keinen umfangreicheren Zugriff auf das Volume hat, als Aufrufer eigentlich haben sollten.

Für weitere Informationen darüber, wie Azure Databricks den autorisierten Benutzer bestimmt, wenn die Ausführung in einen UDF-Körper übergeht, siehe Autorisierter Benutzer und Sitzungsbenutzer.

Nächste Schritte