Processa filer med UDF:er

Important

Den här funktionen finns i Beta. Arbetsyteadministratörer kan styra åtkomsten till den här funktionen från sidan Förhandsversioner . Se Hantera förhandsversioner av Azure Databricks.

Använd en användardefinierad funktion (UDF) för att bearbeta filerna som refereras till i en FILE kolumn med din egen kod och bibliotek. UDF tar emot varje FILE värde som en språkbaserad filreferens. Den kan läsa filens byte eller öppna den som en lokal sökväg, och sedan returnera ett metadatavärde, en härledd fil eller transformerad utdata.

Den här sidan visar filbehandlings-UDF:er i Python, Scala och SQL. För FILE typreferensen, se FILE typ. För allmän UDF-författande, se Python skalar användardefinierade funktioner (UDF),session-scoped Scala och Java UDF samtPython användardefinierade tabellfunktioner (UDTF).

Läs filmetadata i en UDF

Ett FILE värde har metadatafält som du kan läsa utan att öppna filen. Följande tabell innehåller de tillgängliga fälten:

Accessor Description
uri URI:n för filen.
offset En offset in i filen, i bytes.
size Filens storlek i byte.
content_type MIME-typen på filen, när den är känd.
checksum En kontrollsumma som används för att identifiera filversionen, som <algorithm>:<value>.

Åtkomst till dessa fält med punktnotation på FILE värdet, som visas i följande kod:

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;

Läs filinnehållet i en UDF

Ett FILE värde har två metoder för att läsa den underliggande filen:

  • as_local_file(): Returnerar en lokal sökväg som du kan skicka till vilket bibliotek som helst som accepterar en filväg, såsom ett bild- eller mediebibliotek.
  • open(): Returnerar en binär ström som bara läser de bytes du begär, istället för att materialisera hela filen.

Båda kräver Azure Databricks compute (en notebook- eller UDF-worker) och finns inte tillgängliga på en Azure Databricks Connect-klient. Du kan deklarera FILE som en UDF-parameter eller return-typ i Python, Scala och SQL UDF. För hela API:et, se FileType.

Extrahera bilddimensioner

Du kan använda en skalär UDF för att returnera en bilds dimensioner som en width x height sträng. UDF-anropen as_local_file() för att hämta en lokal sökväg skickar sedan den vägen till ett standardbildbibliotek (PILi Python, ImageIO i Scala), som visas i följande kod:

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

Upptäck en fils typ från dess bytes

Följande UDF läser endast de första åtta bytena av varje fil med open() och detekterar filtypen utifrån dess magiska tal, utan att materialisera hela filen:

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

Generera flera filer med en tabell-UDF (UDTF)

För att omvandla en inmatningsfil till många utdatafiler, till exempel när man delar upp en video i bildrutor, använd en tabell-UDF (UDTF). UDTF tar en FILE som indata och ger en rad per utdatafil, och skapar varje fil med FileRef.from_bytes(). Deklarera filkolumnen som FILE i UDTF:s returnType schema. För allmän UDTF-författande, se Python användardefinierade tabellfunktioner (UDTF).

När en UDTF (eller någon UDF) skriver nya filer med FileRef.from_bytes, måste din kod uppfylla följande krav:

  • Skapa målvolymen innan du kör UDTF. En Python-arbetare kan inte skapa en volym på toppnivå. Skapa den med CREATE VOLUME IF NOT EXISTS. Inuti en befintlig volym kan man os.makedirs() skapa underkataloger, men inte volymen själv.
  • Gå en absolut dbfs: väg. Att returnera a FileRef till en Delta Lake-tabell kräver en dbfs: URI, såsom dbfs:/Volumes/my_catalog/my_schema/frames/frame_00000.jpg. En bar stig reser DELTA_VIOLATE_CONSTRAINT_WITH_VALUESsig.
  • Verifiera att skrivningar är idempotenta. Radera eller hoppa över filer som redan finns innan du skriver. Eftersom FileRef.from_bytes skriver med exklusiv-skapa-flaggor, höjer FileExistsErrorskrivande över en befintlig fil .

Exempel: Extrahera videobilder

Följande UDTF läser en video FILE, extraherar varje bildruta av med (PyAV)-biblioteket, skriver den till en volym och ger en rad per bildruta:

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)

Skapa måltabellen med en FILE EXTERNAL kolumn, och anropa sedan UDTF med LATERAL för att utöka varje video till en rad per bildruta:

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;

Styr FILE-kolumner med radfilter

Styr en FILE kolumn med radfilter baserat på anroparens identitet eller filens metadata.

Radfilter

Ett radfilter är en UDF som returnerar en BOOLEAN. Rader där den returnerar false utelämnas från frågeresultaten.

Följande radfilter behåller endast rader med filer som refererar till ett Excel-kalkylblad, baserat på filens content_type metadata:

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

För mer information om hur man applicerar och hanterar radfilter, inklusive Catalog Explorer-stegen och begränsningarna, se Manuell tillämpa radfilter och kolumnmasker.

Registrera en UDF i Unity-katalogen

Registrera en filbehandlings-UDF i Unity Catalog för att styra den med katalogbehörigheter och återanvända den i anteckningsböcker, frågor och användare. Att registrera och köra en UDF kräver följande privilegier:

  • För att skapa en UDF: USAGE och CREATE på schemat, och USAGE på katalogen.
  • För att köra en UDF: EXECUTE på UDF, samt USAGE på schemat och katalogen.

Följande exempel registrerar en SQL UDF som returnerar filens filändelse och sedan anropar UDF för att skapa en ny kolumn:

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;

För att registrera en Python- eller Scala-UDF i Unity Catalog, se SQL och Python användardefinierade funktioner (UDF) i Unity Catalog och Python användardefinierade tabellfunktioner (UDTF) i Unity Catalog.

Säkerhet: UDF:er körs med ägarens rättigheter

UDF-kod körs med funktionsägarens privilegier, inte funktionsanroparens. Ägarens rättigheter gäller för att läsa bytes i en FILE. En anropare med endast EXECUTE behörigheter på UDF, och utan direkt åtkomst till den underliggande volymen, kan ändå trigga läsningar av de refererade filerna.

Eftersom en filbehandlings-UDF är en styrd åtkomstväg till filinnehåll, bör följande säkerhets- och styrningseffekter beaktas:

  • Användare kan komma åt filinnehåll med hjälp av UDF. Ge EXECUTE endast behörigheter till användare som du avser att indirekt få tillgång till filinnehållet.
  • Uppringare ärver ägarens filåtkomst. Verifiera att UDF:s ägare inte har volymåtkomst som är bredare än vad uppringare borde ha.

För mer information om hur Azure Databricks avgör den auktoriserade användaren när exekveringen korsar in i en UDF-kropp, se Auktoriserad användare och sessionsanvändare.

Nästa steg