以檔案類型來導入檔案

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) 函式來發現路徑上可用的檔案。 它會回傳每個檔案一列,包含其 pathsizemodification_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 MANAGEDFILE 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 這類來源逐步串流檔案:

  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 CDCSTREAM read_files(..., readChangeFeed => true) 重讀每個下游流的來源變更導流,可能會增加處理成本。

分兩步來接收變更訂閱。 以下範例是從 SharePoint 匯入變更資料流,然後將其套用到目標串流表中,作為 SCD 類型 1:

  1. 將變更資料寫入帶有受管理檔案的串流表,如下程式碼所示。 設定 readChangeFeed => trueread_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 該資料表的變更套用到目標串流資料表,如下程式碼所示。 以鍵_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_pathcreate_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()

下一步