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. 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
FILEcon Auto Loader. - Analizza ogni documento con
ai_parse_documentfunzione e classificalo conai_classifyfunzione. - 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
FILEper 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
FileSpacedella 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
FileSpacedella tabella, mentre la riga della tabella memorizza solo un riferimento leggeroFILE(uri,size,content_type,checksum). Confronta questo con ilBINARYtipo, che inlinea i byte nella riga. -
Ciclo di vita dei file gestiti: Unity Catalog copia ogni file acquisito nella
FileSpacedella 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.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 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
Python
Pipeline di elaborazione file Python 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
- tipo
- Acquisisci i file come file di tipo FILE
- Funzione FILE quickstart
- Altre informazioni sul caricatore automatico. Consulta Che cos'è il caricatore automatico?.