Bestanden verwerken met UDF's

Important

Deze functie bevindt zich in de bètaversie. Werkruimtebeheerders kunnen de toegang tot deze functie beheren vanaf de pagina Previews . Zie Azure Databricks previews beheren.

Gebruik een door de gebruiker gedefinieerde functie (UDF) om de bestanden te verwerken waarnaar FILE een kolom verwijst, met je eigen code en bibliotheken. De UDF ontvangt elke FILE waarde als een taal-native bestandsreferentie. Het kan de bytes van het bestand lezen of het als lokaal pad openen, waarna een metadatawaarde, een afgeleid bestand of getransformeerde output wordt teruggegeven.

Deze pagina toont bestandsverwerkende UDF's in Python, Scala en SQL. Voor de FILE typereferentie, zie FILE type. Voor algemene UDF-authoring, zie Python scalar user-defined functions (UDF's),Session-scoped Scala en Java UDF's, en Python user-defined table functions (UDTFs).

Lees bestandsmetadata in een UDF

Een FILE waarde bevat metadatavelden die je kunt lezen zonder het bestand te openen. De volgende tabel bevat de beschikbare velden:

Accessor Description
uri De URI van het bestand.
offset Een offset in het bestand, in bytes.
size De grootte van het bestand, in bytes.
content_type Het MIME-type van het bestand, wanneer bekend.
checksum Een controlesum die wordt gebruikt om de bestandsversie te identificeren, als <algorithm>:<value>.

Toegang tot deze velden met puntnotatie op de FILE waarde, zoals weergegeven in de volgende code:

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;

Lees de inhoud van bestanden in een UDF

Een FILE waarde heeft twee methoden om het onderliggende bestand te lezen:

  • as_local_file(): Geeft een lokaal pad terug dat je kunt doorgeven aan elke bibliotheek die een bestandspad accepteert, zoals een afbeeldings- of mediabibliotheek.
  • open(): Geeft een binaire stroom terug die alleen de bytes leest die je vraagt, in plaats van het hele bestand te materialiseren.

Beide vereisen Azure Databricks compute (een notebook- of UDF-worker) en zijn niet beschikbaar op een Azure Databricks Connect-client. Je kunt declareren FILE als UDF-parameter of returntype in Python-, Scala- en SQL-UDF's. Voor de volledige API, zie FileType.

Afbeeldingsafmetingen extraheren

Je kunt een scalair UDF gebruiken om de afmetingen van een afbeelding als string width x height terug te geven. De UDF-aanroepen as_local_file() om een lokaal pad te krijgen, geeft dat pad vervolgens door aan een standaard imagebibliotheek (PILin Python, ImageIO in Scala), zoals weergegeven in de volgende code:

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

Detecteer het type van een bestand aan de hand van de bytes ervan

De volgende UDF leest alleen de eerste acht bytes van elk bestand en open() detecteert het bestandstype aan de hand van het magische getal, zonder het hele bestand te materialiseren:

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

Genereer meerdere bestanden met een tabel UDF (UDTF)

Om één invoerbestand om te zetten in meerdere uitvoerbestanden, bijvoorbeeld bij het splitsen van een video in frames, gebruik je een table UDF (UDTF). De UDTF neemt een FILE als invoer en levert één rij per uitvoerbestand op, waarbij elk bestand wordt aangemaakt met FileRef.from_bytes(). Verklaar de bestandskolom zoals FILE in het returnType schema van de UDTF. Voor algemene UDTF-authoring, zie Python user-defined table functions (UDTFs).

