Nota
L'accesso a questa pagina richiede l'autorizzazione. È possibile provare ad accedere o modificare le directory.
L'accesso a questa pagina richiede l'autorizzazione. È possibile provare a modificare le directory.
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à:
- Ingerire incrementalmente PDF contrattuali da un volume come riferimenti esterni
FILEcon Auto Loader. - Analizza ogni documento con
ai_parse_documentfunzione e classificalo conai_classifyfunzione. - Estrai i campi strutturati per ogni tipo di accordo con
ai_extractfunzione.
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 ilBINARYtipo, 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.contractsdataset 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 conAUTO 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
- tipo
- Ingesti file come tipo FILE
- Funzione FILE quickstart
- Altre informazioni sul caricatore automatico. Consulta Che cos'è il caricatore automatico?.