Tutorial: Costruisci una pipeline di elaborazione file con il tipo FILE

Importante

Questa funzionalità è in versione beta. Per utilizzarlo, un amministratore dello spazio di lavoro deve attivare il Tipo di file dalla pagina Anteprime . Vedere Gestire le anteprime di Azure Databricks.

Scopri come creare una pipeline medallion con Lakeflow pipeline che elabora documenti non strutturati dall'inizio alla fine. Questo esempio utilizza il samples.sec.contracts dataset di esempio, una raccolta di accordi legali depositati presso la SEC memorizzati come PDF in un volume del Catalogo Unity.

La pipeline acquisisce i PDF come riferimenti gestiti FILE con Auto Loader, analizza ogni documento con funzioni AI, lo classifica in un tipo di accordo ed estrae campi strutturati per ciascun tipo.

Per il riferimento tipografico, vedi FILE tipo.

In questa esercitazione si eseguiranno le seguenti attività:

  • Acquisire in modo incrementale i PDF dei contratti da un volume come riferimenti gestiti FILE con Auto Loader.
  • Analizza ogni documento con ai_parse_document funzione e classificalo con ai_classify funzione.
  • Estrai i campi strutturati per ogni tipo di accordo con la funzione ai_extract.

Il risultato è una pipeline in stile medaglione: bronzo (riferimenti grezzi gestiti FILE ), argento (documenti analizzati e classificati) e oro (campi estratti per tipo di accordo). Per ulteriori informazioni, vedere Che cos'è l'architettura del lakehouse medallion? Il livello di bronzo è una tabella di streaming che ingerisce file in modo incrementale, mentre i livelli argento e oro sono viste materializzate che si ricalcolano solo quando i loro input cambiano.

Requisiti

Per completare questa esercitazione, è necessario soddisfare i requisiti seguenti:

  • Effettua l'accesso a uno spazio di lavoro Azure Databricks con Unity Catalog abilitato.
  • Abilita il tipo FILE per il tuo spazio di lavoro. Gli amministratori dello spazio di lavoro possono abilitarlo dalla pagina delle anteprime . Vedere Gestire le anteprime di Azure Databricks.
  • Avere permessi per creare tabelle in uno schema e per creare una pipeline.
  • Avere un volume del Catalogo Unity su cui puoi scrivere. Dichiari questo volume come FileSpace della tabella bronze e Unity Catalog vi copia i file ingeriti come storage gestito.
  • Usa il canale Anteprima.

Il samples.sec.contracts dataset è disponibile di default in tutti gli spazi di lavoro. Questo tutorial memorizza i PDF acquisiti come riferimenti FILE MANAGED: Unity Catalog copia ogni file nel volume che dichiari come FileSpace della tabella e lo gestisce insieme alla tabella, così eliminando le righe i file citati diventano idonei alla garbage collection e la tabella e i suoi file rimangono sincronizzati. Per adattare la pipeline ai tuoi PDF, punta il percorso di origine a un volume che contiene i tuoi file. Per altre opzioni di acquisizione, vedi Acquisire file come tipo FILE.

Crea la pipeline di elaborazione dei file

Il gasdotto elabora i documenti in tre fasi.

Passaggio 1: Bronzo: ingerire PDF grezzi come riferimenti FILE gestiti

Usa Auto Loader per leggere incrementalmente i PDF dei contratti dal volume. La lettura dei file con format => 'file' cattura un riferimento e metadati per ciascun file senza materializzare i suoi byte. Dichiarando la colonna come FILE MANAGED si copia ogni file nel FileSpace della tabella, il volume che imposti con la proprietà della tabella databricks.filespace-preview, così Unity Catalog gestisce i file con la tabella.

SQL

