Ingesti file come tipo FILE

Importante

Questa funzionalità è in versione beta.

Il FILE tipo memorizza e interroga riferimenti a file non strutturati (documenti, immagini e audio) in tabelle. Questa pagina mostra come scoprire file, assumerli come FILE riferimenti e ingerire gradualmente nuovi file man mano che arrivano.

Per il riferimento sul FILE tipo, vedi FILE tipo. Per una panoramica degli approcci per l'assunzione di dati non strutturati, vedi FILE type e dati non strutturati.

Nota

FILE Le colonne non hanno un ordine definito. Non puoi usare una FILE colonna come colonna di partizionamento, colonna di clustering o chiave di ordine Z. Per altre informazioni, vedere Limiti.

Modalità di archiviazione

Un FILE riferimento può essere memorizzato in una delle due modalità:

  • FILE MANAGEDmemorizza copie dei file nello storage gestito da Unity Catalog: i permessi sono gestiti tramite la tabella e la cancellazione delle righe rende i file citati idonei alla garbage collection, così la tabella e i suoi file rimangono sincronizzati. I file provenienti da fonti esterne ai volumi, come SharePoint, Google Drive o SFTP, devono essere acquisiti e memorizzati come FILE MANAGED.
  • FILE EXTERNAL fa riferimento a file già esistenti in un volume di Unity Catalog. Databricks non supporta la memorizzazione FILE EXTERNAL di riferimenti per file archiviati al di fuori dei volumi.

Azure Databricks raccomanda FILE MANAGED per carichi di lavoro che beneficiano di permessi a livello di file e conformità integrata. Per un confronto tra governance e comportamento del ciclo di vita, vedi FILE type e dati non strutturati.

Uso list_files per scoprire file

Usa la list_files funzione a valori di tabella funzione a valori di tabella per scoprire i file disponibili in un percorso. Restituisce una riga per file con i relativi path, size, modification_time e riferimento FILE:

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

Per scoprire file in una sorgente che richiede una connessione Unity Catalog, come SharePoint, Google Drive o SFTP, aggiungi il parametroconnection:

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

list_files scopre i file ricorsivamente di default. Per saperne di più, vedi list_files funzione a valori di tabella.

Ingestimento dei file come riferimenti FILE

Seleziona un approccio di ingestione in base a dove conservi i tuoi file. Per assorbire file da una sorgente esterna, copiali nello storage gestito come FILE MANAGED. Per fare riferimento a file già presenti in un volume del Catalogo Unity senza copiarli, usa FILE EXTERNAL.

Importare file di origine esterni come FILE MANAGED

Per generare FILE riferimenti per file in una sorgente come SharePoint, Google Drive o SFTP, ingerisci prima i file e memorizzali come FILE MANAGED. FILE EXTERNAL Non è supportata per file archiviati fuori dai volumi.

Il seguente esempio integra file da SharePoint in una FILE MANAGED tabella:

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

Ingestire file di volume come FILE ESTERNO

Per ingerire file già esistenti in un volume del Catalogo Unity, usa un'istruzione CREATE TABLE AS SELECT (CTAS) con list_files. Questo crea una tabella con una FILE EXTERNAL colonna che fa riferimento a ogni file in posizione, senza copiarne il contenuto. Il seguente esempio crea una documents tabella con il nome del file, i metadati e un FILE riferimento per ciascun file:

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

Usa pipeline per ingerire gradualmente nuovi file

Per ingerire nuovi file man mano che arrivano, usa una tabella di streaming in una pipeline Lakeflow che legge la sorgente con STREAM read_files(..., format => 'file'). Ogni aggiornamento della pipeline elabora solo i file aggiunti dopo l'ultimo aggiornamento. Vedi read_files e Pipeline dichiarative di Spark.

Per trasmettere file in modo incrementale da una sorgente come Google Drive:

  1. Imposta il canale della pipeline su PREVIEW. Ingerire FILE riferimenti in una pipeline richiede il PREVIEW canale.

  2. Definiamo una tabella di streaming che legge la sorgente con STREAM read_files(..., format => 'file'), come nel seguente codice:

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

Applica aggiornamenti e cancellazioni con AUTO CDC

Un ingest streaming aggiunge nuovi file ma non cattura aggiornamenti o cancellazioni dalla fonte. Per applicare queste modifiche, leggi il flusso di modifica sorgente con AUTO CDC.

Avvertimento

Databricks consiglia di caricare inizialmente i dati delle modifiche in una tabella gestita, come nell'esempio seguente, quindi di applicare AUTO CDC a tale tabella. Applicare AUTO CDC direttamente a STREAM read_files(..., readChangeFeed => true) comporta la rilettura del feed delle modifiche di origine per ogni flusso a valle, il che potrebbe aumentare i costi di elaborazione.

Ingerire il feed di cambiamento in due fasi. Il seguente esempio assorbe il feed di modifiche da SharePoint, poi lo applica a una tabella di streaming target come SCD tipo 1:

  1. Scrivi i dati di modifica in una tabella di streaming con file gestiti, come nel codice seguente. Imposta readChangeFeed => true su read_files per ottenere il feed delle modifiche, che include le colonne di metadati _file_id, _sequence e _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. Usare AUTO CDC per applicare le modifiche da quella tabella a una tabella di streaming target, come nel codice seguente. Usalo _file_id come chiave, _sequence come colonna di sequenza e _is_deleted per identificare le cancellazioni.

    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
    )
    

Converti dati binari inline in riferimenti FILE

Se una tabella già memorizza il contenuto del file come dati binari inline, usa create_file la funzione per scrivere quei dati da archiviare e produrre un FILE riferimento.

I seguenti esempi utilizzano una tabella generata dall'utente, raw_documents, con una name colonna e una content colonna che contengono i dati binari.

Scrivi dati binari nello storage gestito come FILE MANAGED

Per memorizzare i file come file gestiti, chiama create_file solo con il contenuto binario. Quando ometti destination_path, Unity Catalog carica il contenuto nella posizione di archiviazione gestita:

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

Scrivi dati binari in un volume come FILE EXTERNAL

Per scrivere invece i file in un volume del Catalogo Unity come file esterni, passa destination_path a create_file, come nel seguente codice:

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

Passaggi successivi