以檔案類型來導入檔案

Important

這項功能位於 測試版 (Beta) 中。 工作區管理員可以從 「預覽 」頁面控制對此功能的存取。 請參閱 管理 Azure Databricks 預覽。

此 FILE 類型用於儲存並查詢非結構化檔案(文件、圖片及音訊)的參考資料。 本頁說明如何發現檔案、將其作為 FILE 參考資料擷取,以及隨著新檔案到來逐步擷取。

關於該 FILE 類型的參考,請參見 FILE 類型。 關於擷取非結構化資料的方法概述,請參見 FILE 類型與非結構化資料。

Note

FILE 欄位沒有明確的排序。 你不能用欄位 FILE 作為分割欄位、叢集欄位或 Z 階鍵。 如需其他資訊,請參閱限制。

儲存模式

FILE參考資料可以儲存在兩種模式之一:

  • FILE MANAGED將檔案副本儲存在 Unity 目錄管理的儲存空間中:權限由資料表管理,刪除資料列後,參考檔案可進行垃圾回收,因此資料表與檔案保持同步。來自卷外來源(如 SharePoint、Google Drive 或 SFTP)的檔案必須以 FILE MANAGED.
  • FILE EXTERNAL 參考檔案中已存在於 Unity 目錄卷中。 Databricks 不支援儲存 FILE EXTERNAL 存放在卷外檔案的參考資料。

Azure Databricks 建議FILE MANAGED使用檔案層級權限與內建合規性的工作負載。 關於治理與生命週期行為的比較,請參見 FILE 類型與非結構化資料。

用 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 參考

根據你存放檔案的位置選擇擷取方式。 若要從外部來源擷取檔案,請將它們複製到管理儲存裝置中。FILE MANAGED 若要參考已在 Unity 目錄卷中的檔案而不複製檔案,請使用 FILE EXTERNAL。

以檔案管理方式匯入外部原始碼檔案

若要在 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()

將磁碟檔案匯入為 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/');

使用管線逐步擷取新檔案

要在新檔案到達時即時擷取,請使用 Lakeflow 管線中的串流表,該表讀取來源資料。STREAM read_files(..., format => 'file') 每次管線更新只處理上次更新後新增的檔案。 請參閱 read_files 並 啟動宣告式管線。

要從像 Google Drive 這類來源逐步串流檔案:

  1. 將管線的通道設定為 PREVIEW。 在管線中攝FILE取參考需要通道。PREVIEW

  2. 定義一個串流表,讀取來源, 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:

  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/")
      )
    
  2. 請使用 AUTO CDC 該資料表的變更套用到目標串流資料表,如下程式碼所示。 以鍵_file_id、序列欄位_sequence、_is_deleted識別刪除為鍵。

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

下一步