Ficheiros de ingesta como tipo de ficheiro

Importante

Este recurso está em versão Beta. Os administradores do espaço de trabalho podem controlar o acesso a esse recurso na página Visualizações . Ver Gerir as pré-visualizações de Azure Databricks.

O FILE tipo armazena e consulta a ficheiros não estruturados (documentos, imagens e áudio) em tabelas. Esta página mostra como descobrir ficheiros, integrá-los como FILE referências e ingerir gradualmente novos ficheiros à medida que 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 ficheiro e dados não estruturados.

Observação

FILE As colunas não têm uma ordem definida. Não podes usar uma FILE coluna como coluna de partição, coluna de agrupamento ou chave de ordem Z. Para obter mais informações, consulte Limites.

Modos de armazenamento

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

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

Usar list_files para descobrir ficheiros

Utilize a list_files função de tabela para descobrir os ficheiros disponíveis num caminho. Devolve uma linha por ficheiro com os seus path, size, modification_time, e uma FILE referência:

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

Para descobrir ficheiros numa fonte que requer uma ligação ao 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 ficheiros recursivamente por defeito. Para saber mais, veja list_files função com valores de tabela.

Ingestão de ficheiros como referências de ficheiro

Selecione uma abordagem de ingestão com base no local onde armazena os seus ficheiros. Para referenciar ficheiros já num volume do Unity Catalog, use FILE EXTERNAL. Para ingerir ficheiros de uma fonte externa, copie-os para armazenamento gerido como FILE MANAGED.

Ingerir ficheiros de volume como FICHEIRO EXTERNO

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

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

Ingerir ficheiros fonte externos como FILE MANAGED

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

O exemplo seguinte ingere ficheiros do SharePoint numa 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()

Usar pipelines para ingerir novos ficheiros de forma incremental

Para ingerir novos ficheiros à medida que chegam, use uma tabela de streaming numa pipeline Lakeflow que lê a fonte com STREAM read_files(..., format => 'file'). Cada atualização do pipeline processa apenas os ficheiros adicionados após a última atualização. Ver read_files e Iniciar Pipelines Declarativos.

Para transmitir ficheiros de forma incremental a partir de uma fonte como o Google Drive:

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

  2. Defina uma tabela de streaming que lê 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")
      )
    

Aplicar atualizações e eliminações com o AUTO CDC

Uma ingestão de streaming adiciona novos ficheiros mas não captura atualizações ou eliminações da fonte. Para aplicar essas alterações, leia o feed de alterações de origem com AUTO CDC.

Advertência

O Databricks recomenda que coloque primeiro os dados de alteração numa tabela gerida, como no exemplo seguinte, e depois aplique AUTO CDC a essa tabela. Aplicar AUTO CDC diretamente para STREAM read_files(..., readChangeFeed => true) reler o feed de alterações de origem para cada fluxo a jusante, o que pode aumentar os custos de processamento.

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

  1. Escreva os dados de alteração numa tabela de streaming com ficheiros geridos, como no código seguinte. Defina readChangeFeed => true para read_files devolver 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, como no código seguinte. Use _file_id como chave, _sequence como coluna de sequência, e _is_deleted para identificar eliminaçõ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 FICHEIRO

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

Os exemplos seguintes utilizam uma tabela gerada pelo utilizador, raw_documents, com uma name coluna e uma content coluna que armazenam os dados binários.

Escrever dados binários num volume como FICHEIRO EXTERNO

Para escrever os ficheiros num volume do Unity Catalog como ficheiros externos, passe a destination_path para create_file, como 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()

Escrever dados binários para armazenamento gerido como FILE MANAGED

Para armazenar os ficheiros como ficheiros geridos, chame create_file apenas com o conteúdo binário. Quando omites destination_path, o Unity Catalog carrega o conteúdo para o local de armazenamento gerido:

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

Passos seguintes