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. Um es zu verwenden, muss ein Workspace-Administrator den Dateityp auf der Vorschauen-Seite aktivieren. 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 verwaltete 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:
- Importiere schrittweise Vertrags-PDFs aus einem Volume als verwaltete
FILEReferenzen mit Auto Loader. - Analysiere jedes Dokument mit der
ai_parse_documentFunktion und klassifiziere es mit derai_classifyFunktion. - Extrahiere strukturierte Felder für jeden Vereinbarungstyp mit
ai_extractFunktion.
Das Ergebnis ist eine Pipeline im Medaillon-Stil: Bronze (rohe verwaltete FILE Referenzen), Silber (analysierte und klassifizierte Dokumente) und Gold (extrahierte Felder pro Vereinbarungstyp). Weitere Informationen finden Sie unter "Was ist die Medallion Lakehouse-Architektur? Die Bronze-Ebene ist eine Streaming-Tabelle, die Dateien inkrementell einliest, 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.
- Aktiviere den
FILE-Typ für deinen Arbeitsbereich. Workspace-Administratoren können sie auf der Seite Vorschauen aktivieren. Siehe Manage Azure Databricks Previews. - Sie benötigen Berechtigungen, um Tabellen in einem Schema zu erstellen und eine Pipeline zu erstellen.
- Sie benötigen ein Unity Catalog-Volume, in das Sie schreiben können. Du deklarierst diesen Band als
FileSpaceder Bronze-Tabelle, und Unity Catalog kopiert die ingestierten Dateien als verwalteten Speicher darin. - Nutzen Sie den Vorschaukanal.
Der samples.sec.contracts Datensatz ist standardmäßig in allen Arbeitsbereichen verfügbar. Dieses Tutorial speichert die eingelesenen PDFs als FILE MANAGED Referenzen: Unity Catalog kopiert jede Datei in das von dir deklarierte Volume als FileSpace der Tabelle und verwaltet es mit der Tabelle, sodass das Löschen von Zeilen die referenzierten Dateien für die Garbage Collection infrage kommen lässt und die Tabelle und ihre Dateien synchron bleiben. Um die Pipeline an deine eigenen PDFs anzupassen, verweise den Quellpfad auf ein Volumen, das deine Dateien enthält. Für andere Aufnahmeoptionen siehe Dateien als DATEITYP aufnehmen.
Erstellen Sie die Dateiverarbeitungspipeline
Die Pipeline verarbeitet Dokumente in drei Phasen.
Schritt 1. Bronze: Roh-PDFs als verwaltete FILE-Referenzen importieren
Nutze Auto Loader, um die Vertrags-PDFs inkrementell vom Volume zu lesen. Das Lesen von Dateien mit format => 'file' erfasst eine Referenz und Metadaten für jede Datei, ohne deren Bytes zu materialisieren. Wenn man die Spalte als FILE MANAGED deklariert, kopiert man jede Datei in das FileSpace der Tabelle, das Volumen, das man mit der databricks.filespace-preview Tabelleneigenschaft setzt, sodass Unity Catalog die Dateien mit der Tabelle verwaltet.
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/")
)
-
Funktioniert für große Dateien: Ein großes PDF befindet sich im
FileSpaceder Tabelle, während die Tabellenzeile nur eine leichteFILEReferenz speichert (uri,size,content_type,checksum). Vergleichen Sie dies mit dem TypBINARY, der die Bytes in der Zeile inlinet. -
Verwalteter Dateilebenszyklus: Unity Catalog kopiert jede ingestierte Datei in das
FileSpaceder Tabelle und verwaltet sie mit der Tabelle: Das Löschen von Zeilen macht die referenzierten Dateien für die Garbage Collection geeignet, sodass die Tabelle und ihre Dateien synchron bleiben. Details finden Sie unter FILE MANAGED und FILE EXTERNAL. -
Inkrementelle Verarbeitung: Die Streaming-Tabelle nimmt neue Dateien schrittweise auf, sobald sie in der Quelle 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 mitAUTO CDCeingelesen werden. Siehe Aktualisierungen und Löschungen mit AUTO CDC anwenden.
Schritt 2. Silber: Dokumente parsen und klassifizieren
Übergeben Sie jede FILE an die ai_parse_document Funktion, um das Roh-PDF in ein strukturiertes VARIANT PDF umzuwandeln, das Dokumentelemente, Layout-Metadaten und Text enthält. Da ai_parse_document eine FILE Spalte akzeptiert, liest ai_parse_document 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 führt zu einer inkrementellen Berechnung. Jedes Pipeline-Update wird ai_parse_document nur auf die seit dem letzten Update hinzugefügten Dateien angewendet, nicht auf die gesamte Tabelle. Da ai_parse_document dies der teuerste Schritt ist, vermeidet dies das erneute Parsen 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 die Funktion ai_classify 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, also lesen Sie das Label aus dem Schlüssel 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
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, übergib den geparsten Inhalt an die ai_extract Funktion mit einem Schema der gewünschten Felder und flache die Antwort dann in typisierte Spalten auf. Dieses Beispiel legt ai_extract auf Version 2.1 fest, 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, liest Auto Loader sie als verwaltete FILE Referenzen ein, ai_parse_document und ai_classify leiten jedes Dokument weiter, und die consulting_agreements goldmaterialisierte Ansicht zeigt die extrahierten Felder an.
Beispiel-Notebooks
Die folgenden Notizbücher enthalten die vollständige Pipeline aus diesem Tutorial. Diese Notebooks sind Pipeline-Quellcode, keine ausführbaren Notebooks. Importiere das Notebook für deine Sprache und gib dann seinen Pfad im Quellcode-Feld an, wenn du die Pipeline konfigurierst. Siehe Konfigurieren von Pipelines.
SQL
SQL-Notebook für Dateiverarbeitungspipeline
Python
Dateiverarbeitungs-Pipeline Python-notebook
Erkunde auf eigene Faust
Die Pipeline klassifiziert Dokumente in fünf Vereinbarungstypen, extrahiert jedoch nur für consulting_agreement Felder. Um sie zu erweitern, wiederholen Sie den Goldschritt für jeden verbleibenden Typ, wobei der contract_type Filter und das Schema ai_extract angepasst werden, um ihn an die für diesen Typ relevanten Felder anzupassen. 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 importieren
- FILE-Funktionen Schnellstart
- Erfahren Sie mehr über das automatische Laden. Siehe Was ist Autoloader?.