Обрабатывайте файлы с UDF

Important

Эта функция доступна в бета-версии. Администраторы рабочей области могут управлять доступом к этой функции на странице "Предварительные версии ". См. статью "Управление предварительными версиями Azure Databricks".

Используйте пользовательскую функцию (UDF) для обработки файлов, на которые ссылается столбец FILE, с помощью собственного кода и библиотек. UDF получает каждое FILE значение как ссылку на файл, родной для языка. Он может читать байты файла или открывать его как локальный путь, затем возвращать значение метаданных, производный файл или преобразованный выход.

На этой странице представлены UDF для обработки файлов на Python, Scala и SQL. Сведения о типе FILE см. в разделе FILE тип. Общие сведения о создании UDF см. в разделах скалярные пользовательские функции Python (UDF), UDF Scala и Java с областью действия в пределах сеанса и пользовательские табличные функции Python (UDTF).

Чтение метаданных файла в UDF

Значение FILE имеет поля метаданных, которые можно прочитать, не открывая файл. В следующей таблице приведены доступные поля:

Средство доступа Description
uri URI файла.
offset Смещение в файле, в байтах.
size Размер файла в байтах.
content_type Тип MIME файла, когда он известен.
checksum Контрольная сумма, используемая для определения версии файла, например <algorithm>:<value>.

Обращайтесь к этим полям с помощью точечной нотации у значения FILE, как показано в следующем коде:

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;

Чтение содержимого файла в UDF

У значения FILE есть два метода для чтения базового файла:

  • as_local_file(): Возвращает локальный путь, который можно передать любой библиотеке, принимающей путь к файлу, например, к изображению или медиатеке.
  • open(): Возвращает двоичный поток, который читает только запрошенные вами байты, не загружая весь файл целиком.

Оба требуют вычислительных ресурсов Azure Databricks (ноутбук или исполнитель UDF) и недоступны в клиенте Azure Databricks Connect. Вы можете объявить FILE как параметр UDF или тип возвращаемого значения в UDF на Python, Scala и SQL. Полный API см. FileType.

Извлечь размеры изображения

Вы можете использовать скалярный UDF, чтобы вернуть размеры изображения в виде width x height строки. UDF вызывает as_local_file(), чтобы получить локальный путь, а затем передаёт этот путь стандартной библиотеке для работы с изображениями (PIL в Python, ImageIO в Scala), как показано в следующем коде:

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

Определите тип файла по его байтам

Следующий UDF с помощью open() считывает только первые восемь байт каждого файла и определяет тип файла по его сигнатуре, без материализации всего файла:

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

Генерируйте несколько файлов с помощью таблицы UDF (UDTF)

Чтобы превратить один входный файл в множество выходных, например, при разделении видео на кадры, используйте таблицу UDF (UDTF). UDTF принимает FILE в качестве входных данных и возвращает по одной строке для каждого выходного файла, создавая каждый файл с помощью FileRef.from_bytes(). Объявить столбец файла, как FILE в схеме returnType UDTF. Общие сведения о создании UDTF см. в разделе Пользовательские табличные функции (UDTF) в Python.

Когда UDTF (или любой другой UDF) записывает новые файлы с FileRef.from_bytes, ваш код должен соответствовать следующим требованиям:

  • Создайте целевой том до запуска UDTF. Сотрудник Python не может создать том верхнего уровня. Создайте это с помощью CREATE VOLUME IF NOT EXISTS. Внутри существующего тома os.makedirs() может создавать подкаталоги, но не сам том.
  • Укажите абсолютный путь dbfs:. Чтобы вернуть FileRef в таблицу Delta Lake, требуется URI dbfs:, например dbfs:/Volumes/my_catalog/my_schema/frames/frame_00000.jpg. Простой путь вызывает DELTA_VIOLATE_CONSTRAINT_WITH_VALUES.
  • Убедитесь, что операции записи идемпотентны. Удаляйте или пропускайте уже существующие файлы до написания. Поскольку FileRef.from_bytes записывает с флагами эксклюзивного создания, попытка записи поверх существующего файла вызывает FileExistsError.

Пример: Извлечение видеокадров

Следующий UDTF считывает видео FILE, извлекает каждый кадр с помощью библиотеки av (PyAV), записывает его в том и возвращает по одной строке на кадр:

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)

Создайте целевую таблицу со столбцом FILE EXTERNAL, затем вызовите UDTF с LATERAL, чтобы развернуть каждое видео в по одной строке на каждый кадр:

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 с помощью фильтров строк

Управляйте FILE столбцем с помощью фильтров строк на основе идентичности вызывающего или метаданных файла.

Фильтр строк

Строковый фильтр — это UDF, который возвращает BOOLEAN. Строки, для которых возвращается false, исключаются из результатов запроса.

Следующий фильтр строк содержит только строки с файлами, ссылающимися на Excel-таблицу, на основе метаданных файлаcontent_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)
}

Для получения дополнительной информации о применении и управлении фильтрами строк, включая шаги и ограничения Каталога Explorer, см. раздел «Ручное применение фильтров строк и масок столбцов».

Зарегистрировать UDF в каталоге Unity

Зарегистрируйте UDF для обработки файлов в Unity Catalog, чтобы управлять им с разрешениями каталога и повторно использовать его между ноутбуками, запросами и пользователями. Регистрация и запуск UDF требуют следующих привилегий:

  • Чтобы создать UDF: USAGE и CREATE для схемы и USAGE для каталога.
  • Для запуска UDF: EXECUTE на UDF и USAGE для схемы и каталога.

В следующем примере регистрируется SQL UDF, который возвращает расширение файла, затем вызывает UDF для создания нового столбца:

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;

Чтобы зарегистрировать Python- или Scala-UDF в Unity Catalog, см. пользовательские функции (UDF) SQL и Python в Unity Catalog, а также пользовательские табличные функции (UDTF) на Python в Unity Catalog.

Безопасность: UDF работают с привилегиями владельца

UDF-код работает с привилегиями владельца функции, а не вызывающего функции. Права владельца применяются к чтению байтов FILE. Вызывающий, имеющий только EXECUTE права на UDF и не имеющий прямого доступа к базовому тому, всё равно может запускать чтение ссылаемых файлов.

Поскольку UDF обработки файлов является регулируемым путем доступа к содержимому файлов, рассмотрим следующие побочные эффекты безопасности и управления:

  • Пользователи могут получать доступ к содержимому файлов с помощью UDF. Предоставляйте EXECUTE права только тем пользователям, которым вы намерены предоставить косвенный доступ к содержимому файлов.
  • Звонящие наследуют доступ к файлам владельца. Проверьте, что у владельца UDF есть доступ к громкости не шире, чем должен быть у звонящих.

Для получения дополнительной информации о том, как Azure Databricks определяет авторизованного пользователя при переходе выполнения в тело UDF, см. раздел «Авторизованный пользователь и пользователь сессии».

Дальнейшие действия