使用 UDF 处理文件

Important

此功能在 Beta 版中。 工作区管理员可以从 预览 页控制对此功能的访问。 请参阅 Manage Azure Databricks 预览版

使用用户定义函数(UDF),借助您自己的代码和库来处理由 FILE 列引用的文件。 UDF接收每个 FILE 值作为语言原生文件引用。 在 Python 中,该引用是一个可从 pyspark.sql.types 导入的 FileRef 对象。 UDF可以读取文件的字节或作为本地路径打开,然后返回元数据值、派生文件或转换后的输出。

本页展示了Python、Scala和SQL文件处理UDF的过程。 关于 FILE 类型参考,请参见 FILE 类型。 有关常规 UDF 编写,请参见 Python 标量用户定义函数(UDF)会话范围内的 Scala 和 Java UDF,以及 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, 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;

在 UDF 中读取文件内容

一个 FILE 值有两种读取底层文件的方法:

  • as_local_file()返回一条本地路径,你可以传递给任何接受文件路径的库,比如图像库或媒体库。
  • open():返回一个二进制流,它仅读取所请求的字节,而不是将整个文件具体化。

两者都需要 Azure Databricks 计算资源(笔记本或 UDF 工作器),并且在 Azure Databricks Connect 客户端上不可用。 你可以在 Python、Scala 和 SQL 的 UDF 中声明FILE为 UDF 参数或返回类型。 完整 API 请参见 文件类型

提取图像尺寸

你可以用标量UDF来返回图像的尺寸,作为 width x height 字符串。 UDF 调用as_local_file()获取本地路径,然后将该路径传递到标准图像库(PILPython 和 ImageIO Scala 格式),如下代码所示:

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

通过字节检测文件类型

以下 UDF 仅读取带有 open() 的每个文件的前八个字节,并根据其魔数检测文件类型,而无需将整个文件物化:

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

使用表 UDF 生成多个文件(UDTF)

要将一个输入文件转换为多个输出文件,例如将视频拆分为帧时,可以使用表UDF(UDTF)。 UDTF 接受 FILE 作为输入,并为每个输出文件生成一行,使用 FileRef.from_bytes() 创建每个文件。 在 UDTF 的 returnType 架构中将文件列声明为 FILE。 有关一般 UDTF 编写的信息,请参见 Python 用户定义表函数 (UDTF)

当UDTF(或任何UDF)写入带有 FileRef.from_bytes的新文件时,你的代码必须满足以下要求:

  • 在运行UDTF之前先设定目标体积。 Python 工作者无法创建顶层卷。 用 CREATE VOLUME IF NOT EXISTS 创建它。 在现有卷内, os.makedirs() 可以创建子目录,但不能创建卷本身。
  • 传递绝对dbfs:路径。FileRef 返回到 Delta Lake 表中需要使用 dbfs: URI,例如 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: 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)

创建带有 FILE EXTERNAL 列的目标表,然后使用 LATERAL 调用 UDTF,将每个视频展开为每帧一行:

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对某个列进行管控。

行筛选器

行过滤器是一种返回 BOOLEAN 的 UDF。 返回的 false 行从查询结果中被省略。

以下行过滤器会根据文件的content_type元数据,仅保留包含引用 Excel 电子表格的文件的行:

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

有关应用和管理行筛选的更多信息,包括目录资源管理器的步骤和限制,请参见 “手动应用行筛选器和列遮罩”。

在 Unity 目录中注册 UDF

在 Unity 目录中注册一个文件处理 UDF,以便通过目录权限管理,并在笔记本、查询和用户间重复使用。 注册和运行UDF需要以下权限:

  • 要创建 UDF:需要在架构上具有 USAGECREATE,并在目录上具有 USAGE
  • 要运行 UDF:需要在该 UDF 上具有 EXECUTE,并在架构和目录上具有 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;

要在 Unity 目录中注册 Python 或 Scala UDF,请参见 Unity 目录中的 SQL 和 Python 用户定义函数(UDFs)以及 Unity 目录中的 Python 用户定义表函数(UDTF)。

安全性:UDF 在运行时使用所有者的权限

UDF 代码以函数 所有者的权限运行,而非函数调用者的权限。 所有者的权限适用于读取 FILE 的字节。 仅拥有对 UDF 的 EXECUTE 权限且无法直接访问底层卷的调用方,仍可触发对被引用文件的读取。

由于文件处理UDF是对文件内容的受控访问路径,因此考虑以下安全和治理的副作用:

  • 用户可以使用 UDF 访问文件内容。 仅将 EXECUTE 权限授予那些你打算向其授予文件内容间接访问权限的用户。
  • 来电者继承了所有者的文件访问权限。 验证 UDF 的所有者拥有的卷访问权限没有超出调用者所应有的范围。

关于 Azure Databricks 如何在执行跨入 UDF 实体时确定授权用户的更多信息,请参见授权用户和会话用户

后续步骤