Вводите файлы как тип FILE

Important

Эта функция доступна в бета-версии. Администраторы рабочей области могут управлять доступом к этой функции на странице "Предварительные версии ". См. статью "Управление предварительными версиями 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-valued , чтобы найти доступные файлы на определённом пути. Он возвращает одну строку на каждый файл с path, size, modification_timeи ссылкой FILE :

SELECT * FROM list_files('/Volumes/my_catalog/my_schema/raw_files/');

Чтобы найти файлы в исходном коде, требующем подключения к Unity Catalog, например, 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 Catalog без их копирования, используйте FILE EXTERNAL.

Вводите внешние исходные файлы как FILE MANAGED

Чтобы сгенерировать FILE ссылки на файлы в исходном коде, таком как SharePoint, Google Drive или SFTP, сначала загрузите файлы и храните их в 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()

Вводите файлы томов как ФАЙЛ 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

Предупреждение

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 Catalog загружает контент в управляемое хранилище:

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 как внешние, передайте 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()

Дальнейшие действия