Tutoriál: Vytvořte pipeline pro zpracování souborů s typem FILE

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 externí 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:

Výsledkem je medailonový proces: bronz (surové externí FILE reference), stříbro (analyzované a utajované dokumenty) a zlato (extrahované pole podle typu dohody). 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ít oprávnění vytvářet tabulky ve schématu a vytvářet pipeline.
  • Použijte kanál Náhled.

Datová samples.sec.contracts sada je ve výchozím nastavení dostupná ve všech pracovních prostorech, takže není potřeba žádné další nastavení. Protože soubory již jsou uloženy v svazku Unity Catalog, tento tutoriál je ukládá jako reference bez FILE EXTERNAL kopírování jejich obsahu. 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. Bronz: ingest raw PDF jako externí 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 odkazuje FILE EXTERNAL na každý soubor na místě, aniž by se kopíroval jeho obsah.

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/")
  )
  • Funguje pro velké soubory: velké PDF zůstává ve svazku, zatímco tabulka ukládá pouze lehkou referenci FILE (uri, size, , content_type). checksum Porovnejte to s BINARY typem, který zařazuje bajty v řádku.
  • 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.contracts sada 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ýpisy máte plně inkrementální pipeline: jakmile do svazku přicházejí nové PDF kontrakty, Auto Loader je nabírá jako externí FILE odkazy, ai_parse_document směruje ai_classify každý dokument a pohled consulting_agreements gold materialized zobrazuje extrahovaná pole.

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