将文件作为 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 函数来发现路径上可用的文件。 它为每个文件返回一行,其中包含其 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 引用导入

根据你存放文件的位置选择一种摄取方式。 要从外部来源导入文件,将文件复制到托管存储中。FILE MANAGED 要引用Unity目录卷中已有的文件而不复制,请使用 FILE EXTERNAL

以文件管理方式导入外部源文件

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

以 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_filesSpark 声明式管道

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

后续步骤