Przetwarzanie plików za pomocą funkcji UDF

Important

Ta funkcja jest dostępna w wersji beta. Administratorzy obszaru roboczego mogą kontrolować dostęp do tej funkcji ze strony Podglądy . Zobacz Zarządzanie wersjami zapoznawczami usługi Azure Databricks.

Użyj funkcji zdefiniowanej przez użytkownika (UDF), aby przetwarzać pliki odwołane kolumną FILE z własnym kodem i bibliotekami. Funkcja UDF otrzymuje każdą wartość FILE w postaci natywnego dla języka odwołania do pliku. W Pythonie ta referencja to obiekt, FileRef który można zaimportować z .pyspark.sql.types UDF może odczytać bajty pliku lub otworzyć go jako lokalną ścieżkę, a następnie zwrócić wartość metadanych, wyprowadzony plik lub przekształcony wynik.

Ta strona przedstawia pliki UDF do przetwarzania plików w Python, Scala i SQL. Informacje o typie FILE można znaleźć w sekcji FILE type. Ogólne informacje na temat tworzenia funkcji UDF można znaleźć w artykułach Skalarne funkcje Python definiowane przez użytkownika (UDF), Funkcje UDF języków Scala i Java o zakresie sesji oraz Funkcje tabelaryczne Python definiowane przez użytkownika (UDTF).

Odczytuj metadane plików w UDF

Wartość zawiera pola metadanych, które można odczytać bez otwierania FILE pliku. Poniższa tabela zawiera dostępne pola:

Akcesor Description
uri URI pliku.
offset Przesunięcie w pliku, w bajtach.
size Rozmiar pliku w bajtach.
content_type Typ pliku MIME, gdy jest znany.
checksum Suma kontrolna używana do identyfikacji wersji pliku, jako <algorithm>:<value>.

Uzyskaj dostęp do tych pól za pomocą notacji kropkowej dla wartości FILE, jak pokazano w poniższym kodzie:

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;

Odczytuj zawartość plików w UDF

Wartość FILE ma dwie metody odczytywania pliku źródłowego:

  • as_local_file(): Zwraca lokalną ścieżkę, którą możesz przekazać dowolnej bibliotece akceptującej ścieżkę pliku, takiej jak biblioteka obrazów lub multimediów.
  • open(): Zwraca strumień binarny, który odczytuje tylko bajty, o które prosisz, zamiast materializować cały plik.

Obie funkcje wymagają zasobów obliczeniowych Azure Databricks (notebooka lub workera UDF) i nie są dostępne w kliencie Azure Databricks Connect. Możesz zadeklarować FILE jako parametr UDF lub typ return w UDF-ach w Python, Scala i SQL. Pełne API można znaleźć w artykule FileType.

Wyodrębnij wymiary obrazu

Możesz użyć skalarnego UDF, aby zwrócić wymiary obrazu jako width x height ciąg znaków. Funkcja UDF wywołuje as_local_file(), aby uzyskać ścieżkę lokalną, a następnie przekazuje tę ścieżkę do standardowej biblioteki do obsługi obrazów (PIL w Pythonie, ImageIO w Scali), jak pokazano w poniższym kodzie:

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

Wykrywanie typu pliku na podstawie jego bajtów

Poniższy UDF odczytuje tylko pierwsze osiem bajtów każdego pliku za pomocą open() i wykrywa typ pliku na podstawie jego sygnatury, bez materializowania całego pliku:

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

Generuj wiele plików za pomocą tabeli UDF (UDTF)

Aby przekształcić jeden plik wejściowy w wiele plików wyjściowych, na przykład podczas dzielenia wideo na klatki, użyj tabeli UDF (UDTF). UDTF przyjmuje FILE jako dane wejściowe i zwraca jeden wiersz dla każdego pliku wyjściowego, tworząc każdy plik przy użyciu FileRef.from_bytes(). Zadeklaruj kolumnę pliku jako w FILE schemacie UDTF returnType . Ogólne informacje o tworzeniu funkcji UDTF można znaleźć w sekcji Funkcje tabelaryczne zdefiniowane przez użytkownika w języku Python (UDTF).

