Proses file dengan UDF

Important

Fitur ini ada di Beta. Admin ruang kerja dapat mengontrol akses ke fitur ini dari halaman Pratinjau . Lihat Kelola Pratinjau Azure Databricks.

Gunakan fungsi yang didefinisikan pengguna (UDF) untuk memproses file yang dirujuk oleh FILE kolom dengan kode dan pustaka Anda sendiri. UDF menerima setiap FILE nilai sebagai referensi file native bahasa. Ia dapat membaca byte file atau membukanya sebagai jalur lokal, lalu mengembalikan nilai metadata, file turunan, atau output yang ditransformasi.

Halaman ini menunjukkan UDF pemrosesan file dalam Python, Scala, dan SQL. Untuk FILE referensi tipe, lihat FILE tipe. Untuk pembuatan UDF umum, lihat fungsi skalar pengguna (UDF) skalar Python, Skala sesi dan UDF Java, serta fungsi tabel yang didefinisikan pengguna Python (UDTF).

Metadata file baca di UDF

Sebuah FILE nilai memiliki bidang metadata yang dapat Anda baca tanpa membuka file. Tabel berikut berisi kolom yang tersedia:

Aksesori Description
uri URI dari file.
offset Offset ke dalam file, dalam byte.
size Ukuran file, dalam satuan byte.
content_type Jenis MIME dari file, jika diketahui.
checksum Checksum yang digunakan untuk mengidentifikasi versi file, sebagai <algorithm>:<value>.

Akses kolom-kolom ini dengan notasi titik pada FILE nilai, seperti yang ditunjukkan dalam kode berikut:

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;

Baca isi file dalam UDF

Sebuah FILE nilai memiliki dua metode untuk membaca file dasar:

  • as_local_file(): Mengembalikan jalur lokal yang dapat Anda teruskan ke pustaka mana pun yang menerima jalur file, seperti gambar atau perpustakaan media.
  • open(): Mengembalikan aliran biner yang hanya membaca byte yang Anda minta, bukan seluruh file yang dimaterialisasikan.

Keduanya memerlukan komputasi Azure Databricks (notebook atau pekerja UDF) dan tidak tersedia di klien Azure Databricks Connect. Anda dapat mendeklarasikan FILE sebagai parameter UDF atau tipe return di Python, Scala, dan SQL UDF. Untuk API lengkap, lihat FileType.

Ekstrak dimensi gambar

Anda dapat menggunakan UDF skalar untuk mengembalikan dimensi gambar sebagai width x height string. UDF memanggil as_local_file() untuk mendapatkan jalur lokal, lalu meneruskan jalur tersebut ke perpustakaan citra standar (PILdi Python, ImageIO di Scala), seperti yang ditunjukkan dalam kode berikut:

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

Mendeteksi tipe file dari byte-nya

UDF berikut hanya membaca delapan byte pertama dari setiap file dengan dan open() mendeteksi tipe file dari angka ajaibnya, tanpa mematerialisasi seluruh file:

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

Menghasilkan beberapa file dengan tabel UDF (UDTF)

Untuk mengubah satu file input menjadi banyak file output, seperti saat membagi video menjadi frame, gunakan tabel UDF (UDTF). UDTF mengambil sebagai FILE input dan menghasilkan satu baris per file output, membuat setiap file dengan FileRef.from_bytes(). Deklarasikan kolom file seperti FILE dalam skema UDTF returnType . Untuk pembuatan UDTF umum, lihat fungsi tabel yang didefinisikan pengguna (UDTFs) Python.

