将文件作为 FILE 类型导入

Important

此功能在 Beta 版中。 工作区管理员可以从 预览 页控制对此功能的访问。 请参阅 Manage Azure Databricks 预览版。

该 FILE 类型存储并查询表格中对非结构化文件(文档、图片和音频)的引用。 本页介绍如何发现文件、将其作为 FILE 引用导入,以及在新文件到达时对其进行增量导入。

关于该 FILE 类型的参考,请参见 FILE 类型。 关于非结构化数据摄取方法的概述,请参见 FILE 类型和非结构化数据。

注释

FILE 列没有固定的顺序。 你不能用列 FILE 作为划分列、聚类列或Z阶键。 有关详细信息,请参阅 限制。

存储模式

FILE参考可以以两种模式之一存储:

  • FILE MANAGED将文件副本存储在 Unity 目录管理的存储中:权限通过表管理,删除行后,引用的文件有资格进行垃圾回收,因此表及其文件保持同步。来自卷外来源的文件,如 SharePoint、Google Drive 或 SFTP,必须被导入并存储为 FILE MANAGED。
  • FILE EXTERNAL 引用已存在在 Unity Catalog 卷中的文件。 Databricks 不支持存储 FILE EXTERNAL 卷外文件的引用。

Azure Databricks 推荐FILE MANAGED,适用于可从文件级权限和内建合规支持中受益的工作负载。 有关治理和生命周期行为的比较,请参见 文件类型和非结构化数据。

使用 list_files 来发现文件

使用 list_files 表值函数 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 带有 list_files的 (CTAS) 语句。 这会创建一个包含 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 和 Spark 声明式管道。

要从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 读取源更改馈送。

警告

Databricks 建议你先将变更数据落在一个受管理的表中,如下例所示,然后应用 AUTO CDC 到该表。 将 AUTO CDC 直接应用于 STREAM read_files(..., readChangeFeed => true) 会针对每个下游流重新读取源更改馈送,这可能会增加处理成本。

分两步导入变更信息流。 以下示例从 SharePoint 摄取变更源,然后将其以 SCD 类型 1 的方式应用于目标流式表:

  1. 将变更数据写入带有管理文件的流表中,如下代码所示。 在 read_files 上将 readChangeFeed => true 设置为返回更改源,其中包括 _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 Catalog 卷,请将 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()

后续步骤