Wanneer een UDTF (of een andere UDF) nieuwe bestanden schrijft met FileRef.from_bytes, moet je code aan de volgende eisen voldoen:

  • Maak het doelvolume aan voordat je de UDTF uitvoert. Een Python-worker kan geen top-level volume maken. Maak het met CREATE VOLUME IF NOT EXISTS. Binnen een bestaand volume kun je os.makedirs() subdirectories aanmaken, maar niet het volume zelf.
  • Ga een absoluut dbfs: pad door. Het terugsturen van a FileRef naar een Delta Lake-tabel vereist een dbfs: URI, zoals dbfs:/Volumes/my_catalog/my_schema/frames/frame_00000.jpg. Een kale weg rijst DELTA_VIOLATE_CONSTRAINT_WITH_VALUESomhoog.
  • Controleer of schrijfopdrachten idempotent zijn. Verwijder of sla bestanden over die al bestaan voordat je schrijft. Omdat FileRef.from_bytes schrijft met exclusive-create-flags, genereert FileExistsErrorhet overschrijven van een bestaand bestand .

Voorbeeld: Extraheren videoframes

De volgende UDTF leest een video FILE, haalt elk frame uit met de av (PyAV)-bibliotheek, schrijft het naar een volume en levert één rij per frame op:

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)

Maak de doeltabel met een FILE EXTERNAL kolom, roep vervolgens de UDTF aan LATERAL om elke video uit te breiden tot één rij per frame:

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;

Bestuur FILE-kolommen met rijfilters

Bestuur een FILE kolom met rijfilters op basis van de identiteit van de aanroeper of de metadata van het bestand.

Rijfilter

Een rijfilter is een UDF die een BOOLEANretour geeft. Rijen waarvoor het teruggeeft false worden weggelaten uit de zoekresultaten.

Het volgende rijfilter bewaart alleen rijen met bestanden die een Excel-spreadsheet verwijzen, gebaseerd op de content_type metadata van het bestand:

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

Voor meer informatie over het toepassen en beheren van rijfilters, inclusief de stappen en beperkingen van Catalog Explorer, zie Handmatig rijfilters en kolommaskers toepassen.

Registreer een UDF in de Unity Catalog

Registreer een bestandsverwerkende UDF in Unity Catalog om deze te beheren met catalogusrechten en hergebruik deze over notitieboeken, queries en gebruikers. Het registreren en uitvoeren van een UDF vereist de volgende privileges:

  • Om een UDF te maken: USAGE en CREATE op het schema, en USAGE op de catalogus.
  • Om een UDF te draaien: EXECUTE op de UDF, en USAGE op het schema en de catalogus.

Het volgende voorbeeld registreert een SQL UDF die de extensie van een bestand teruggeeft, waarna de UDF wordt aangeroepen om een nieuwe kolom aan te maken:

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;

Om een Python- of Scala-UDF te registreren in de Unity Catalog, zie SQL- en Python-gebruikersgedefinieerde functies (UDF's) in de Unity Catalog en Python user-defined table functions (UDTFs) in de Unity Catalog.

Beveiliging: UDF's werken met de rechten van de eigenaar

UDF-code draait met de rechten van de eigenaar van de functie, niet van de functiecaller. De rechten van de eigenaar gelden voor het lezen van de bytes van een FILE. Een caller met alleen EXECUTE rechten op de UDF en geen directe toegang tot het onderliggende volume, kan nog steeds het lezen van de genoemde bestanden triggeren.

Omdat een bestandsverwerkings-UDF een gereguleerd toegangspad is naar de inhoud van bestanden, overweeg de volgende beveiligings- en governance-neveneffecten:

  • Gebruikers kunnen de bestandsinhoud raadplegen via de UDF. Verleen EXECUTE alleen rechten aan gebruikers die je indirecte toegang tot de inhoud van bestanden wilt geven.
  • Bellers erven de bestandstoegang van de eigenaar. Controleer of de eigenaar van de UDF toegang heeft tot het volume, niet breder dan wat bellers zouden moeten hebben.

Voor meer informatie over hoe Azure Databricks de geautoriseerde gebruiker bepaalt wanneer de uitvoering overgaat in een UDF-lichaam, zie Geautoriseerde gebruiker en sessiegebruiker.

Volgende stappen