Tutorial: Membangun pipeline pemrosesan file dengan tipe FILE

Important

Fitur ini ada di Beta. Admin ruang kerja dapat mengontrol akses ke fitur ini dari halaman Pratinjau . Lihat Kelola Pratinjau Azure Databricks.

Pelajari cara membangun pipeline medallion dengan pipeline Lakeflow yang memproses dokumen tidak terstruktur secara end-to-end. Contoh ini menggunakan samples.sec.contracts dataset contoh, kumpulan perjanjian hukum yang diajukan SEC dan disimpan sebagai PDF dalam volume Unity Catalog.

Pipeline mengolah PDF sebagai referensi terkelola FILE dengan Auto Loader, mengurai setiap dokumen dengan fungsi AI, mengklasifikasikannya ke dalam jenis perjanjian, dan mengekstrak field terstruktur untuk setiap tipe.

Untuk referensi tipe, lihat FILE tipe.

Dalam tutorial ini, Anda akan:

Hasilnya adalah pipeline bergaya medali: perunggu (referensi terkelola FILE mentah), perak (dokumen yang telah dianalisis dan diklasifikasikan), dan emas (bidang yang diekstrak berdasarkan jenis perjanjian). Lihat Apa arsitektur medali lakehouse? untuk informasi selengkapnya. Lapisan perunggu adalah tabel streaming yang secara bertahap mengingest file, dan lapisan perak serta emas adalah tampilan materialisasi yang hanya menghitung ulang ketika inputnya berubah.

Persyaratan

Untuk menyelesaikan tutorial ini, Anda harus memenuhi persyaratan berikut:

  • Masuk ke ruang kerja Azure Databricks dengan Unity Catalog diaktifkan.
  • Aktifkan FILE tipe untuk ruang kerja Anda. Admin workspace dapat mengaktifkannya dari halaman Pratinjau . Lihat Kelola Pratinjau Azure Databricks.
  • Memiliki izin untuk membuat tabel dalam skema dan membuat pipeline.
  • Miliki volume Unity Catalog yang bisa Anda tulis. Anda mendeklarasikan volume ini sebagai tabel FileSpaceperunggu , dan Unity Catalog menyalin file yang dimasukkan ke dalamnya sebagai penyimpanan terkelola.
  • Gunakan saluran Pratinjau.

samples.sec.contracts Dataset tersedia secara default di semua workspace. Tutorial ini menyimpan PDF yang dimasukkan sebagai FILE MANAGED referensi: Unity Catalog menyalin setiap file ke volume yang Anda deklarasikan sebagai milik tabel FileSpace dan mengelolanya bersama tabel, sehingga menghapus baris membuat file yang direferensikan memenuhi syarat untuk pengumpulan sampah dan tabel beserta file-filenya tetap sinkron. Untuk menyesuaikan pipeline ke PDF Anda sendiri, arahkan jalur sumber ke volume yang berisi file Anda. Untuk opsi ingestion lainnya, lihat Ingest file sebagai tipe FILE.

Buat pipeline pemrosesan file

Pipeline memproses dokumen dalam tiga tahap.

Langkah 1. Perunggu: menginget PDF mentah sebagai referensi FILE yang dikelola

Gunakan Auto Loader untuk secara bertahap membaca PDF kontrak dari volume. Membaca file dengan format => 'file' menangkap referensi dan metadata untuk setiap file tanpa menghasilkan byte-nya. Mendeklarasikan kolom sebagai FILE MANAGED menyalin setiap file ke FileSpace, volume yang Anda atur dengan databricks.filespace-preview properti tabel, sehingga Unity Catalog mengelola file dengan tabel tersebut.

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/")
  )
  • Berfungsi untuk file besar: PDF besar berada di FileSpacetabel , sementara baris tabel hanya menyimpan referensi ringan FILE (uri, size, content_type, checksum). Bandingkan dengan BINARY tipe yang menginline byte dalam baris.
  • Siklus hidup file yang dikelola: Unity Catalog menyalin setiap file yang dimasukkan ke dalam tabel FileSpace dan mengelolanya bersama tabel: menghapus baris membuat file yang direferensikan memenuhi syarat untuk pengumpulan sampah, sehingga tabel dan file-filenya tetap sinkron. Untuk detailnya, lihat FILE MANAGED dan FILE EXTERNAL.
  • Pemrosesan inkremental: tabel streaming secara bertahap menghisap file baru saat tiba di sumber, tanpa memproses ulang file yang sudah ada. Dataset samples.sec.contracts dalam contoh ini statis, tetapi dengan sumber live, file baru diambil pada setiap pembaruan pipeline. Untuk juga menyebarkan perubahan dan penghapusan sumber, masukkan feed perubahan dengan AUTO CDC. Lihat Terapkan pembaruan dan penghapusan dengan AUTO CDC.