CREATE OR REFRESH STREAMING TABLE raw_contracts (
  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(
    '/Volumes/samples/sec/contracts/',
    format => 'file');

Python

from pyspark import pipelines as dp

@dp.table(
  name="raw_contracts",
  schema="path STRING, size BIGINT, modification_time TIMESTAMP, file FILE MANAGED",
  table_properties={"databricks.filespace-preview": "/Volumes/my_catalog/my_schema/filespace/"}
)
def raw_contracts():
  return (
    spark.readStream.format("cloudFiles")
      .option("cloudFiles.format", "file")
      .load("/Volumes/samples/sec/contracts/")
  )
  • Funziona per file grandi: un PDF grande si trova nel FileSpace della tabella, mentre la riga della tabella memorizza solo un riferimento leggero FILE (uri, size, content_type, checksum). Confronta questo con il BINARY tipo, che inlinea i byte nella riga.
  • Ciclo di vita dei file gestiti: Unity Catalog copia ogni file acquisito nella FileSpace della tabella e lo gestisce insieme alla tabella: eliminando le righe, i file citati diventano idonei per la garbage collection, così la tabella e i suoi file rimangono sincronizzati. Per dettagli, vedi FILE MANAGED e FILE EXTERNAL.
  • Elaborazione incrementale: la tabella di streaming ingesta gradualmente nuovi file man mano che arrivano nel sorgente, senza rielaborare quelli esistenti. Il samples.sec.contracts dataset in questo esempio è statico, ma con una fonte in tempo reale, nuovi file vengono rilevati ad ogni aggiornamento della pipeline. Per propagare anche le modifiche e le cancellazioni della sorgente, ingerire il feed di modifica con AUTO CDC. Vedi Applica aggiornamenti e cancellazioni con AUTO CDC.

Passaggio 2: Argento: analizzano e classificano i documenti

Passa ogni FILE alla ai_parse_document funzione per convertire il PDF grezzo in un VARIANT strutturato contenente elementi del documento, metadati di layout e testo. Poiché ai_parse_document accetta una FILE colonna, legge il documento direttamente dall'archiviazione e non carica mai i byte nella memoria del cluster.

SQL

CREATE OR REFRESH MATERIALIZED VIEW parsed_contracts AS
  SELECT
    path,
    ai_parse_document(file) AS parsed
  FROM raw_contracts;

Python

@dp.materialized_view(name="parsed_contracts")
def parsed_contracts():
  return (
    spark.read.table("raw_contracts")
      .selectExpr("path", "ai_parse_document(file) AS parsed")
  )

Nota

Definendo la fase di parsing come una vista materializzata sulla tabella di streaming raw_contracts, la computazione diventa incrementale. Ogni aggiornamento della pipeline viene eseguito ai_parse_document solo sui file aggiunti dall'ultimo aggiornamento, non sull'intera tabella. Poiché ai_parse_document è l'operazione più costosa, questo evita di analizzare nuovamente i documenti che hai già elaborato. L'aggiornamento incrementale delle viste materializzate richiede l'elaborazione serverless; esegui la pipeline in modalità serverless. Vedi pipeline dichiarative di Spark.

Successivamente, passa l'output analizzato alla ai_classify funzione per assegnare a ogni documento uno dei cinque tipi di accordo. I documenti con errori di analisi vengono filtrati prima della classificazione. Questo esempio blocca ai_classify alla versione 2.1, che restituisce la classificazione sotto forma di oggetto per etichetta, quindi ricava l'etichetta dalla chiave value.

SQL

CREATE OR REFRESH MATERIALIZED VIEW classified_contracts AS
  SELECT
    path,
    parsed,
    ai_classify(
      parsed,
      '["affiliate_agreement", "marketing_agreement", "consulting_agreement", "hosting_agreement", "escrow_agreement"]',
      map('version', '2.1')
    ):response[0].value::STRING AS contract_type
  FROM parsed_contracts
  WHERE is_variant_null(parsed:error_status);

Python

@dp.materialized_view(name="classified_contracts")
def classified_contracts():
  return (
    spark.read.table("parsed_contracts")
      .filter("is_variant_null(parsed:error_status)")
      .selectExpr(
        "path",
        "parsed",
        """ai_classify(
             parsed,
             '["affiliate_agreement", "marketing_agreement", "consulting_agreement", "hosting_agreement", "escrow_agreement"]',
             map('version', '2.1')
           ):response[0].value::STRING AS contract_type""")
  )

Tip

Per migliorare l'accuratezza della classificazione, aggiungi descrizioni delle etichette e un'opzione instructions a ai_classify. Vedi la funzione ai_classify.

Passaggio 3. Oro: estrazione dei campi per tipo di accordo

Ogni tipo di accordo ha il proprio insieme di campi rilevanti. Filtra i documenti classificati per un solo tipo, passa il contenuto analizzato alla funzione ai_extract con uno schema dei campi desiderati, quindi appiattisci la risposta in colonne tipizzate. Questo esempio si collega ai_extract alla versione 2.1, in cui ogni campo estratto è un oggetto, quindi leggi la sua value chiave.

Il seguente esempio costruisce la tavola d'oro per gli accordi di consulenza:

SQL

CREATE OR REFRESH MATERIALIZED VIEW consulting_agreements AS
  WITH extracted AS (
    SELECT
      path,
      ai_extract(
        parsed,
        '["company_name", "consultant_name", "compensation_amount", "effective_date"]',
        map('version', '2.1')
      ) AS fields
    FROM classified_contracts
    WHERE contract_type = 'consulting_agreement'
  )
  SELECT
    path,
    fields:response.company_name.value::STRING AS company_name,
    fields:response.consultant_name.value::STRING AS consultant_name,
    fields:response.compensation_amount.value::STRING AS compensation_amount,
    fields:response.effective_date.value::STRING AS effective_date
  FROM extracted;

Python

@dp.materialized_view(name="consulting_agreements")
def consulting_agreements():
  return (
    spark.read.table("classified_contracts")
      .filter("contract_type = 'consulting_agreement'")
      .selectExpr(
        "path",
        """ai_extract(
             parsed,
             '["company_name", "consultant_name", "compensation_amount", "effective_date"]',
             map('version', '2.1')
           ) AS fields""")
      .selectExpr(
        "path",
        "fields:response.company_name.value::STRING AS company_name",
        "fields:response.consultant_name.value::STRING AS consultant_name",
        "fields:response.compensation_amount.value::STRING AS compensation_amount",
        "fields:response.effective_date.value::STRING AS effective_date")
  )

Con queste dichiarazioni, si ottiene una pipeline completamente incrementale: man mano che nuovi PDF contrattuali arrivano nel volume, Auto Loader li acquisisce come riferimenti gestiti FILE, ai_parse_document e ai_classify instradano ogni documento, e la consulting_agreements materialized view gold fa emergere i campi estratti.

Notebook di esempio

I seguenti quaderni contengono l'intera pipeline di questo tutorial. Questi notebook sono codice sorgente pipeline, non notebook eseguibili. Importa il notebook per la tua lingua, poi specifica il suo percorso nel campo Codice sorgente quando configuri la pipeline. Vedi Configurare le pipeline.

SQL

Notebook SQL per la pipeline di elaborazione dei file

Ottieni il notebook

Python

Pipeline di elaborazione file Python notebook

Ottieni il notebook

Esplora da solo

La pipeline classifica i documenti in cinque tipi di accordo, ma estrae i campi solo per consulting_agreement. Per estenderlo, ripeti il passo oro per ogni tipo rimanente, cambiando il contract_type filtro e lo ai_extract schema per adattarli ai campi rilevanti per quel tipo. Per esempio:

  • affiliate_agreement: party_1_name, party_2_name, commission_ratepayment_frequency
  • marketing_agreement: party_1_name, party_2_name, effective_dateterritory
  • hosting_agreement: provider_name, customer_name, effective_dateterm_length
  • escrow_agreement: owner_name, licensee_name, escrow_agent_namesoftware_name

Risorse aggiuntive