Poznámka:
Přístup k této stránce vyžaduje autorizaci. Můžete se zkusit přihlásit nebo změnit adresáře.
Přístup k této stránce vyžaduje autorizaci. Můžete zkusit změnit adresáře.
Important
Tato funkce je v beta verzi. Správci pracovního prostoru můžou řídit přístup k této funkci ze stránky Previews . Viz Manage Azure Databricks preview.
Naučte se, jak vytvořit medailonový pipeline s Lakeflow pipeline, který zpracovává nestrukturované dokumenty od začátku do konce. Tento příklad využívá ukázkovou datovou sadu samples.sec.contracts , sbírku právních dohod podaných SEC uložených jako PDF ve svazku Unity Catalog.
Pipeline zpracovává PDF jako spravované FILE reference pomocí Auto Loaderu, analyzuje každý dokument pomocí AI funkcí, klasifikuje jej do typu dohody a extrahuje strukturovaná pole pro každý typ.
Pro typovou referenci viz FILE typ.
V tomto kurzu:
- Postupně spotřebovávat PDF kontraktů z jednoho svazku jako spravované
FILEreference pomocí Auto Loaderu. - Parsujte každý dokument pomocí
ai_parse_documentfunkce a klasifikujte jej pomocíai_classifyfunkce. - Extrahujte strukturovaná pole pro každý typ dohody s
ai_extractfunkcí.
Výsledkem je medailonový pipeline: bronz (surové spravované FILE reference), stříbro (analyzované a utajované dokumenty) a zlato (extrahované pole podle typu smlouvy). Další informace najdete v tématu Co je architektura jezera medallion? Bronzová vrstva je streamovací tabulka , která postupně spotřebovává soubory, a stříbrná a zlatá vrstva jsou materializované pohledy , které se znovu přepočítají až při změně vstupů.
Požadavky
K dokončení tohoto kurzu musíte splnit následující požadavky:
- Buďte přihlášeni do pracovního prostoru Azure Databricks s povoleným Unity Catalog.
- Mějte typ zapnutý
FILEpro váš pracovní prostor. Správci pracovního prostoru to mohou povolit ze stránky Náhledy . Viz Manage Azure Databricks preview. - Mít oprávnění vytvářet tabulky ve schématu a vytvářet pipeline.
- Mějte svazek Unity Catalog, do kterého můžete psát. Tento svazek označíte za bronzovou tabulku
FileSpacea Unity Catalog do něj zkopíruje ingestované soubory jako spravované úložiště. - Použijte kanál Náhled.
Datová samples.sec.contracts sada je ve výchozím nastavení dostupná ve všech pracovních prostorech. Tento tutoriál ukládá ingestované PDF jako FILE MANAGED reference: Unity Catalog zkopíruje každý soubor do svazku, který deklarujete jako tabulku FileSpace , a spravuje ho spolu s tabulkou, takže smazání řádků činí odkazované soubory způsobilými pro garbage collection a tabulka i její soubory zůstávají synchronizované. Pro přizpůsobení pipeline vlastním PDF nasměrujte zdrojovou cestu na svazek, který obsahuje vaše soubory. Pro další možnosti ingestu viz Ingest soubory jako typ SOUBORU.
Vytvořte pipeline pro zpracování souborů
Pipeline zpracovává dokumenty ve třech fázích.
Krok 1. Bronze: ingest raw PDF jako spravované FILE reference
Použijte Auto Loader pro postupné čtení PDF smluv ze svazku. Čtení souborů s format => 'file' zachycuje referenci a metadata pro každý soubor, aniž by materializovalo jeho bajty. Deklarace sloupce jako FILE MANAGED se zkopíruje každý soubor do svazku FileSpacetabulky, který nastavíte vlastností databricks.filespace-preview tabulky, takže Unity Catalog spravuje soubory spolu s tabulkou.
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/")
)
-
Funguje to pro velké soubory: velké PDF je umístěno v tabulce
FileSpace, zatímco řádek tabulky uchovává pouze lehkouFILEreferenci (uri,size, ,content_type).checksumPorovnejte to sBINARYtypem, který zařazuje bajty v řádku. -
Životní cyklus spravovaných souborů: Unity Catalog zkopíruje každý ingestovaný soubor do tabulky a spravuje ho spolu s tabulkou: smazání řádků činí odkazované soubory způsobilými
FileSpacepro sběr odpadu, takže tabulka a její soubory zůstávají synchronizované. Podrobnosti viz FILE MANAGED a FILE EXTERNAL. -
Inkrementální zpracování: tabulka streamování postupně přijímá nové soubory, jakmile dorazí do zdroje, aniž by znovu zpracovávala ty stávající. Datová
samples.sec.contractssada v tomto příkladu je statická, ale s živým zdrojem se při každé aktualizaci pipeline objevují nové soubory. Pro propagaci změn a mazání zdrojů také ingestujte feed změn pomocíAUTO CDC. Viz Aplikace aktualizací a vymazání u AUTO CDC.
Krok 2. Stříbro: parsování a klasifikace dokumentů
Předejte každé FILEai_parse_document funkci pro převod surového PDF do strukturovaného VARIANT dokumentu obsahujícího prvky dokumentu, metadata rozvržení a text. Protože ai_parse_document přijímá sloupec FILE , čte dokument přímo z úložiště a nikdy nenačítá bajty do clusterové paměti.
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")
)
Note
Definování kroku parsování jako materializovaného pohledu přes streamovací tabulku raw_contracts inkrementalizuje výpočty. Každá aktualizace pipeline běží ai_parse_document pouze na souborech přidaných od poslední aktualizace, ne na celé tabulce. Protože ai_parse_document je to nejdražší krok, vyhnete se opravě dokumentů, které jste již zpracovali. Postupné obnovování materializovaných pohledů vyžaduje serverless výpočetní výkon; Spouštějte pipeline na serverless. Viz deklarativní kanály Sparku.
Dále předejte parsovaný výstup ai_classify funkci, která každému dokumentu přiřadí jeden z pěti typů dohod. Dokumenty s chybami analýzy se před klasifikací odfiltrují. Tento příklad je připojen ai_classify k verzi 2.1, která vrací klasifikaci jako objekt pro každý štítek, takže štítek čtěte z klíče 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
Pro zlepšení přesnosti klasifikace přidejte popisy štítků a instructions možnost .ai_classify Viz ai_classify funkce.
Krok 3. Zlato: extrahujte pole podle typu dohody
Každý typ dohody má svou vlastní sadu relevantních polí. Filtrujte klasifikované dokumenty na jeden typ, předejte parsovaný obsah ai_extract funkci se schématem požadovaných polí, a pak odpověď zploťte do typovaných sloupců. Tento příklad je připojen ai_extract k verzi 2.1, kde je každé extrahované pole objektem, takže čtěte jeho value klíč.
Následující příklad vytváří zlatý stůl pro konzultační smlouvy:
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")
)
S těmito výkazy máte plně inkrementální pipeline: jakmile do svazku přicházejí nové PDF kontrakty, Auto Loader je zpracovává jako spravované FILE reference, ai_parse_document směruje ai_classify každý dokument a zlatý consulting_agreements materializovaný pohled zobrazí extrahovaná pole.
Ukázkové notebooky
Následující zápisníky obsahují kompletní pipeline z tohoto tutoriálu. Tyto zápisníky jsou zdrojový kód pipeline, ne běžitelné notebooky. Importujte zápisník pro svůj jazyk a pak zadejte jeho cestu v poli Zdrojový kód při konfiguraci pipeline. Viz Konfigurace kanálů.
SQL
SQL notebook pro zpracování souborů
Python
Notebook pro zpracování souborů v Python pipeline
Prozkoumávej to sám
Pipeline klasifikuje dokumenty do pěti typů dohod, ale extrahuje pole pouze consulting_agreementpro . Pro rozšíření opakujte zlatý krok pro každý zbývající typ, měňte contract_type filtr a ai_extract schéma tak, aby odpovídaly polím relevantním pro daný typ. Příklady:
-
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
Dodatečné zdroje
-
FILEtyp - Ingest souborů jako typ SOUBORU
- Rychlý start funkcí FILE
- Přečtěte si další informace o Auto Loaderu. Podívejte se na Co je to Auto Loader?