Langkah 2. Perak: mengurai dan mengklasifikasikan dokumen

Teruskan masing-masing FILE ke fungsi untuk ai_parse_document mengonversi PDF mentah menjadi struktur VARIANT yang berisi elemen dokumen, metadata tata letak, dan teks. Karena ai_parse_document menerima kolom FILE , ia membaca dokumen langsung dari penyimpanan dan tidak pernah memuat byte ke memori cluster.

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

Catatan

Mendefinisikan langkah parse sebagai tampilan materialisasi di atas raw_contracts tabel streaming akan meningkatkan perhitungan. Setiap pembaruan pipeline hanya berjalan ai_parse_document pada file yang ditambahkan sejak pembaruan terakhir, bukan pada seluruh tabel. Karena ai_parse_document ini adalah langkah yang paling mahal, ini menghindari perbaikan dokumen yang sudah Anda proses. Penyegaran inkremental tampilan materialisasi membutuhkan komputasi serverless; Jalankan pipeline di serverless. Lihat Spark Declarative Pipelines.

Selanjutnya, teruskan output yang telah diurai ke ai_classify fungsi untuk menetapkan setiap dokumen salah satu dari lima jenis perjanjian. Dokumen dengan kesalahan penguraian difilter sebelum klasifikasi. Contoh ini mempin ai_classify ke versi 2.1, yang mengembalikan klasifikasi sebagai objek per-label, jadi baca label dari value kunci.

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

Untuk meningkatkan akurasi klasifikasi, tambahkan deskripsi label dan instructions opsi ke ai_classify. Silakan lihat fungsi ai_classify.

Langkah 3. Gold: field ekstrak per jenis perjanjian

Setiap jenis perjanjian memiliki set bidang relevan sendiri. Saring dokumen rahasia ke satu tipe, teruskan konten yang telah diurai agar ai_extract berfungsi dengan skema field yang Anda inginkan, lalu ratakan respons ke dalam kolom yang diketik. Contoh ini dipinkan ai_extract ke versi 2.1, di mana setiap field yang diekstrak adalah objek, jadi baca kuncinya value .

Contoh berikut membangun tabel emas untuk perjanjian konsultasi:

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

Dengan pernyataan ini, Anda memiliki pipeline yang sepenuhnya inkremental: saat PDF kontrak baru tiba di volume, Auto Loader mengingesnya sebagai referensi terkelola FILE , ai_parse_document lalu ai_classify mengarahkan setiap dokumen, dan consulting_agreements tampilan materialisasi emas akan menampilkan field yang diekstrak.

Contoh buku catatan

Buku catatan berikut berisi seluruh pipeline dari tutorial ini. Notebook ini adalah kode sumber pipeline, bukan notebook yang dapat dijalankan. Impor notebook untuk bahasa Anda, lalu tentukan jalurnya di kolom Source code saat Anda mengonfigurasi pipeline. Lihat Mengonfigurasi alur.

SQL

Notebook SQL pipeline pemrosesan file

Dapatkan buku catatan

Python

Notebook Python pipeline pemrosesan file

Dapatkan buku catatan

Jelajahi sendiri

Pipeline mengklasifikasikan dokumen ke dalam lima tipe perjanjian tetapi hanya mengekstrak field untuk consulting_agreement. Untuk memperluasnya, ulangi langkah emas untuk setiap tipe yang tersisa, ubah contract_type filter dan ai_extract skema agar sesuai dengan bidang yang relevan dengan tipe tersebut. Contohnya:

  • 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

Sumber daya tambahan