Samouczek: Zbuduj potok przetwarzania plików z typem PLIKU

Important

Ta funkcja jest dostępna w wersji beta. Administratorzy obszaru roboczego mogą kontrolować dostęp do tej funkcji ze strony Podglądy . Zobacz Zarządzanie wersjami zapoznawczami usługi Azure Databricks.

Dowiedz się, jak zbudować pipeline medallion z Lakeflow, który przetwarza nieustrukturyzowane dokumenty od początku do końca. Ten przykład wykorzystuje samples.sec.contracts przykładowy zestaw danych, zbiór umów prawnych złożonych przez SEC przechowywanych jako pliki PDF w tomie Unity Catalog.

Potok pobiera pliki PDF jako FILE zewnętrzne referencje za pomocą Auto Loadera, analizuje każdy dokument za pomocą funkcji AI, klasyfikuje go do typu umowy i wyodrębnia pola strukturalne dla każdego typu.

Aby uzyskać informacje o typie, zobacz FILE typ.

Ten samouczek obejmuje następujące kroki:

Efektem jest pipeline w stylu medalionowym: brąz (surowe FILE zewnętrzne odniesienia), srebro (przeanalizowane i tajne dokumenty) oraz złoto (wyodrębnione pola według typu umowy). Aby uzyskać więcej informacji, zobacz Co to jest architektura medallion lakehouse? Warstwa brązu to tabela strumieniowa , która stopniowo pobiera pliki, a warstwy srebrnej i złotej to zmaterializowane widoki , które przetwarzają się dopiero wtedy, gdy ich dane się zmieniają.

Requirements

Aby ukończyć ten samouczek, musisz spełnić następujące wymagania:

  • Zaloguj się do przestrzeni roboczej Azure Databricks z włączonym Unity Catalog.
  • Posiadanie uprawnień do tworzenia tabel w schemacie oraz do tworzenia pipeline.
  • Użyj kanału Podgląd.

Zbiór samples.sec.contracts danych jest domyślnie dostępny we wszystkich przestrzeniach roboczych, więc nie jest wymagana dodatkowa konfiguracja. Ponieważ pliki te już znajdują się w woluminie Unity Catalog, ten samouczek przechowuje je jako FILE EXTERNAL referencje bez kopiowania ich treści. Aby dostosować pipeline do własnych PDF-ów, skieruj ścieżkę źródłową na wolumin zawierający twoje pliki. Aby poznać inne opcje pobierania, zobacz pliki pobierające pliki jako typ PLIKU.

Utwórz potok przetwarzania plików

Pipeline przetwarza dokumenty w trzech etapach.

Krok nr 1. Brąz: pobieranie surowych plików PDF jako zewnętrznych odniesień PLIKÓW

Użyj Auto Loadera, aby stopniowo czytać PDF-y kontraktów z tomu wolumu. Odczyt plików z plikami format => 'file' rejestruje referencję i metadane dla każdego pliku bez materializacji jego bajtów. Deklarowanie kolumny jako odnosi FILE EXTERNAL się do każdego pliku w miejscu, bez kopiowania jego treści.

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/")
  )
  • Działa dla dużych plików: duży PDF pozostaje w woluminie, podczas gdy wiersz tabeli przechowuje tylko lekkie FILE źródło (uri, size, , content_type). checksum Porównaj to z typem BINARY , który wplata bajty w wierszu.
  • Przetwarzanie przyrostowe: tabela strumieniowa stopniowo pobiera nowe pliki, gdy dotrą do źródła, bez ponownego przetwarzania istniejących. Zbiór samples.sec.contracts danych w tym przykładzie jest statyczny, ale przy live source nowe pliki są pobierane przy każdej aktualizacji pipeline. Aby również propagować zmiany i usuwanie źródeł, należy pobierać feed zmian z .AUTO CDC Zobacz Wprowadź aktualizacje i usunięcia z AUTO CDC.

Krok nr 2. Srebro: analizować i klasyfikować dokumenty

Przekaż każdy do FILEai_parse_document funkcji , aby przekonwertować surowy PDF na ustrukturyzowany VARIANT dokument zawierający elementy dokumentu, metadane układu i tekst. Ponieważ ai_parse_document akceptuje kolumnę FILE , odczytuje dokument bezpośrednio z pamięci i nigdy nie ładuje bajtów do pamięci klastrowej.

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")
  )

Uwaga / Notatka

Zdefiniowanie kroku parsowania jako zmaterializowanego widoku nad tabelą raw_contracts strumieniową inkrementalizuje obliczenia. Każda aktualizacja potoku uruchamia ai_parse_document się tylko na plikach dodanych od ostatniej aktualizacji, a nie na całej tabeli. Ponieważ ai_parse_document jest to najdroższy krok, unika to naprawiania już przetworzonych dokumentów. Stopniowe odświeżanie zmaterializowanych widoków wymaga obliczeń bez serwera; Uruchom pipeline na serwerless. Zobacz Potoki deklaratywne platformy Spark.

Następnie przekaż rozłożone wyjście ai_classify do funkcji , aby przypisać każdemu dokumentu jeden z pięciu typów umowy. Dokumenty z błędami analizy są filtrowane przed klasyfikacją. Ten przykład przypina ai_classify do wersji 2.1, która zwraca klasyfikację jako obiekt dla każdej etykiety, więc odczytuje etykietę z klucza 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""")
  )

Wskazówka

Aby poprawić dokładność klasyfikacji, dodaj opisy etykiet oraz instructions opcję do ai_classify. Zobacz ai_classify funkcję.

Krok nr 3. Złoto: pola ekstrakcyjne według typu umowy

Każdy typ umowy ma własny zestaw odpowiednich pól. Przefiltruj dokumenty tajne do jednego typu, przekaż analizowaną treść ai_extract do pracy ze schematem wybranych pól, a następnie spłaszcz odpowiedź do kolumn typowych. Ten przykład przypina ai_extract do wersji 2.1, w której każde wyodrębnione pole jest obiektem, więc odczytaj jego value klucz.

Poniższy przykład buduje "złoty stół dla umów konsultingowych":

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")
  )

Dzięki tym stacjom masz w pełni inkrementalny pipeline: gdy nowe PDF-y kontraktów pojawiają się w woluminie, AutoLoader pobiera je jako FILE zewnętrzne referencje ai_parse_document i kieruje ai_classify każdy dokument, a consulting_agreements widok gold materialized wyświetla wyodrębnione pola.

Eksploruj na własną rękę

Potok klasyfikuje dokumenty na pięć typów umów, ale wyodrębnia pola tylko consulting_agreementdla . Aby ją rozszerzyć, powtórz krok złoty dla każdego pozostałego typu, zmieniając contract_type filtr i schemat tak ai_extract , aby dopasowały pola istotne dla danego typu. Przykład:

  • 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

Dodatkowe zasoby