Eğitim: FILE tipiyle bir dosya işleme boru hattı oluşturun

Important

Bu özellik Beta sürümündedir. Çalışma alanı yöneticileri Bu özelliğe erişimi Önizlemeler sayfasından denetleyebilir. Bkz. Azure Databricks önizlemelerini yönetme.

Lakeflow boru hattıyla yapılandırılmamış belgeleri uçtan uca işleyen bir madalyon boru hattı nasıl inşa edileceğini öğrenin. Bu örnek, samples.sec.contracts SEC tarafından doldurulmuş yasal anlaşmaların PDF olarak Unity Kataloğu ciltinde saklanan örnek veri setini kullanır.

Boru hattı, PDF'leri Auto Loader ile yönetilen FILE referanslar olarak alır, her belgeyi yapay zeka fonksiyonlarıyla ayrıştırır, bir anlaşma türüne sınıflandırır ve her tür için yapılandırılmış alanlar çıkarır.

Tip referansı için bkz.FILE

Bu kılavuzda şunları yapacaksınız:

Sonuç olarak madalyon tarzı bir boru hattı ortaya çıkar: bronz (ham yönetilen FILE referanslar), gümüş (ayrıştırılmış ve gizli belgeler) ve altın (anlaşma türüne göre çıkarılmış alanlar). Daha fazla bilgi için bkz. Madalyon göl evi mimarisi nedir? Bronz katman, dosyaları kademeli olarak alan bir akış tablosudur ve gümüş ve altın katmanlar yalnızca girdileri değiştiğinde yeniden hesaplanan maddi görünümlerdir .

Requirements

Bu öğreticiyi tamamlamak için aşağıdaki gereksinimleri karşılamanız gerekir:

  • Unity Catalog etkin iken bir Azure Databricks çalışma alanına giriş yapın.
  • Çalışma alanınız için bu FILE tip etkin olsun. Çalışma alanı yöneticileri bunu Önizlemeler sayfasından etkinleştirebilir. Bkz. Azure Databricks önizlemelerini yönetme.
  • Bir şemada tablo oluşturma ve bir pipeline oluşturma iznlerine sahip olmak.
  • Yazabileceğiniz bir Unity Kataloğu cildi olsun. Bu cildi bronz tablo FileSpaceolarak ilan ediyorsunuz ve Unity Catalog alınan dosyaları yönetilen depolama olarak ona kopyalıyor.
  • Önizleme kanalını kullanın.

Veri samples.sec.contracts seti varsayılan olarak tüm çalışma alanlarında mevcuttur. Bu eğitim, alınan PDF'leri referans olarak FILE MANAGED saklar: Unity Catalog, her dosyayı tablo FileSpace olarak ilan ettiğiniz hacme kopyalar ve tablo ile yönetir, böylece satır silmek, referans edilen dosyaları çöp toplamaya uygun hale getirir ve tablo ile dosyaları senkronize kalır. Boru hattını kendi PDF'lerinize uyarlamak için, kaynak yolunu dosyalarınızı içeren bir hacme yönlendirin. Diğer alım seçenekleri için FILE tipi olarak Ingest files sayfasına bakınız.

Dosya işleme boru hattı oluşturun

Boru hattı belgeleri üç aşamada işliyor.

Adım 1. Bronze: ham PDF'leri yönetilen DOSYA referansları olarak alın

Ciltten sözleşme PDF'lerini kademeli olarak okumak için Auto Loader kullanın. Dosyaları okumak format => 'file' , her dosyanın baytlarını oluşturmadan bir referans ve meta veri yakalar. Sütunu olarak FILE MANAGED ilan etmek, her dosyayı tablonun FileSpacehacmine kopyalar, tablo özelliğiyle ayarladığınız databricks.filespace-preview hacm, böylece Unity Kataloğu, dosyaları tablo ile yönetir.

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/")
  )
  • Büyük dosyalar için çalışır: büyük bir PDF tabloda FileSpaceyer alırken, tablo satırı yalnızca hafif FILE bir referans (uri, size, content_type, checksum). Bunu, sıradaki baytların iç çizimini oluşturan type BINARY ile karşılaştırın.
  • Yönetilen dosya yaşam döngüsü: Unity Catalog, alınan her dosyayı tabloya FileSpace kopyalar ve tablo ile yönetir: satır silmek, referans edilen dosyaları çöp toplamaya uygun hale getirir, böylece tablo ve dosyaları senkronize kalır. Detaylar için DOSYA YÖNETİLİLİR ve DOSYA HARİŞİ bkz.
  • Artan işleme: Akış tablosu, mevcut dosyaları yeniden işlemeden kaynağa ulaştıkça yeni dosyaları kademeli olarak alır. Bu örnekteki samples.sec.contracts veri seti statiktir, ancak canlı bir kaynakla her boru hattı güncellemesinde yeni dosyalar alınır. Kaynak değişiklikleri ve silmeleri de yaymak için, değişim beslemesini ile AUTO CDCalın. Bkz. Güncellemeleri ve silmeleri AUTO CDC ile uygula.

