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:需要在架构上具有
USAGE和CREATE,并在目录上具有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 实体时确定授权用户的更多信息,请参见授权用户和会话用户。