Important
這項功能位於 測試版 (Beta) 中。 工作區管理員可以從 「預覽 」頁面控制對此功能的存取。 請參閱 管理 Azure Databricks 預覽。
此 FILE 類型用於儲存並查詢非結構化檔案(文件、圖片及音訊)的參考資料。 本頁說明如何發現檔案、將其作為 FILE 參考資料擷取,以及隨著新檔案到來逐步擷取。
關於該 FILE 類型的參考,請參見 FILE 類型。 關於擷取非結構化資料的方法概述,請參見 FILE 類型與非結構化資料。
Note
FILE 欄位沒有明確的排序。 你不能用欄位 FILE 作為分割欄位、叢集欄位或 Z 階鍵。 如需其他資訊,請參閱限制。
儲存模式
FILE參考資料可以儲存在兩種模式之一:
-
FILE EXTERNAL參考檔案中已存在於 Unity 目錄卷中。 Databricks 不支援儲存FILE EXTERNAL存放在卷外檔案的參考資料。 -
FILE MANAGED將檔案副本儲存在 Unity 目錄管理的儲存空間中。 來自卷外來源(如 SharePoint、Google Drive 或 SFTP)的檔案必須以FILE MANAGED.
用 list_files 來發現檔案
使用 list_files table-value(table-value) 函式來發現路徑上可用的檔案。 它會回傳每個檔案一列,包含其 path、 size、 modification_time及一個 FILE 參考:
SELECT * FROM list_files('/Volumes/my_catalog/my_schema/raw_files/');
要在需要 Unity 目錄連線的來源中找到檔案,例如 SharePoint、Google Drive 或 SFTP,請新增connection參數:
SELECT * FROM list_files('https://example.sharepoint.com/sites/my-site/', connection => 'my_sharepoint_connection');
list_files 預設以遞迴方式發現檔案。 欲了解更多,請參閱 list_files 表值函數。
將檔案匯入為 FILE 參考
根據你存放檔案的位置選擇擷取方式。 若要參考 Unity 目錄中已有的檔案,請使用 FILE EXTERNAL. 若要從外部來源擷取檔案,請將它們複製到管理儲存裝置中。FILE MANAGED
將磁碟檔案匯入為 FILE EXTERNAL
若要匯入已存在於 Unity 目錄卷中的檔案,請使用 CREATE TABLE AS SELECT (CTAS) 陳述式。list_files 這會建立一個帶有 FILE EXTERNAL 欄位的表格,該欄位會參考每個檔案,且不會複製其內容。 以下範例建立 documents 一個包含檔案名稱、元資料及 FILE 每個檔案參考的資料表:
CREATE TABLE documents AS
SELECT _metadata.file_name, *
FROM list_files('/Volumes/my_catalog/my_schema/raw_files/');
以檔案管理方式匯入外部原始碼檔案
若要在 SharePoint、Google Drive 或 SFTP 等來源產生FILE檔案參考,請先擷取檔案並儲存為 FILE MANAGED。
FILE EXTERNAL 不支援存放在磁碟區外的檔案。
以下範例是將 SharePoint FILE MANAGED 的檔案匯入資料表:
SQL
CREATE TABLE managed_documents (
file_name STRING,
path STRING,
size BIGINT,
modification_time TIMESTAMP,
file FILE MANAGED
) USING DELTA
TBLPROPERTIES ('databricks.filespace-preview' = '/Volumes/my_catalog/my_schema/filespace/');
INSERT INTO managed_documents
SELECT _metadata.file_name, *
FROM read_files(
'https://example.sharepoint.com/sites/my-site/',
connection => 'my_sharepoint_connection',
format => 'file');
Python
(spark.read.format("file")
.option("databricks.connection", "my_sharepoint_connection")
.load("https://example.sharepoint.com/sites/my-site/")
.selectExpr("_metadata.file_name", "*")
.writeTo("managed_documents").append())
Scala
spark.read.format("file")
.option("databricks.connection", "my_sharepoint_connection")
.load("https://example.sharepoint.com/sites/my-site/")
.selectExpr("_metadata.file_name", "*")
.writeTo("managed_documents").append()
使用管線逐步擷取新檔案
要在新檔案到達時即時擷取,請使用 Lakeflow 管線中的串流表,該表讀取來源資料。STREAM read_files(..., format => 'file') 每次管線更新只處理上次更新後新增的檔案。 請參閱 read_files 並 啟動宣告式管線。
要從像 Google Drive 這類來源逐步串流檔案:
將管線的通道設定為
PREVIEW。 在管線中攝FILE取參考需要通道。PREVIEW定義一個串流表,讀取來源,
STREAM read_files(..., format => 'file')如下程式碼所示:SQL
CREATE STREAMING TABLE streaming_documents ( path STRING, size BIGINT, modification_time TIMESTAMP, file FILE MANAGED ) TBLPROPERTIES ('databricks.filespace-preview' = '/Volumes/my_catalog/my_schema/filespace/') AS SELECT * FROM STREAM read_files( 'https://drive.google.com/drive/folders/my-folder-id', connection => 'my_gdrive_connection', format => 'file');Python
from pyspark import pipelines as dp @dp.table( name="streaming_documents", schema="path STRING, size BIGINT, modification_time TIMESTAMP, file FILE MANAGED", table_properties={"databricks.filespace-preview": "/Volumes/my_catalog/my_schema/filespace/"} ) def streaming_documents(): return ( spark.readStream.format("cloudFiles") .option("cloudFiles.format", "file") .option("databricks.connection", "my_gdrive_connection") .load("https://drive.google.com/drive/folders/my-folder-id") )
透過 AUTO CDC 套用更新與刪除
串流擷取會新增檔案,但不會擷取來源的更新或刪除。 要套用這些變更,請讀取來源變更訂閱源。AUTO CDC
Warning
Databricks 建議你先將變更資料放入受管理的資料表,如以下範例所示,然後再套用 AUTO CDC 該資料表。 直接應用 AUTO CDC 於 STREAM read_files(..., readChangeFeed => true) 重讀每個下游流的來源變更導流,可能會增加處理成本。
分兩步來接收變更訂閱。 以下範例是從 SharePoint 匯入變更資料流,然後將其套用到目標串流表中,作為 SCD 類型 1:
將變更資料寫入帶有受管理檔案的串流表,如下程式碼所示。 設定
readChangeFeed => true為read_files返回變更資料,包含_file_id、_sequence和_is_deleted元資料欄位。SQL
CREATE OR REFRESH STREAMING TABLE documents_changes ( _file_id STRING, _sequence BIGINT, _is_deleted BOOLEAN, path STRING, size BIGINT, modification_time TIMESTAMP, file FILE MANAGED ) TBLPROPERTIES ('databricks.filespace-preview' = '/Volumes/my_catalog/my_schema/filespace/') AS SELECT * FROM STREAM read_files( 'https://example.sharepoint.com/sites/my-site/', connection => 'my_sharepoint_connection', format => 'file', readChangeFeed => true);Python
from pyspark import pipelines as dp @dp.table( name="documents_changes", table_properties={"databricks.filespace-preview": "/Volumes/my_catalog/my_schema/filespace/"} ) def documents_changes(): return ( spark.readStream.format("cloudFiles") .option("cloudFiles.format", "file") .option("databricks.connection", "my_sharepoint_connection") .option("cloudFiles.readChangeFeed", "true") .load("https://example.sharepoint.com/sites/my-site/") )請使用
AUTO CDC該資料表的變更套用到目標串流資料表,如下程式碼所示。 以鍵_sequence、序列欄位_is_deleted、_file_id識別刪除為鍵。SQL
CREATE OR REFRESH STREAMING TABLE documents TBLPROPERTIES ('databricks.filespace-preview' = '/Volumes/my_catalog/my_schema/filespace/'); CREATE FLOW documents_cdc AS AUTO CDC INTO documents FROM STREAM documents_changes KEYS (_file_id) APPLY AS DELETE WHEN _is_deleted = true SEQUENCE BY _sequence COLUMNS * EXCEPT (_is_deleted, _sequence) STORED AS SCD TYPE 1;Python
from pyspark import pipelines as dp from pyspark.sql.functions import col, expr dp.create_streaming_table( name="documents", table_properties={"databricks.filespace-preview": "/Volumes/my_catalog/my_schema/filespace/"} ) dp.create_auto_cdc_flow( target = "documents", source = "documents_changes", keys = ["_file_id"], sequence_by = col("_sequence"), apply_as_deletes = expr("_is_deleted = true"), except_column_list = ["_is_deleted", "_sequence"], stored_as_scd_type = 1 )
將內嵌二進位資料轉換為 FILE 參考
如果資料表已經以內嵌二進位資料儲存檔案內容,請使用 create_file 函式 將該資料寫入儲存並產生 FILE 參考資料。
以下範例使用使用者產生的表格 raw_documents,包含一 name 欄與 content 一欄存放二進位資料。
將二進位資料寫入磁碟區,格式為 FILE EXTERNAL
若要將檔案寫入 Unity 目錄卷,作為外部檔案,請如以下程式碼傳送 a destination_path 至 create_file, :
SQL
CREATE TABLE documents (name STRING, file FILE EXTERNAL) USING DELTA;
INSERT INTO documents (name, file)
SELECT
name,
create_file(
content => content,
destination_path => '/Volumes/my_catalog/my_schema/my_volume/' || name
)
FROM raw_documents;
Python
(spark.read.table("raw_documents")
.selectExpr(
"name",
"create_file(content => content, destination_path => '/Volumes/my_catalog/my_schema/my_volume/' || name) AS file")
.writeTo("documents").append())
Scala
spark.read.table("raw_documents")
.selectExpr(
"name",
"create_file(content => content, destination_path => '/Volumes/my_catalog/my_schema/my_volume/' || name) AS file")
.writeTo("documents").append()
將二進位資料寫入受管理儲存裝置為 FILE MANAGED
若要將檔案存為受管理檔案,請只用二進位內容呼叫 create_file 。 當你省略 destination_path時,Unity 目錄會將內容上傳到受管理的儲存位置:
SQL
CREATE TABLE managed_documents (name STRING, file FILE MANAGED) USING DELTA
TBLPROPERTIES ('databricks.filespace-preview' = '/Volumes/my_catalog/my_schema/filespace/');
INSERT INTO managed_documents (name, file)
SELECT name, create_file(content => content)
FROM raw_documents;
Python
(spark.read.table("raw_documents")
.selectExpr("name", "create_file(content => content) AS file")
.writeTo("managed_documents").append())
Scala
spark.read.table("raw_documents")
.selectExpr("name", "create_file(content => content) AS file")
.writeTo("managed_documents").append()
下一步
-
FILE類型 - FILE 類型與非結構化資料
- 教學:建立一個帶有 FILE 類型的檔案處理管線
- 了解更多關於自動裝填機的資訊。 請參閱 什麼是自動載入器?。