Gdy UDTF (lub dowolny UDF) zapisuje nowe pliki z FileRef.from_bytes, Twój kod musi spełniać następujące wymagania:

  • Utwórz docelowy wolumin, zanim uruchomisz UDTF. Pracownik Python nie może utworzyć woluminu najwyższego poziomu. Stwórz go za pomocą CREATE VOLUME IF NOT EXISTS. W obrębie istniejącego woluminu os.makedirs() może tworzyć podkatalogi, ale nie sam wolumin.
  • Przejdź ścieżką absolutną dbfs: . Przywrócenie FileRef do tabeli Delta Lake wymaga identyfikatora URI dbfs:, takiego jak dbfs:/Volumes/my_catalog/my_schema/frames/frame_00000.jpg. Goła ścieżka podnosi DELTA_VIOLATE_CONSTRAINT_WITH_VALUES.
  • Weryfikuj zapisy idempotentne. Usuń lub pomiń pliki, które już istnieją, zanim go zapiszesz. Ponieważ FileRef.from_bytes zapisuje się z flagami wyłącznego tworzenia, zapisywanie przez istniejący plik podnosi FileExistsError.

Przykład: Wyodrębnianie klatek wideo

Poniższy UDTF odczytuje plik wideo FILE, wyodrębnia każdą klatkę za pomocą biblioteki av (PyAV), zapisuje ją w woluminie i zwraca jeden wiersz dla każdej klatki:

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)

Utwórz tabelę docelową z kolumną FILE EXTERNAL, a następnie wywołaj funkcję UDTF za pomocą LATERAL, aby rozwinąć każdy film do postaci jednego wiersza na klatkę:

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;

Zarządzaj kolumnami FILE za pomocą filtrów wierszowych

Zarządzaj kolumną FILEza pomocą filtrów wierszowych opartych na tożsamości wywołującego lub metadanych pliku.

Filtr wiersza

Filtr wierszy to funkcja UDF, która zwraca BOOLEAN. Wiersze, dla których zwracane jest false, są pomijane w wynikach zapytań.

Następujący filtr wierszowy przechowuje tylko wiersze z plikami odwołującymi się do arkusza Excel, na podstawie metadanych plikucontent_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, 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)
}

Aby uzyskać więcej informacji na temat stosowania i zarządzania filtrami wierszowymi, w tym krokami i ograniczeniami Eksploratora katalogu, zobacz Ręczne stosowanie filtrów wierszowych i masek kolumnowych.

Zarejestruj UDF w katalogu Unity

Zarejestruj UDF przetwarzający pliki w Unity Catalog, aby zarządzać nim uprawnieniami katalogowymi i ponownie używać go na notatnikach, zapytaniach i użytkownikach. Rejestracja i uruchomienie UDF wymaga następujących uprawnień:

  • Aby utworzyć UDF: USAGE i CREATE na schemacie oraz USAGE w katalogu.
  • Aby wykonać funkcję UDF: EXECUTE dla funkcji UDF oraz USAGE dla schematu i katalogu.

Poniższy przykład rejestruje UDF SQL, który zwraca rozszerzenie pliku, a następnie wywołuje UDF, aby utworzyć nową 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;

Aby zarejestrować funkcję UDF napisaną w języku Python lub Scala w Unity Catalog, zobacz funkcje SQL i Python zdefiniowane przez użytkownika (UDF) w Unity Catalog oraz funkcje tabelaryczne Python zdefiniowane przez użytkownika (UDTF) w Unity Catalog.

Bezpieczeństwo: UDF działają z uprawnieniami właściciela

Kod UDF działa z uprawnieniami właściciela funkcji, a nie wywołującego funkcję. Uprawnienia właściciela mają zastosowanie do odczytu bajtów obiektu FILE. Wywołujący, mający jedynie uprawnienia EXECUTE do UDF i niemający bezpośredniego dostępu do bazowego woluminu, może nadal spowodować odczyt wskazanych plików.

Ponieważ UDF przetwarzający pliki jest zarządzaną ścieżką dostępu do zawartości plików, rozważmy następujące skutki uboczne bezpieczeństwa i zarządzania:

  • Użytkownicy mogą uzyskać dostęp do zawartości plików za pomocą UDF. Udzielaj EXECUTE uprawnień tylko użytkownikom, którym zamierzasz dać pośredni dostęp do zawartości plików.
  • Wywołujący dziedziczą dostęp właściciela do plików. Sprawdź, czy właściciel UDF ma dostęp do głośności nie szerszy niż powinien mieć dzwoniący.

Aby uzyskać więcej informacji o tym, jak usługa Azure Databricks określa użytkownika autoryzowanego, gdy wykonywanie przechodzi do treści funkcji UDF, zobacz Użytkownik autoryzowany i użytkownik sesji.

Następne kroki