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 sprachnatives Dateireferenz. Es 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 in die Datei, in Bytes.
size Die Größe der Datei in Byte.
content_type Der MIME-Typ der Datei, wenn bekannt.
checksum Eine Kontrollsumme, die zur Identifikation der Dateiversion verwendet wird, als <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

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

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. Der UDF-Aufruf, as_local_file() um einen lokalen Pfad zu erhalten, übergibt diesen dann an eine Standard-Image-Bibliothek (PILin Python, ImageIO in Scala), wie im folgenden Code gezeigt:

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

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

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 liefert eine Zeile pro Ausgabedatei, wobei jede Datei mit erstellt wird FileRef.from_bytes(). Deklarieren Sie die Dateispalte wie FILE im UDTF-Schema.returnType 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.
  • Geh einen absoluten dbfs: Weg durch. 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 kahler Weg erhebt DELTA_VIOLATE_CONSTRAINT_WITH_VALUESsich.
  • Überprüfen Sie, ob Schreibvorgänge idempotent sind. Lösche oder überspringe Dateien, die bereits existieren, bevor du schreibst. Da FileRef.from_bytes mit exklusiven Erstellungs-Flags geschrieben wird, führt FileExistsErrordas Überschreiben einer bestehenden Datei zu .

Beispiel: Videoframes extrahieren

Das folgende UDTF liest ein Video FILEaus, extrahiert jeden Frame mit der av (PyAV)-Bibliothek, schreibt ihn in ein Volume und liefert pro Frame eine Zeile:

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)

Erstellen Sie die Zieltabelle mit einer Spalte FILE EXTERNAL und rufen Sie dann das UDTF mit auf LATERAL , um jedes Video auf eine Zeile pro Bild 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

Lenken Sie eine FILE Spalte mit Zeilenfiltern basierend auf der Identität des Aufrufers oder den Metadaten der Datei.

Zeilenfilter

Ein Zeilenfilter ist ein UDF, der ein Zeilenfilter BOOLEANzurü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

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

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.
  • Um einen UDF auszuführen: EXECUTE auf dem UDF sowie 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. Überprüfen Sie, dass der Eigentümer des UDF keinen größeren Zugriff auf das Volumen hat, als Anrufer 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