Tutorial: Bouw een bestandsverwerkingspijplijn met het type FILE

Important

Deze functie bevindt zich in de bètaversie. Werkruimtebeheerders kunnen de toegang tot deze functie beheren vanaf de pagina Previews . Zie Azure Databricks previews beheren.

Leer hoe je een medallion-pijplijn bouwt met Lakeflow-pijplijn die ongestructureerde documenten van begin tot eind verwerkt. Dit voorbeeld gebruikt de samples.sec.contracts voorbeelddataset, een verzameling door SEC ingediende juridische overeenkomsten die als PDF's in een Unity Catalog-volume zijn opgeslagen.

De pipeline neemt de PDF's als externe FILE referenties op met Auto Loader, parseert elk document met AI-functies, classificeert het in een overeenkomsttype en extraht gestructureerde velden voor elk type.

Voor de typereferentie, zie FILE type.

In deze tutorial, zul je:

Het resultaat is een medaillonachtige pijpleiding: brons (ruwe externe FILE referenties), zilver (geanalyseerde en geclassificeerde documenten) en goud (extraherde velden per overeenkomsttype). Zie Wat is de medallion lakehouse-architectuur? voor meer informatie. De bronze-laag is een streaming-tabel die bestanden incrementeel invoert, en de silver- en gold-lagen zijn gematerialiseerde weergaven die pas opnieuw berekenen wanneer hun invoer verandert.

Requirements

Als u deze zelfstudie wilt voltooien, moet u aan de volgende vereisten voldoen:

  • Wees ingelogd op een Azure Databricks-werkruimte met Unity Catalog ingeschakeld.
  • Heb permissies om tabellen in een schema te maken en een pipeline te maken.
  • Gebruik het Preview-kanaal.

De samples.sec.contracts dataset is standaard beschikbaar in alle werkruimtes, dus er is geen extra setup nodig. Omdat de bestanden al in een Unity Catalog-volume staan, slaat deze tutorial ze op als FILE EXTERNAL referenties zonder hun inhoud te kopiëren. Om de pipeline aan te passen aan je eigen PDF's, wijs je het bronpad naar een volume dat je bestanden bevat. Voor andere invoeropties, zie Ingest files als het FILE-type.

Maak de bestandsverwerkingspijplijn aan

De pijplijn verwerkt documenten in drie fasen.

Stap 1. Brons: voer ruwe PDF's in als externe BESTANDSREFERENTIES

Gebruik Auto Loader om de contract-PDF's stapsgewijs van het volume te lezen. Het lezen van bestanden met format => 'file' vangt een referentie en metadata voor elk bestand zonder de bytes te materialiseren. Het verklaren van de kolom verwijst naar FILE EXTERNAL elk bestand ter plaatse, zonder de inhoud te kopiëren.

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/")
  )
  • Werkt voor grote bestanden: een grote PDF blijft in het volume, terwijl de tabelrij alleen een lichtgewicht FILE referentie opslaat (uri, size, content_type, checksum). Vergelijk dit met het BINARY type, dat de bytes in de rij inlijnt.
  • Incrementele verwerking: de streamingtabel voert incrementeel nieuwe bestanden in zodra ze in de bron aankomen, zonder bestaande bestanden opnieuw te verwerken. De samples.sec.contracts dataset in dit voorbeeld is statisch, maar met een live bron worden bij elke pipeline-update nieuwe bestanden opgepikt. Om ook bronwijzigingen en -verwijderingen te verspreiden, consumeer je de change feed met AUTO CDC. Zie Updates en verwijderingen toepassen bij AUTO CDC.

Stap 2. Zilver: documenten analyseren en classificeren

Geef elke FILEai_parse_document functie door om de ruwe PDF om te zetten in een structuur VARIANT met documentelementen, lay-outmetadata en tekst. Omdat ai_parse_document het een FILE kolom accepteert, leest het het document direct uit het geheugen en laadt het de bytes nooit in het clustergeheugen.

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

Het definiëren van de parsestap als een gematerialiseerde weergave over de raw_contracts streamingtabel verhoogt de berekening. Elke pipeline-update draait ai_parse_document alleen op de bestanden die sinds de laatste update zijn toegevoegd, niet op de hele tabel. Omdat ai_parse_document dit de duurste stap is, voorkom je het opnieuw parpareren van documenten die je al hebt verwerkt. Incrementele verversing van gematerialiseerde weergaven vereist serverless compute; Draai de pipeline serverloos. Zie Spark Declarative Pipelines.

Vervolgens geef je de geparseerde output door aan ai_classify de functie om elk document een van de vijf overeenkomsttypes toe te wijzen. Documenten met parseringsfouten worden gefilterd vóór classificatie. Dit voorbeeld pint ai_classify naar versie 2.1, die de classificatie als een per-label object teruggeeft, dus het label uit de value sleutel leest.

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

Om de classificatienauwkeurigheid te verbeteren, voeg labelbeschrijvingen toe en een instructions optie aan ai_classify. Zie ai_classify de functie.

Stap 3. Goud: extractievelden per overeenkomsttype

Elk overeenkomsttype heeft zijn eigen set relevante velden. Filter de geclassificeerde documenten op één type, geef de geanalyseerde inhoud door om te ai_extract functioneren met een schema van de velden die je wilt, en plat het antwoord vervolgens in getypte kolommen. Dit voorbeeld sluit ai_extract aan op versie 2.1, waarin elk geëxtraheerd veld een object is, dus de value sleutel wordt gelezen.

Het volgende voorbeeld vormt de gouden tabel voor adviesovereenkomsten:

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

Met deze statements heb je een volledig incrementele pipeline: zodra nieuwe contract-PDF's in het volume verschijnen, voert Auto Loader ze in als externe FILE referenties ai_parse_document en ai_classify routeert elk document, waarna de consulting_agreements goudgematerialiseerde weergave de geëxtraheerde velden naar boven brengt.

Ontdek het zelf

De pijplijn classificeert documenten in vijf overeenkomsttypes, maar extraheert velden voor slechts consulting_agreement. Om deze uit te breiden, herhaal je de gouden stap voor elk overgebleven type, waarbij je het contract_type filter en het ai_extract schema aanpast zodat ze overeenkomen met de velden die relevant zijn voor dat type. Voorbeeld:

  • affiliate_agreement party_1_name: , party_2_name, , commission_rate,payment_frequency
  • marketing_agreement party_1_name: , party_2_name, , effective_date,territory
  • hosting_agreement provider_name: , customer_name, , effective_date,term_length
  • escrow_agreement owner_name: , licensee_name, , escrow_agent_name,software_name

Aanvullende bronnen