Notatka
Dostęp do tej strony wymaga autoryzacji. Może spróbować zalogować się lub zmienić katalogi.
Dostęp do tej strony wymaga autoryzacji. Możesz spróbować zmienić katalogi.
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:
- Stopniowe pobieranie kontraktowych PDF-ów z woluminu jako
FILEzewnętrznych referencji za pomocą Auto Loadera. - Parsuj każdy dokument z funkcją
ai_parse_documenti klasyfikuj go według funkcjiai_classify. - Wyodrębnij pola strukturalne dla każdego typu umowy z
ai_extractfunkcją.
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).checksumPorównaj to z typemBINARY, 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.contractsdanych 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 CDCZobacz 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
-
FILEtyp - Pobieraj pliki jako typ pliku
- Funkcje szybkiego startu FILE
- Dowiedz się więcej o automatycznym ładowaniu. Zobacz Co to jest moduł automatycznego ładowania?.