Útmutató: Fájlfeldolgozó csővezeték építése a FILE típussal

Important

Ez a funkció bétaverzióban érhető el. A munkaterület rendszergazdái az Előnézetek lapon szabályozhatják a funkcióhoz való hozzáférést. Lásd: Az Azure Databricks előzetes verziójának kezelése.

Tanuld meg, hogyan építs egy medál pipeline-t a Lakeflow pipeline-rel, amely a strukturálatlan dokumentumokat végétől végéig feldolgozza. Ez a példa a samples.sec.contracts mintaadatkészletet használja, amely egy SEC által benyújtott jogi megállapodások gyűjteménye, amelyeket PDF-ként tárolnak egy Unity Catalog kötetben.

A csővezeték Auto Loaderrel kezelt FILE hivatkozásként fogadja be a PDF-eket, minden dokumentumot AI funkciókkal elemzi, egyezségi típusba sorolja őket, és minden típushoz strukturált mezőket bont ki.

A típusreferenciáért lásd: FILE típus.

Ebben az oktatóanyagban a következőket meg fogja tanulni:

Az eredmény egy medál-stílusú csővezeték: bronz (nyers kezelt FILE hivatkozások), ezüst (parzált és titkosított dokumentumok), és arany (kinyert mezők megállapodás típuson). További információért lásd: Mi a medallion lakehouse architektúra? A bronzréteg egy streaming tábla , amely fokozatosan veszi fel a fájlokat, míg az ezüst és arany rétegek materializált nézetek , amelyeket csak akkor számolnak újra, ha bemenetük változik.

Requirements

Az oktatóanyag elvégzéséhez meg kell felelnie a következő követelményeknek:

  • Be kell jelentkezni egy Azure Databricks munkaterületre, ahol a Unity Catalog engedélyezve van.
  • Legyen engedélyezve a FILE típus a munkaterületen. A munkaterület adminei engedélyezhetik ezt az Előnézetek oldalán. Lásd: Az Azure Databricks előzetes verziójának kezelése.
  • Legyen jogosultsága tábla létrehozására egy sémában és pipeline létrehozására.
  • Legyen egy Unity Catalog kötet, amibe írhatsz. Ezt a kötetet a bronz tábla FileSpacerészének hirdeted, és a Unity Catalog a bejegyzett fájlokat kezelt tárhelyként másolja be.
  • Használd a Preview csatornát.

Az samples.sec.contracts adathalmaz alapértelmezés szerint minden munkaterületen elérhető. Ez a tutorial a feltöltött PDF-eket FILE MANAGED hivatkozásként tárolja: a Unity Catalog minden fájlt lemásol a tábláként bejelentett FileSpace kötetbe, és kezeli a táblával, így a sorok törlése miatt a hivatkozott fájlok alkalmasak a szemétgyűjtésre, és a táblázat és fájlok szinkronban maradnak. Ha a pipeline-t a saját PDF-jeidhez igazítsd, irányítsd a forrásútvonalat egy olyan kötetre, amely tartalmazza a fájljait. További felvételi opciókért lásd: Ingest fájlok FILE típusként.

Hozd létre a fájlfeldolgozó csővezetéket

A csővezeték három szakaszban dolgozza fel a dokumentumokat.

1. lépés. Bronze: nyers PDF-ek feldolgozása menedzselt FÁJL hivatkozásként

Használd az Auto Loadert, hogy fokozatosan olvasd el a szerződéses PDF-eket a kötetből. A fájlok olvasása format => 'file' minden fájlhoz hivatkozást és metaadatot rögzít anélkül, hogy a bájtjait megvalósítaná. Az oszlop bevallása FILE MANAGED minden fájlt a tábla FileSpace, a tábla által beállított databricks.filespace-preview térfogatba a tábla tulajdonságával beállított volumenbe, így a Unity Catalog kezeli a fájlokat a táblával.

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/")
  )
  • Nagy fájloknál működik: egy nagy PDF a táblákban FileSpacetalálható, míg a táblázat sor csak egy könnyű FILE hivatkozást tárol (uri, size, content_type, checksum). Hasonlítsuk össze a BINARY sort besorolt bájtokat besoroló típussal.
  • Kezelt fájl életciklus: A Unity Catalog minden feltöltött fájlt lemásol a táblákba FileSpace , és a táblával kezeli: a sorok törlése miatt a hivatkozott fájlok alkalmassá teszik a hulladékgyűjtésre, így a tábla és fájlok szinkronban maradnak. Részletekért lásd: FÁJLKEZELŐ és FÁJL KÜLSŐ oldal.
  • Inkrementális feldolgozás: a streaming table fokozatosan vesz fel új fájlokat, amint azok megérkeznek a forrásba, anélkül, hogy újrafeldolgozná a meglévőket. A példa adathalmaza samples.sec.contracts statikus, de élő forrással minden pipeline frissítéskor új fájlokat vesznek fel. A forrásváltozások és törlések terjesztéséhez is lenyeljük be a változási feedet .AUTO CDC Lásd : Frissítések és törlések alkalmazása az AUTO CDC-vel.

