Tutorial: Erstellen Sie eine Dateiverarbeitungspipeline mit dem FILE-Typ

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:

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 FILE Referenz speichert (uri, size, content_type, checksum). Vergleichen Sie dies mit dem Typ BINARY , 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.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 eingetragen werden AUTO 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