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 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
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 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 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.
- Lass den
FILETyp für deinen Arbeitsbereich aktiviert sein. Workspace-Administratoren können sie über die Vorschauenseite aktivieren. Siehe Manage Azure Databricks Previews. - Habe Berechtigungen, um Tabellen in einem Schema zu erstellen und eine Pipeline zu erstellen.
- Habe ein Unity-Katalog-Volume, in das du schreiben kannst. Du deklarierst diesen Band als Bronze-Tabelle,
FileSpaceund Unity Catalog kopiert die eingenommenen 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 aufgenommenen PDFs als FILE MANAGED Referenzen: Unity Catalog kopiert jede Datei in das von dir deklarierte Volume als Tabellenbuch FileSpace und verwaltet es mit der Tabelle, sodass das Löschen der referenzierten Dateien für die Garbage Collection geeignet ist 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 Ingest files als DATEITYP.
Erstellen Sie die Dateiverarbeitungspipeline
Die Pipeline verarbeitet Dokumente in drei Phasen.
Schritt 1. Bronze: Roh-PDFs als verwaltete FILE-Referenzen eintragen
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. Wenn man die Spalte als FILE MANAGED deklariert, kopiert man jede Datei in die Tabelle FileSpace, 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 in der Tabelle
FileSpace, 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 eingetragene Datei in die Tabelle
FileSpaceund 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 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, speichert Auto Loader sie als verwaltete FILE Referenzen ai_parse_document und ai_classify routet jedes Dokument, und die consulting_agreements goldmaterialisierte Ansicht zeigt die extrahierten Felder.
Beispiel-Notebooks
Die folgenden Notizbücher enthalten die vollständige Pipeline aus diesem Tutorial. Diese Notebooks sind Pipeline-Quellcode, keine ausführenden Notebooks. Importiere das Notizbuch für deine Sprache und gib dann seinen Pfad im Quellcode-Feld an, wenn du die Pipeline konfigurierst. Siehe Konfigurieren von Pipelines.
SQL
SQL-Notizbuch für Dateiverarbeitungspipeline
Python
Dateiverarbeitungs-Pipeline Python Notebook
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?.