Ingesti file come tipo FILE

Importante

Questa funzionalità è in versione beta. Gli amministratori dell'area di lavoro possono controllare l'accesso a questa funzionalità dalla pagina Anteprime . Vedere Gestire le anteprime di Azure Databricks.

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 EXTERNAL riferimenti che già esistono in un volume del Catalogo Unity. Databricks non supporta la memorizzazione FILE EXTERNAL di riferimenti per file archiviati al di fuori dei volumi.
  • FILE MANAGED memorizza copie dei file in uno storage gestito dal catalogo Unity. I file provenienti da fonti esterne ai volumi, come SharePoint, Google Drive o SFTP, devono essere ingeriti e memorizzati come FILE MANAGED.

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 suoi path, size, modification_time, e un FILE riferimento:

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 fare riferimento a file già presenti in un volume del Catalogo Unity, usa FILE EXTERNAL. Per assorbire file da una sorgente esterna, copiali nello storage gestito come FILE MANAGED.

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/');

Ingerire file sorgente 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()

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 accendi pipeline dichiarative.

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 raccomanda di inserire prima i dati di cambiamento in una tabella gestita, come nel seguente esempio, e poi applicare AUTO CDC la tabella a quella tabella. Applicare AUTO CDC direttamente a STREAM read_files(..., readChangeFeed => true) rileggere il flusso di cambiamento sorgente 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 restituire il feed di modifica, che include le _file_idcolonne , _sequence, e _is_deleted i metadati.

    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 su un volume come FILE ESTERNO

Per scrivere i file in un volume del Catalogo Unity come file esterni, passa a destination_path a create_filea , 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()

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

Passaggi successivi