Ingeste arquivos como o tipo FILE

Importante

Esse recurso está em Beta. Os administradores do workspace podem controlar o acesso a esse recurso na página Visualizações . Consulte Gerenciar visualizações do Azure Databricks.

O FILE tipo armazena e consulta referências a arquivos não estruturados (documentos, imagens e áudio) em tabelas. Esta página mostra como descobrir arquivos, ingeri-los como FILE referências e ingerir novos arquivos gradualmente à medida que eles chegam.

Para a referência sobre o FILE tipo, veja FILE tipo. Para uma visão geral das abordagens para ingerir dados não estruturados, veja TIPO de ARQUIVO e dados não estruturados.

Observação

FILE Colunas não têm uma ordem definida. Você não pode usar uma FILE coluna como coluna de particionamento, coluna de agrupamento ou chave de ordem Z. Para obter mais informações, confira Limites.

Modos de armazenamento

Uma FILE referência pode ser armazenada em um de dois modos:

  • FILE EXTERNAL referencias que já existem em um volume do Catálogo Unity. O Databricks não suporta armazenar FILE EXTERNAL referências para arquivos armazenados fora de volumes.
  • FILE MANAGED armazena cópias de arquivos em armazenamento gerenciado pelo Catálogo Unity. Arquivos de fontes fora de volumes, como SharePoint, Google Drive ou SFTP, devem ser ingeridos e armazenados como FILE MANAGED.

Uso list_files para descobrir arquivos

Use a list_files função tabela para descobrir os arquivos disponíveis em um caminho. Ele retorna uma linha por arquivo com seu path, size, modification_time, e uma FILE referência:

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

Para descobrir arquivos em uma fonte que requer uma conexão com o Unity Catalog, como SharePoint, Google Drive ou SFTP, adicione o connection parâmetro:

SELECT * FROM list_files('https://example.sharepoint.com/sites/my-site/', connection => 'my_sharepoint_connection');

list_files descobre arquivos recursivamente por padrão. Para saber mais, veja list_files função com valores de tabela.

Ingesta como referências de ARQUIVO

Selecione uma abordagem de ingestão com base em onde você armazena seus arquivos. Para referenciar arquivos já em um volume do Unity Catalog, use FILE EXTERNAL. Para ingerir arquivos de uma fonte externa, copie-os para o armazenamento gerenciado como FILE MANAGED.

Ingir arquivos de volume como ARQUIVO EXTERNO

Para ingerir arquivos que já existem em um volume do Unity Catalog, use uma CREATE TABLE AS SELECT instrução (CTAS) com list_files. Isso cria uma tabela com uma FILE EXTERNAL coluna que faz referência a cada arquivo no local, sem copiar seu conteúdo. O exemplo a seguir cria uma documents tabela com o nome do arquivo, metadados e uma FILE referência para cada arquivo:

CREATE TABLE documents AS
  SELECT _metadata.file_name, *
  FROM list_files('/Volumes/my_catalog/my_schema/raw_files/');

Ingira arquivos fonte externos como FILE MANAGED

Para gerar FILE referências para arquivos em uma fonte, como SharePoint, Google Drive ou SFTP, ingira-se primeiro os arquivos e armazene-os como FILE MANAGED. FILE EXTERNAL Não é suportado para arquivos armazenados fora de volumes.

O exemplo a seguir ingire arquivos do SharePoint em uma FILE MANAGED tabela:

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

Use pipelines para ingerir novos arquivos de forma incremental

Para ingerir novos arquivos conforme eles chegam, use uma tabela de streaming em um pipeline Lakeflow que leia a fonte com STREAM read_files(..., format => 'file'). Cada atualização do pipeline processa apenas os arquivos adicionados após a última atualização. Veja read_files e Estimule Pipelines Declarativos.

Para transmitir arquivos incrementalmente de uma fonte como o Google Drive:

  1. Defina o canal do pipeline para PREVIEW. Ingerir FILE referências em um pipeline requer o PREVIEW canal.

  2. Defina uma tabela de streaming que leia a fonte com STREAM read_files(..., format => 'file'), como no seguinte código:

    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")
      )
    

Aplique atualizações e excluções com o AUTO CDC

Uma ingesta de streaming adiciona novos arquivos, mas não captura atualizações ou excluções da fonte. Para aplicar essas mudanças, leia o feed de alterações de origem com AUTO CDC.

Aviso

O Databricks recomenda que você coloque os dados de alteração em uma tabela gerenciada primeiro, como no exemplo a seguir, e depois aplique AUTO CDC a essa tabela. Aplicar AUTO CDC diretamente para STREAM read_files(..., readChangeFeed => true) reler o feed de mudança de fonte para cada fluxo a jusante, o que pode aumentar os custos de processamento.

Ingire o feed de mudança em dois passos. O exemplo a seguir ingere o feed de alterações do SharePoint e o aplica a uma tabela de streaming alvo como SCD tipo 1:

  1. Escreva os dados de alteração em uma tabela de streaming com arquivos gerenciados, conforme no código a seguir. Defina readChangeFeed => true em read_files Ativo para retornar o feed de alterações, que inclui as _file_idcolunas , _sequence, e _is_deleted metadados.

    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. Use AUTO CDC para aplicar as alterações dessa tabela a uma tabela de streaming alvo, conforme no código a seguir. Use _file_id como chave, _sequence como coluna de sequência, e _is_deleted para identificar exclusões.

    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
    )
    

Converter dados binários inline em referências de ARQUIVO

Se uma tabela já armazena o conteúdo do arquivo como dados binários inline, use create_file a função para gravar esses dados no armazenamento e produzir uma FILE referência.

Os exemplos a seguir usam uma tabela gerada pelo usuário, raw_documents, com uma name coluna e uma content coluna que armazenam os dados binários.

Gravar dados binários em um volume como ARQUIVO EXTERNO

Para escrever os arquivos em um volume do Unity Catalog como arquivos externos, passe a destination_path para create_file, conforme no seguinte código:

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

Escreva dados binários no armazenamento gerenciado como FILE MANAGED

Para armazenar os arquivos como arquivos gerenciados, chame create_file apenas com o conteúdo binário. Quando você omite destination_path, o Unity Catalog envia o conteúdo para o local de armazenamento gerenciado:

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

Próximas Etapas