Important
這項功能位於 測試版 (Beta) 中。 工作區管理員可以從 「預覽 」頁面控制對此功能的存取。 請參閱 管理 Azure Databricks 預覽。
使用使用者自訂函式(UDF)處理欄位 FILE 所參考的檔案,並搭配您自己的程式碼與函式庫。 UDF 會以語言原生檔案參考的形式接收每個 FILE 值。 它可以讀取檔案的位元組,或以本地路徑開啟,然後回傳元資料值、衍生檔案或轉換後的輸出。
本頁展示了以 Python、Scala 和 SQL 處理的檔案處理 UDF。 關於 FILE 類型參考,請參見 FILE 類型。 關於一般的 UDF 撰寫,請參見 Python 純量使用者定義函數(UDFs)、會話範圍的 Scala 與 Java UDFs,以及 Python 使用者定義表格函式(UDTFs)。
在 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 用戶端上使用。 你可以在 Python、Scala 和 SQL 的 UDF 中宣告FILE為 UDF 參數或回傳類型。 完整 API 請參見 FileType。
擷取影像尺寸
你可以用純量 UDF 來回傳影像的尺寸,以字 width x height 串形式呈現。 UDF 呼叫as_local_file()取得本地路徑,然後將該路徑傳入標準影像函式庫(PILPython 與 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() 建立每個檔案。 在 UDTF 的 FILE 結構描述中,將 file 欄位宣告為 returnType。 關於一般的 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會使用 exclusive-create 旗標寫入,所以覆寫現有檔案時會引發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 欄位的目標資料表,然後以 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 欄位
列篩選
資料列篩選器是一種會傳回 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
@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)
}
欲了解更多關於套用與管理列篩選器的資訊,包括目錄檔案總管的步驟與限制,請參閱 「手動套用列篩選器與欄遮罩」。
在 Unity Catalog 中註冊 UDF
在 Unity Catalog 註冊一個檔案處理 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 使用者定義表格函式(UDTFs)。
安全性:UDF 會以擁有者權限運作
UDF 程式碼以函式 擁有者權限執行,而非函式呼叫者。 擁有者的權限適用於讀取 FILE 的位元組。 即使呼叫者對 UDF 僅具有 EXECUTE 權限,且無法直接存取底層磁碟區,仍可觸發對所參照檔案的讀取。
由於檔案處理 UDF 是對檔案內容的受控存取路徑,請考慮以下安全與治理的副作用:
- 使用者可透過 UDF 存取檔案內容。 僅將
EXECUTE權限授予你打算讓其間接存取檔案內容的使用者。 - 來電者會繼承擁有者的檔案存取權。 確認 UDF 擁有者擁有的音量存取權不會超過來電者應有的範圍。
欲了解更多關於 Azure Databricks 如何在執行跨入 UDF 實體時判定授權使用者的資訊,請參見授權使用者與會話使用者。