Ketika sebuah UDTF (atau UDF apa pun) menulis file baru dengan FileRef.from_bytes, kode Anda harus memenuhi persyaratan berikut:

  • Buat volume target sebelum Anda menjalankan UDTF. Pekerja Python tidak dapat membuat volume tingkat atas. Buat dengan CREATE VOLUME IF NOT EXISTS. Di dalam volume yang sudah ada, os.makedirs() dapat membuat subdirektori, tetapi tidak volume itu sendiri.
  • Lewati jalur mutlak dbfs: . Mengembalikan a FileRef ke tabel Delta Lake memerlukan dbfs: URI, seperti dbfs:/Volumes/my_catalog/my_schema/frames/frame_00000.jpg. Jalur kosong terangkat.DELTA_VIOLATE_CONSTRAINT_WITH_VALUES
  • Verifikasi penulisan adalah idempotent. Hapus atau lewati file yang sudah ada sebelum menulis. Karena FileRef.from_bytes menulis dengan flag exclusive-create, menulis di atas file yang sudah ada akan menimbulkan FileExistsError.

Contoh: Ekstrak frame video

UDTF berikut membaca video FILE, mengekstrak setiap frame dengan av pustaka (PyAV), menulisnya ke volume, dan menghasilkan satu baris per frame:

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)

Buat tabel target dengan kolom FILE EXTERNAL , lalu panggil UDTF dengan LATERAL untuk memperluas setiap video menjadi satu baris 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;

Kelola kolom FILE dengan filter baris

Kelola kolom FILE dengan filter baris berdasarkan identitas penelepon atau metadata file.

Filter baris

Filter baris adalah UDF yang mengembalikan .BOOLEAN Baris yang dikembalikan false akan dihilangkan dari hasil kueri.

Filter baris berikut hanya menyimpan baris dengan file yang merujuk pada spreadsheet Excel, berdasarkan metadata filecontent_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)
}

Untuk informasi lebih lanjut tentang menerapkan dan mengelola filter baris, termasuk langkah dan batasan Catalog Explorer, lihat Menerapkan filter baris dan masker kolom secara manual.

Daftarkan UDF di Unity Catalog

Daftarkan UDF pemrosesan file di Unity Catalog untuk mengaturnya dengan izin katalog dan gunakan kembali di buku catatan, kueri, dan pengguna. Mendaftarkan dan menjalankan UDF memerlukan hak istimewa berikut:

  • Untuk membuat UDF: USAGE dan CREATE pada skema, serta USAGE pada katalog.
  • Untuk menjalankan UDF: EXECUTE pada UDF, serta USAGE pada skema dan katalog.

Contoh berikut mendaftarkan SQL UDF yang mengembalikan ekstensi file, lalu memanggil UDF untuk membuat kolom baru:

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;

Untuk mendaftarkan UDF Python atau Scala di Unity Catalog, lihat fungsi SQL dan Python user-defined functions (UDFs) di Unity Catalog dan Python user-defined table functions (UDTFs) di Unity Catalog.

Keamanan: UDF dijalankan dengan hak milik pemilik

Kode UDF berjalan dengan hak milik pemilik fungsi, bukan pemanggil fungsi. Hak pemilik berlaku untuk membaca byte dari sebuah FILE. Pemanggil yang hanya EXECUTE memiliki izin pada UDF, dan tidak memiliki akses langsung ke volume dasar, masih dapat memicu pembacaan file yang direferensikan.

Karena UDF pemrosesan file adalah jalur akses yang diatur ke konten file, pertimbangkan efek samping keamanan dan tata kelola berikut:

  • Pengguna dapat mengakses isi file menggunakan UDF. Berikan EXECUTE izin hanya kepada pengguna yang Anda maksudkan untuk memberikan akses tidak langsung ke isi file.
  • Pemanggil mewarisi akses file milik pemilik. Verifikasi bahwa pemilik UDF memiliki akses volume yang tidak lebih luas dari yang seharusnya dimiliki penelpon.

Untuk informasi lebih lanjut tentang bagaimana Azure Databricks menentukan pengguna yang diizinkan saat eksekusi melintasi ke badan UDF, lihat Pengguna yang diotorisasi dan pengguna sesi.

Langkah berikutnya