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. 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 Dosyaya bayt cinsinden bir ofset.
size Dosyanın bayt cinsinden boyutu.
content_type Dosyanın MIME türü, bilindiğinde.
checksum Dosya versiyonunu tanımlamak için kullanılan bir kontrol toplamı, .<algorithm>:<value>

Bu alanlara, aşağıdaki kodda gösterildiği gibi değer üzerinde nokta gösterimiyle FILE erişin:

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;

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ı çıkar

Bir görüntünün boyutlarını width x height döndürmek için 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 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()

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

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, bir giriş olarak bir FILE satır alır ve her çıkış dosyası için bir satır verir; böylece her dosya ile FileRef.from_bytes()oluşturulur. 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) yeni dosya yazdığında FileRef.from_bytes, kodunuz aşağıdaki gereksinimleri karşılamalıdır:

  • UDTF'yi çalıştırmadan önce hedef hacmi oluşturun. Bir Python çalışanı üst seviye hacim oluşturamaz. Bunu ile CREATE VOLUME IF NOT EXISTSoluşturun. Mevcut bir hacim os.makedirs() içinde alt dizinler oluşturulabilir, ancak hacim kendisi oluşturulamaz.
  • Mutlak dbfs: bir yol geç. A'yı FileRef Delta Lake tablosuna geri döndürmek için bir dbfs: URI gereklidir, örneğin dbfs:/Volumes/my_catalog/my_schema/frames/frame_00000.jpg. Çıplak bir yol yükselir DELTA_VIOLATE_CONSTRAINT_WITH_VALUES.
  • Yazıların aynı güçlüydüğünü doğrulayın. Yazmadan önce zaten var olan dosyaları silin veya atlayın. FileRef.from_bytes Çünkü özel oluştur bayraklarıyla yazıldığında, mevcut bir dosyanın üzerine yazmak .FileExistsError

Ö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):
        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 FILE yönetin.

Satır filtresi

Bir satır filtre, bir UDF'dir ve bir BOOLEAN. Döndürdüğü false satırlar sorgu sonuçlarından çıkarılmıştır.

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

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

Satır filtrelerinin uygulanması ve yönetimi hakkında, Katalog Explorer adımları ve sınırlamaları dahil olmak üzere daha fazla bilgi için El ile satır filtreleri ve sütun maskeleri uygula bkz.

Unity Kataloğunda UDF Kaydı

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 yaratmak: USAGE şema CREATE ve USAGE katalog üzerine.
  • UDF çalıştırmak için: EXECUTE UDF üzerine, USAGE şema ve katalog üzerinde.

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 ayrıcalıklarıyla çalışır

UDF kodu, fonksiyon çağırıcısının değil, fonksiyon sahibinin ayrıcalıklarıyla çalışır. Sahibin ayrıcalıkları, bir FILEbaytların okunmasına uygulanır. Yalnızca EXECUTE UDF üzerindeki izinleri olan ve altta yatan hacme doğrudan erişimi olmayan bir arayan, referans edilen dosyaların okumaları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, arayanların olması gereken hacmden daha geniş bir hacmi 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