UDF'lerle dosya işleme

Important

Bu özellik Beta sürümündedir. Çalışma alanı yöneticileri Bu özelliğe erişimi Önizlemeler sayfasından denetleyebilir. Bkz. Azure Databricks önizlemelerini yönetme.

Bir sütunda referans FILE verilen dosyaları kendi kodu ve kütüphanelerinizle işlemek için kullanıcı tanımlı bir fonksiyon (UDF) kullanın. UDF, her değeri FILE dil tabanlı bir dosya referansı olarak alır. Python'da bu referans, FileRef'den içe aktarabileceğiniz bir pyspark.sql.types nesnesidir. UDF, dosyanın baytlarını okuyabilir veya yerel bir yol olarak açabilir, ardından meta veri değeri, türetilmiş bir dosya veya dönüştürülmüş çıktı döndürebilir.

Bu sayfa, Python, Scala ve SQL'de dosya işleme UDF'lerini gösterir. Tip FILE referansı için bkz.FILE Genel UDF yazarlığı için Python skaler kullanıcı tanımlı fonksiyonlar (UDF'ler), Oturum kapsamlı Scala ve Java UDF'ler ile Python kullanıcı tanımlı tablo fonksiyonları (UDTF'ler) sayfalarına bakınız.

UDF'de dosya meta verilerini okuma

Bir değer, FILE dosyayı açmadan okuyabileceğiniz meta veri alanlarına sahiptir. Aşağıdaki tablo mevcut alanları içerir:

Erişimci Description
uri Dosyanın URI'si.
offset Dosya içindeki bir ofset, bayt cinsinden.
size Dosyanın bayt cinsinden boyutu.
content_type Dosyanın MIME türü, bilindiğinde.
checksum Dosya sürümünü tanımlamak için kullanılan, <algorithm>:<value> gibi bir kontrol toplamı.

Aşağıdaki kodda gösterildiği gibi, FILE değeri üzerinde nokta gösterimini kullanarak bu alanlara erişin:

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;

Dosya içeriğini UDF'de okuyun

Bir değerin FILE , altta yatan dosyayı okumak için iki yöntemi vardır:

  • as_local_file(): Dosya yolunu kabul eden herhangi bir kütüphaneye, örneğin bir görüntü veya medya kütüphanesine iletebileceğiniz yerel bir yol döndürür.
  • open(): Dosyanın tamamını maddeleştirmek yerine sadece istediğiniz baytları okuyan ikili bir akış döndürür.

Her ikisi de Azure Databricks hesaplama (notebook veya UDF çalışanı) gerektirir ve Azure Databricks Connect istemcisinde mevcut değildir. Python, Scala ve SQL UDF'lerde UDF parametresi veya dönüş tipi olarak bildirebilirsinizFILE. Tam API için bkz. FileType.

Görüntü boyutlarını ayıkla

Bir görüntünün boyutlarını width x height dizesi olarak döndürmek için bir skaler UDF kullanabilirsiniz. UDF, yerel bir yol almak için çağrı as_local_file() yapar, ardından bu yolu standart bir görüntü kütüphanesine (PILPython, ImageIO Scala'da) iletir, aşağıdaki kodda gösterildiği gibi:

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

Bir dosyanın tipini baytlarından tespit et

Aşağıdaki UDF, her dosyanın yalnızca ilk sekiz baytını okur open() ve dosya tipini sihirli numarasından tespit eder, dosyanın tamamını oluşturmaz:

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

Bir tablo UDF (UDTF) ile birden fazla dosya üretin

Bir giriş dosyasını birçok çıktı dosyasına dönüştürmek için, örneğin bir videoyu karelere bölerken, bir tablo UDF (UDTF) kullanın. UDTF, girdi olarak bir FILE alır ve her çıktı dosyası için bir satır üretir; her dosyayı da FileRef.from_bytes() ile oluşturur. Dosya sütununu UDTF'nin FILEreturnType şemasında olduğu gibi ilan edin. Genel UDTF yazarlığı için bkz. Python kullanıcı tanımlı tablo fonksiyonları (UDTF'ler).

Bir UDTF (veya herhangi bir UDF) FileRef.from_bytes ile yeni dosyalar yazdığında, kodunuzun aşağıdaki gereksinimleri karşılaması gerekir:

  • UDTF'yi çalıştırmadan önce hedef hacmi oluşturun. Bir Python çalışanı üst seviye hacim oluşturamaz. Bunu CREATE VOLUME IF NOT EXISTS ile oluşturun. Mevcut bir hacim os.makedirs() içinde alt dizinler oluşturulabilir, ancak hacim kendisi oluşturulamaz.
  • Mutlak bir dbfs: yol belirtin. Bir FileRef öğesini Delta Lake tablosuna geri döndürmek için dbfs: URI'si gerekir; örneğin dbfs:/Volumes/my_catalog/my_schema/frames/frame_00000.jpg. Çıplak bir yol yükselir DELTA_VIOLATE_CONSTRAINT_WITH_VALUES.
  • Yazma işlemlerinin idempotent olduğunu doğrulayın. Yazmadan önce zaten var olan dosyaları silin veya atlayın. FileRef.from_bytes, yalnızca oluşturma bayraklarıyla yazdığı için mevcut bir dosyanın üzerine yazmak FileExistsError hatasına neden olur.

Örnek: Video karelerini çıkarma

Aşağıdaki UDTF bir video FILEokur, her kareyi ( av PyAV) kütüphanesinden çıkarır, bir hacme yazar ve her kare başına bir satır elde eder:

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)

Hedef tabloyu bir FILE EXTERNAL sütunla oluşturun, ardından UDTF'yi LATERAL çağırarak her videoyu kare başına bir satıra genişletin:

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;

FILE sütunlarını satır filtrelerle yönetin

Bir sütunu, arayanın kimliğine veya dosyanın meta verilerine göre satır filtreleriyle yönetinFILE.

Satır filtresi

Satır filtresi, bir BOOLEAN döndüren bir UDF'dir. false döndürdüğü satırlar sorgu sonuçlarına dahil edilmez.

Aşağıdaki satır filtresi, dosyanın content_type meta verilerine dayanarak yalnızca Excel elektronik tablosuna referans veren dosyaları satırları tutar:

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

Satır filtrelerini uygulama ve yönetme hakkında, Catalog Explorer adımları ve sınırlamaları da dahil olmak üzere daha fazla bilgi için Satır filtrelerini ve sütun maskelerini el ile uygulama bölümüne bakın.

Unity Kataloğunda bir UDF kaydedin

Unity Kataloğunda dosya işleme UDF'sini kaydederek katalog izinleriyle yönetin ve defterler, sorgular ve kullanıcılar arasında tekrar kullanın. UDF kaydetmek ve çalıştırmak için aşağıdaki ayrıcalıklar gereklidir:

  • Bir UDF oluşturmak için: şemada USAGE ve CREATE, katalogda ise USAGE.
  • Bir UDF'yi çalıştırmak için: UDF'de EXECUTE, şema ve katalogda USAGE.

Aşağıdaki örnek, bir SQL UDF'yi kaydeder ve dosyanın uzantısını döndürür, ardından UDF'yi çağırarak yeni bir sütun oluşturur:

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;

Bir Python veya Scala UDF'yi Unity Kataloğu'na kaydetmek için Unity Kataloğunda SQL ve Python kullanıcı tanımlı fonksiyonlar (UDF) ve Unity Katalog'da Python kullanıcı tanımlı tablo fonksiyonları (UDTF) bölümlerine bakınız.

Güvenlik: UDF'ler sahibinin yetkileriyle çalışır

UDF kodu, fonksiyon çağırıcısının değil, fonksiyon sahibinin ayrıcalıklarıyla çalışır. Sahibinin ayrıcalıkları, bir FILE nesnesinin baytlarını okumak için geçerlidir. UDF üzerinde yalnızca EXECUTE izinlerine sahip olan ve altta yatan birime doğrudan erişimi bulunmayan bir çağıran taraf, yine de başvurulan dosyaların okunmasını tetikleyebilir.

Dosya işleme UDF, dosya içeriğine yönetilen bir erişim yolu olduğundan, aşağıdaki güvenlik ve yönetişim yan etkilerini göz önünde bulundurun:

  • Kullanıcılar dosya içeriğine UDF ile erişebilirler. Yalnızca dosya içeriğine dolaylı erişim sağlamayı planladığınız kullanıcılara izin verin EXECUTE .
  • Arayanlar, sahibin dosya erişimini devralır. UDF sahibinin, çağıranların sahip olması gerekenden daha geniş bir birim erişimine sahip olmadığını doğrulayın.

Azure Databricks'in yürütmenin UDF gövdesine geçişinde yetkili kullanıcıyı nasıl belirlediğine dair daha fazla bilgi için Yetkili kullanıcı ve oturum kullanıcısı bölümlerine bakınız.

Sonraki Adımlar