Adım 2. Gümüş: belgeleri ayrıştırıp sınıflandır

Her FILEai_parse_document birini fonksiyona aktararak ham PDF'yi belge öğeleri, düzen meta verileri ve metin içeren yapılandırılmış VARIANT bir hale dönüştürün. Bir ai_parse_document sütun kabul ettiği FILE için, belgeyi doğrudan depolamadan okur ve baytları küme belleğine yüklemez.

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

Uyarı

Ayrıştırma adımını, akış raw_contracts tablosu üzerinde maddi bir görünüm olarak tanımlamak, hesaplamayı artırarak artırır. Her pipeline güncellemesi ai_parse_document sadece son güncellemeden sonra eklenen dosyalarda çalışır, tüm tabloda değil. En pahalı adım olduğu ai_parse_document için, daha önce işlediğiniz belgeleri yeniden incelemeden kaçınır. Maddeleştirilmiş görünümlerin kademeli yenilenmesi, sunucusuz hesaplama gerektirir; Pipeline'ı sunucusuz çalıştırın. Bkz. Spark Bildirimsel İşlem Hatları.

Sonra, her belgeye beş anlaşma türünden birini atamak için ayrıştırılmış çıktıyı ai_classify fonksiyona iletin. Ayrıştırma hataları olan belgeler sınıflandırmadan önce filtrelenir. Bu örnek, sınıflandırmayı etiket başına nesne olarak döndüren 2.1 sürümüne sabitler ai_classify , yani etiketi anahtardan value okun.

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

Sınıflandırma doğruluğunu artırmak için, etiket tanımları ve instructions bir seçenek ekleyin ai_classify. bkz. ai_classify işlevi.

Adım 3. Altın: anlaşma türüne göre ekstrakt alanları

Her anlaşma türünün kendine özgü ilgili alanlar seti vardır. Gizli belgeleri tek bir türe filtreleyin, ayrıştırılmış içeriği ai_extract istediğiniz alanların şemasına göre işleyecek şekilde geçirin, ardından cevabı yazılı sütunlara düzleştirin. Bu örnek, her çıkarılan alan bir nesne ai_extract olduğu ve anahtarını okuyacağınız 2.1 sürümüne sabitlervalue.

Aşağıdaki örnek, danışmanlık anlaşmaları için altın tabloyu oluşturur:

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

Bu ifadelerle tamamen kademeli bir boru hattı elde edilir: yeni sözleşme PDF'leri cilme geldiğinde, Auto Loader bunları yönetilen FILE referanslar ai_parse_document olarak alır ve ai_classify her belgeyi yönlendirir, consulting_agreements altın maddeleştirilmiş görünüm çıkarılmış alanları ortaya çıkarır.

Örnek not defterleri

Aşağıdaki defterler, bu eğitimden tüm ürün hattını içermektedir. Bu defterler pipeline kaynak kodudur, çalıştırılabilir not defterleri değildir. Dilinize ait not defterini içe aktarın, ardından boru hattını yapılandırırken Kaynak kodu alanında yolunu belirtin. Bkz . İşlem hatlarını yapılandırma.

SQL

Dosya işleme pipeline SQL notebook

Dizüstü bilgisayar al

Python

Dosya işleme boru hattı Python notebook

Dizüstü bilgisayar al

Kendi başınıza keşfedin

Pipeline belgeleri beş anlaşma türüne sınıflandırır, ancak sadece consulting_agreementiçin alanlar çıkarır. Bunu uzatmak için, kalan her tip için altın adımı tekrarlayın, filtreyi contract_type ve şemayı ai_extract o tipe ait alanlarla eşleştirerek değiştirin. Örneğin:

  • 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

Ek kaynaklar