Procesní soubory s UDF

Important

Tato funkce je v beta verzi. Správci pracovního prostoru můžou řídit přístup k této funkci ze stránky Previews . Viz Manage Azure Databricks preview.

Použijte uživatelem definovanou funkci (UDF) k zpracování souborů, na které odkazuje sloupec, FILE pomocí vlastního kódu a knihoven. UDF přijímá každou FILE hodnotu jako nativní odkaz na soubor v jazyce. Může číst bajty souboru nebo jej otevřít jako lokální cestu, poté vrátit hodnotu metadat, odvozený soubor nebo transformovaný výstup.

Tato stránka ukazuje UDF pro zpracování souborů v Python, Scala a SQL. Pro typovou referenci FILE viz FILE typ. Pro obecné UDF authoring viz Python skalární uživatelsky definované funkce (UDF),relace v Scala a Java UDF a Python uživatelsky definované tabulkové funkce (UDTF).

Čtěte metadata souborů v UDF

FILE Hodnota má metadata, která lze číst bez nutnosti otevírat soubor. Následující tabulka obsahuje dostupná pole:

Přístupové Description
uri URI souboru.
offset Offset do souboru v bajtech.
size Velikost souboru v bajtech.
content_type MIME typ souboru, pokud je znám.
checksum Kontrolní součet používaný k identifikaci verze souboru, jako <algorithm>:<value>.

Přistupujte k těmto polům s tečkovou notací na hodnotě FILE , jak je ukázáno v následujícím kódu:

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;

Čtěte obsah souboru v UDF

Hodnota má FILE dva způsoby čtení podkladového souboru:

  • as_local_file(): Vrací lokální cestu, kterou můžete předat jakékoli knihovně, jež přijímá cestu k souboru, například obrazové nebo mediální knihovně.
  • open(): Vrací binární proud, který čte pouze ty bajty, které požadujete, místo aby materializoval celý soubor.

Oba vyžadují Azure Databricks compute (notebook nebo UDF worker) a nejsou dostupné na klientovi Azure Databricks Connect. Můžete deklarovat FILE jako UDF parametr nebo vrátit typ v Python, Scala a SQL UDF. Pro kompletní API viz FileType.

Extrahování rozměrů obrazu

Můžete použít skalární UDF k vrácení rozměrů obrazu jako width x height řetězce. UDF volá as_local_file() pro získání lokální cesty a poté tuto cestu předá standardní obrazové knihovně (PILv Pythonu, ImageIO ve Scale), jak je ukázáno v následujícím kódu:

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

Detekujte typ souboru podle jeho bajtů

Následující UDF čte pouze prvních osm bajtů každého souboru a open() detekuje typ souboru podle jeho magického čísla, aniž by materializoval celý soubor:

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

Generujte více souborů pomocí tabulky UDF (UDTF)

Pro přeměnu jednoho vstupního souboru na více výstupních souborů, například při rozdělení videa na snímky, použijte tabulku UDF (UDTF). UDTF přijme a FILE jako vstup a dá jeden řádek na výstupní soubor, čímž každý soubor vytvoří s .FileRef.from_bytes() Deklarujte sloupec souboru jako ve FILE schématu UDTF returnType . Pro obecné UDTF authoring viz Python uživatelsky definované tabulkové funkce (UDTF).

Když UDTF (nebo jakýkoli UDF) zapisuje nové soubory s FileRef.from_bytes, váš kód musí splňovat následující požadavky:

  • Vytvořte cílový objem ještě před spuštěním UDTF. Pracovník Python nemůže vytvořit nejvyšší úroveň objemu. Vytvořte ji pomocí CREATE VOLUME IF NOT EXISTS. Uvnitř existujícího svazku lze os.makedirs() vytvářet podadresáře, ale ne samotný svazek.
  • Projděte absolutní dbfs: cestou. Návrat a FileRef do tabulky Delta Lake vyžaduje dbfs: URI, například dbfs:/Volumes/my_catalog/my_schema/frames/frame_00000.jpg. Holá cesta zvedá DELTA_VIOLATE_CONSTRAINT_WITH_VALUES.
  • Ověřte, že zápisy jsou idempotentní. Smažte nebo přeskočte soubory, které už existují, před zápisem. Protože FileRef.from_bytes zápis s exkluzivními příznaky vytváření, zápis přes existující soubor vyvolává FileExistsError.

Příklad: Extrahování video snímků

Následující UDTF čte video FILE, extrahuje každý snímek pomocí knihovny av (PyAV), zapisuje jej do svazku a dává jeden řádek na snímek:

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)

Vytvořte cílovou tabulku ve sloupci FILE EXTERNAL a pak zavolejte UDTF s , LATERAL abyste každé video rozbalili do jednoho řádku na snímek:

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;

Řízení sloupců FILE pomocí řádkových filtrů

Řídit sloupec FILE pomocí řádkových filtrů založených na identitě volajícího nebo metadatech souboru.

Filtr řádků

Řádkový filtr je UDF, který vrací .BOOLEAN Řádky, pro které vrací, false jsou z výsledků dotazu vynechány.

Následující řádkový filtr uchovává pouze řádky se soubory, které odkazují na Excel tabulku, na základě metadat souborucontent_type:

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

Pro více informací o aplikaci a správě řádkových filtrů, včetně kroků a omezení Průzkumníka katalogu, viz Ručně aplikovat řádkové filtry a masky sloupců.

Registrujte UDF v Unity Catalog

Zaregistrujte UDF pro zpracování souborů v Unity Catalog, abyste jej spravovali s katalogovými oprávněními a znovu jej používali napříč zápisníky, dotazy a uživateli. Registrace a provoz UDF vyžaduje následující oprávnění:

  • Vytvořit UDF: USAGE a CREATE na schématu, a USAGE na katalogu.
  • Pro spuštění UDF: EXECUTE na UDF, USAGE na schématu a katalogu.

Následující příklad zaregistruje SQL UDF, které vrátí příponu souboru, a poté zavolá UDF pro vytvoření nového sloupce:

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;

Pro registraci Python nebo Scala UDF v Unity Catalog viz SQL a Python uživatelsky definované funkce (UDF) v Unity Catalog a Python uživatelsky definované tabulkové funkce (UDTF) v Unity Catalog.

Bezpečnost: UDF běží s oprávněními vlastníka

UDF kód běží s oprávněními vlastníka funkce, nikoli volajícího funkce. Oprávnění vlastníka se vztahují na čtení bajtů .FILE Volající s pouze oprávněními na EXECUTE UDF a bez přímého přístupu k podkladovému svazku může stále spustit čtení odkazovaných souborů.

Protože UDF pro zpracování souborů je řízená přístupová cesta k obsahu souborů, zvažte následující vedlejší účinky bezpečnosti a správy:

  • Uživatelé mohou přistupovat k obsahu souborů pomocí UDF. Udělujte EXECUTE oprávnění pouze uživatelům, kterým chcete dát nepřímý přístup k obsahu souborů.
  • Volající dědí přístup k souborům od vlastníka. Ověřte, že vlastník UDF nemá širší přístup k objemu, než by měli mít volající.

Pro více informací o tom, jak Azure Databricks určuje oprávněného uživatele při přechodu do UDF těla, viz Autorizovaný uživatel a uživatel relace.

Další kroky