2. lépés. Ezüst: dokumentumok elemzése és osztályozása

Mindegyiket továbbítsd FILEai_parse_document a funkcióhoz , hogy a nyers PDF-et strukturált VARIANT dokumentumelemeket, elrendezési metaadatokat és szöveget tartalmazó formává alakítsd át. Mivel ai_parse_document elfogad egy FILE oszlopot, közvetlenül a dokumentumot a tárolóból olvassa, és soha nem tölti be a bájtokat a klasztermemóriába.

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

Jegyzet

Az elemzési lépés definiálása a streaming tábla felett materializált nézetként raw_contracts fokozatosan növeli a számítást. Minden pipeline frissítés csak az előző frissítés óta hozzáadott fájlokon fut, ai_parse_document nem az egész táblán. Mivel ai_parse_document ez a legköltségesebb lépés, elkerüli a már feldolgozott dokumentumok újraparzálását. A materializált nézetek inkrementális frissítése szerver nélküli számítást igényel; Futtatd a pipeline-t szerver nélkül. Lásd: Spark deklaratív adatfeldolgozási folyamatok.

Ezután továbbítsuk az eszrendelt kimenetet aai_classify függvénynek, amely minden dokumentumhoz az öt megállapodás közül az egyik típushoz köti. Az elemzési hibákat tartalmazó dokumentumok szűrése a besorolás előtt történik. Ez a példa a 2.1-es verzióhoz köt ai_classify , amely címkénkénti objektumként adja vissza a besorolást, tehát olvasd el a címkét a value kulcsról.

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

Jótanács

A besorolás pontosságának javítása érdekében adj hozzá címkék leírásokat és egy instructions opciót .ai_classify Lásd: ai_classify függvény.

3. lépés. Arany: kivonási mezők megállapodás típus szerint

Minden megállapodástípusnak megvan a maga releváns mezőkészlete. Szűrd be a titkosított dokumentumokat egy típusra, add át az elemzött tartalmat, hogy ai_extract a kívánt mezők sémájával működjön, majd a választ gépelt oszlopokba simítsd le. Ez a példa a 2.1-es verzióhoz köti ai_extract ki, ahol minden kihúzott mező egy objektum, tehát olvasd el a value kulcsát.

Az alábbi példa építi fel a tanácsadói megállapodások aranytáblázatát:

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

Ezekkel az állításokkal teljesen inkrementális folyamatot kapsz: ahogy új szerződéses PDF-ek érkeznek a kötetbe, az Auto Loader kezelt hivatkozásként ai_parse_document veszi fel őketFILE, ai_classify és minden dokumentumot irányít, és az consulting_agreements arany materializált nézet felszínre kerül a kinyert mezők.

Példajegyzetfüzetek

A következő jegyzetfüzetek tartalmazzák ennek a bemutatónak a teljes folyamatát. Ezek a jegyzetfüzetek csővezeték-kódok, nem futtatható jegyzetfüzetek. Importáld a jegyzetfüzetet a nyelvedhez, majd megadd az útját a forráskód mezőben, amikor beállítod a csővezetéket. Lásd: Folyamatok konfigurálása.

SQL

Fájlfeldolgozási pipeline SQL notebook

Jegyzetfüzet szerezz

Python

Fájlfeldolgozó pipeline Python notebook

Jegyzetfüzet szerezz

Fedezz fel magad

A csővezeték öt megállapodási típusba sorolja a dokumentumokat, de csak consulting_agreementa mezőket kinyeri . A bővítéshez ismételjük meg az arany lépést minden maradék típusnál, módosítva a contract_type szűrőt és a ai_extract sémát, hogy illeszkedjen az adott típushoz kapcsolódó mezőkhöz. Például:

  • 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

További erőforrások