Pliki procesowe z UDF-ami

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. UDF otrzymuje każdą FILE wartość jako natywne odniesienie pliku dla języka. Może odczytać bajty pliku lub otworzyć go jako lokalną ścieżkę, a następnie zwrócić wartość metadanych, plik pochodny lub przekształcony wynik.

Ta strona przedstawia pliki UDF do przetwarzania plików w Python, Scala i SQL. Aby uzyskać informacje o typieFILE, zobacz FILE typ. Ogólne autorstwo UDF można znaleźć w Python scalar user defined functions (UDF),Session-scoped Scala i Java UDF oraz Python user-defined table functions (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 do 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 z notacją kropką na FILE wartości, jak pokazano w następującym kodzie:

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;

Odczytuj zawartość plików w UDF

Wartość ma FILE dwa sposoby odczytu pliku bazowego:

  • 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.

Oba wymagają obliczeń Azure Databricks (notatnika lub workera UDF) i nie są dostępne na 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.

Wymiary wyodrębnienia obrazu

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

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

Wykrywanie typu pliku na podstawie jego bajtów

Następujący UDF odczytuje tylko pierwsze osiem bajtów każdego pliku z i open() wykrywa typ pliku na podstawie jego magicznej liczby, nie materializując całego pliku:

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

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 jako wejście i FILE daje jeden wiersz na każdy plik wyjściowy, tworząc każdy plik z FileRef.from_bytes(). Zadeklaruj kolumnę pliku jako w FILE schemacie UDTF returnType . Ogólne tworzenie UDTF można znaleźć w Python user defined table functions (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 istniejącym woluminie można os.makedirs() tworzyć podkatalogi, ale nie w samym tomie.
  • Przejdź ścieżką absolutną dbfs: . Przywrócenie a FileRef do tabeli Delta Lake wymaga dbfs: URI, 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: Ekstrakt klatek wideo

Następujący UDTF odczytuje wideo FILE, wyodrębnia każdą klatkę z biblioteki av (PyAV), zapisuje ją do woluminu i daje jeden wiersz na klatkę:

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)

Stwórz docelową tabelę z kolumną FILE EXTERNAL , a następnie wywołaj UDTF , LATERAL aby rozwinąć każdy film do 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 wierszy

Filtr wierszowy to UDF, który zwraca .BOOLEAN Wiersze, dla których zwraca, 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

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

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 stworzyć UDF: USAGE oraz CREATE na schemacie i USAGE na katalogu.
  • Aby uruchomić UDF: EXECUTE na UDF, USAGE na schemacie 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ć UDF w Python lub Scala w Unity Catalog, zobacz funkcje zdefiniowane przez użytkownika SQL i Python (UDF) w Unity Catalog oraz Python user defined table functions (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 dotyczą odczytu bajtów .FILE Dzwoniący posiadający jedynie EXECUTE uprawnienia do UDF i brak bezpośredniego dostępu do bazowego woluminu może nadal wywołać odczyty odwołanych 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.
  • Dzwoniący dziedziczą dostęp do plików właściciela. 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 Azure Databricks określa użytkownika autoryzowanego podczas przejścia do ciała UDF, zobacz Użytkownik Uprawniony i użytkownik sesji.

Następne kroki