Tutorial: Costruisci una pipeline di elaborazione file con il 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.

Impara a costruire una pipeline medaglione con Lakeflow che elabora documenti non strutturati end-to-end. 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 assume i PDF come riferimenti esterni FILE con Auto Loader, analizza ogni documento con funzioni di IA, 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à:

Il risultato è una pipeline in stile medaglione: bronzo (riferimenti esterni FILE grezzi), 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.
  • Avere permessi per creare tabelle in uno schema e per creare una pipeline.
  • Usa il canale Anteprima.

Il samples.sec.contracts dataset è disponibile di default in tutti gli spazi di lavoro, quindi non è necessaria alcuna configurazione aggiuntiva. Poiché i file sono già presenti in un volume del Catalogo Unity, questo tutorial li memorizza come FILE EXTERNAL riferimenti senza copiarne il contenuto. Per adattare la pipeline ai tuoi PDF, punta il percorso sorgente a un volume che contiene i tuoi file. Per altre opzioni di ingestione, vedi i file Ingest come tipo FILE.

Creare la pipeline di elaborazione file

Il gasdotto elabora i documenti in tre fasi.

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

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. Dichiarare la colonna come FILE EXTERNAL fa riferimento a ogni file in posizione, senza copiarne il contenuto.

SQL

CREATE OR REFRESH STREAMING TABLE raw_contracts (
  path STRING,
  size BIGINT,
  modification_time TIMESTAMP,
  file FILE EXTERNAL
)
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 EXTERNAL"
)
def raw_contracts():
  return (
    spark.readStream.format("cloudFiles")
      .option("cloudFiles.format", "file")
      .load("/Volumes/samples/sec/contracts/")
  )
  • Funziona per file grandi: un PDF grande rimane nel volume, 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.
  • 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 ciascuna FILEfunzioneai_parse_document per convertire il PDF grezzo in una struttura VARIANT che contiene 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

Definire il parse step come una vista materializzata sulla raw_contracts tabella di streaming incrementa il calcolo. Ogni aggiornamento della pipeline viene eseguito ai_parse_document solo sui file aggiunti dall'ultimo aggiornamento, non sull'intera tabella. Poiché ai_parse_document è il passo più costoso, evita di riparare documenti già elaborati. Il refresh incrementale delle viste materializzate richiede calcolo serverless; Esegui la pipeline su 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 si collega ai_classify alla versione 2.1, che restituisce la classificazione come oggetto per etichetta, quindi leggi l'etichetta dalla value chiave.

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 in un solo tipo, passa il contenuto analizzato a ai_extract funzionare con uno schema dei campi che vuoi, poi 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 assume come riferimenti ai_parse_document esterni FILE e ai_classify instrada ogni documento, e la consulting_agreements vista in oro materializzata fa emergere i campi estratti.

Esplora da solo

La pipeline classifica i documenti in cinque tipi di accordo ma estrae i campi solo consulting_agreementper . 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