Hinweis
Für den Zugriff auf diese Seite ist eine Autorisierung erforderlich. Sie können versuchen, sich anzumelden oder das Verzeichnis zu wechseln.
Für den Zugriff auf diese Seite ist eine Autorisierung erforderlich. Sie können versuchen, das Verzeichnis zu wechseln.
Important
Dieses Feature befindet sich in der Betaversion. Arbeitsbereichsadministratoren können den Zugriff auf dieses Feature über die Vorschauseite steuern. Siehe Manage Azure Databricks Previews.
Lernen Sie, wie Sie mit Lakeflow Pipeline eine Medallion-Pipeline aufbauen, die unstrukturierte Dokumente von Anfang bis Ende verarbeitet. Dieses Beispiel verwendet den samples.sec.contracts Beispieldatensatz, eine Sammlung von von SEC eingereichten Rechtsverträgen, die als PDFs in einem Unity-Katalog-Volume gespeichert sind.
Die Pipeline nimmt die PDFs als externe FILE Referenzen mit Auto Loader ein, parst jedes Dokument mit KI-Funktionen, klassifiziert es in einen Vereinbarungstyp und extrahiert strukturierte Felder für jeden Typ.
Für die Typreferenz siehe FILE Typ.
In diesem Tutorial werden Sie Folgendes lernen:
- Inkrementelle Übernahme von Vertrags-PDFs aus einem Volume als externe
FILEReferenzen mit Auto Loader. - Analysiere jedes Dokument mit
ai_parse_documentFunktion und klassifiziere es mitai_classifyFunktion. - Extrahiere strukturierte Felder für jeden Vereinbarungstyp mit Funktion
ai_extract.
Das Ergebnis ist eine Medaillon-ähnliche Pipeline: Bronze (rohe externe FILE Referenzen), Silber (parsierte und klassifizierte Dokumente) und Gold (extrahierte Felder je Vereinbarungstyp). Weitere Informationen finden Sie unter "Was ist die Medallion Lakehouse-Architektur? Die Bronze-Ebene ist eine Streaming-Tabelle , die Dateien schrittweise einschlägt, und die Silber- und Gold-Schichten sind materialisierte Ansichten , die nur neu berechnet werden, wenn sich ihre Eingaben ändern.
Anforderungen
Um dieses Tutorial abzuschließen, müssen Sie die folgenden Anforderungen erfüllen:
- Sei in einem Azure Databricks-Arbeitsbereich mit aktiviertem Unity Catalog eingeloggt.
- Habe Berechtigungen, um Tabellen in einem Schema zu erstellen und eine Pipeline zu erstellen.
- Nutzen Sie den Vorschaukanal.
Der samples.sec.contracts Datensatz ist standardmäßig in allen Arbeitsbereichen verfügbar, sodass keine zusätzliche Einrichtung erforderlich ist. Da die Dateien bereits in einem Unity-Katalog-Volume gespeichert sind, speichert dieses Tutorial sie als FILE EXTERNAL Referenzen, ohne deren Inhalt zu kopieren. Um die Pipeline an deine eigenen PDFs anzupassen, verweise den Quellpfad auf ein Volumen, das deine Dateien enthält. Für andere Aufnahmeoptionen siehe Ingest files als DATEITYP.
Erstellen Sie die Dateiverarbeitungspipeline
Die Pipeline verarbeitet Dokumente in drei Phasen.
Schritt 1. Bronze: Importiere rohe PDFs als externe DATEI-Referenzen
Nutze Auto Loader, um die Vertrags-PDFs schrittweise vom Volume zu lesen. Das Lesen von Dateien erfasst format => 'file' eine Referenz und Metadaten für jede Datei, ohne deren Bytes zu materialisieren. Die Deklaration der Spalte als FILE EXTERNAL referenziert jede Datei an Ort und Stelle, ohne deren Inhalt zu kopieren.
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/")
)
-
Funktioniert für große Dateien: Ein großes PDF bleibt im Band, während die Tabellenzeile nur eine leichte
FILEReferenz speichert (uri,size,content_type,checksum). Vergleichen Sie dies mit dem TypBINARY, der die Bytes in der Zeile inlinet. -
Inkrementelle Verarbeitung: Die Streaming-Tabelle nimmt neue Dateien schrittweise ein, sobald sie im Quellcode eintreffen, ohne bestehende Dateien neu zu verarbeiten. Der
samples.sec.contractsDatensatz in diesem Beispiel ist statisch, aber mit einer Live-Quelle werden bei jedem Update der Pipeline neue Dateien erfasst. Um auch Quellenänderungen und -löschungen weiterzugeben, sollte der Änderungsfeed mit eingetragen werdenAUTO CDC. Siehe Aktualisierungen und Löschungen mit AUTO CDC anwenden.
Schritt 2. Silber: Dokumente analysieren und klassifizieren
Geben Sie jede FILE Funktion zu, ai_parse_document um das Roh-PDF in ein strukturiertes VARIANT PDF umzuwandeln, das Dokumentelemente, Layout-Metadaten und Text enthält. Da ai_parse_document sie eine FILE Spalte akzeptiert, liest sie das Dokument direkt aus dem Speicher und lädt die Bytes niemals in den Clusterspeicher.
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")
)
Hinweis
Die Definition des Parseschritts als materialisierte Ansicht über der raw_contracts Streaming-Tabelle inkrementellisiert die Berechnung. Jedes Pipeline-Update läuft ai_parse_document nur auf den seit dem letzten Update hinzugefügten Dateien, nicht auf der gesamten Tabelle. Da ai_parse_document dies der teuerste Schritt ist, vermeidet dies das Überarbeiten von Dokumenten, die Sie bereits bearbeitet haben. Das inkrementelle Aktualisieren materialisierter Ansichten erfordert serverlose Rechenleistung; Führe die Pipeline serverlos aus. Siehe Spark Declarative Pipelines.
Als Nächstes wird die geparste Ausgabe an ai_classify die Funktion weitergegeben, um jedem Dokument einen von fünf Vereinbarungstypen zuzuweisen. Dokumente mit Analysefehlern werden vor der Klassifizierung herausgefiltert. Dieses Beispiel pinnt ai_classify an Version 2.1, die die Klassifikation als per-Label-Objekt zurückgibt und das Label vom value Schlüssel ablegt.
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
Um die Klassifikationsgenauigkeit zu verbessern, fügen Sie Etikettenbeschreibungen und eine instructions Option für ai_classifyhinzu. Siehe ai_classify Funktion.
Schritt 3: Gold: Felder pro Vereinbarungstyp extrahieren
Jeder Vereinbarungstyp hat seinen eigenen Satz relevanter Felder. Filtere die klassifizierten Dokumente auf einen Typ, übergebe den geparsten Inhalt, damit ai_extract er mit einem Schema der gewünschten Felder funktioniert, und flache die Antwort dann in typisierte Spalten auf. Dieses Beispiel verbindet ai_extract sich mit Version 2.1, in der jedes extrahierte Feld ein Objekt ist, also liest man seinen value Schlüssel.
Das folgende Beispiel bildet die Goldtabelle für Beratungsverträge:
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")
)
Mit diesen Anweisungen haben Sie eine vollständig inkrementelle Pipeline: Wenn neue Vertrags-PDFs im Volume eintreffen, nimmt Auto Loader sie als externe FILE Referenzen ai_parse_document ein und ai_classify routet jedes Dokument, und die consulting_agreements goldmaterialisierte Ansicht zeigt die extrahierten Felder.
Erkunde auf eigene Faust
Die Pipeline klassifiziert Dokumente in fünf Vereinbarungstypen, extrahiert jedoch nur consulting_agreementFelder für . Um sie zu erweitern, wiederhole man den Goldschritt für jeden verbleibenden Typ, wobei der contract_type Filter und das Schema ai_extract angepasst werden, um die für diesen Typ relevanten Felder zu übereinstimmen. Beispiel:
-
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
Weitere Ressourcen
-
FILETyp - Dateien als FILE-Typ eintragen
- FILE-Funktionen Quickstart
- Erfahren Sie mehr über das automatische Laden. Siehe Was ist Autoloader?.