Tutorial: Erstellen Sie eine Dateiverarbeitungspipeline mit dem FILE-Typ

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:

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 FileSpace der 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 FileSpace der Tabelle, während die Tabellenzeile nur eine leichte FILE Referenz speichert (uri, size, content_type, checksum). Vergleichen Sie dies mit dem Typ BINARY , der die Bytes in der Zeile inlinet.
  • Verwalteter Dateilebenszyklus: Unity Catalog kopiert jede ingestierte Datei in das FileSpace der 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.contracts Datensatz 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 AUTO CDC eingelesen 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

Notebook abrufen

Python

Dateiverarbeitungs-Pipeline Python-notebook

Notebook abrufen

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