作为文件类型进行导入文件

Important

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

FILE 类型存储并查询表格中对非结构化文件(文档、图片和音频)的引用。 本页展示了如何发现文件、将其作为 FILE 引用摄取,以及随着新文件到来时逐步导入。

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

Note

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

存储模式

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

  • FILE EXTERNAL 引用已存在于Unity Catalog卷中的文件。 Databricks 不支持存储 FILE EXTERNAL 卷外文件的引用。
  • FILE MANAGED 将文件副本存储在 Unity Catalog 管理的存储中。 来自卷外来源的文件,如 SharePoint、Google Drive 或 SFTP,必须被导入并存储为 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 表值函数

将文件作为文件引用导入

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

将卷文件导入为 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/');

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

要为源代码(如 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()

后续步骤