Important
這項功能位於 測試版 (Beta) 中。 工作區管理員可以從 「預覽 」頁面控制對此功能的存取。 請參閱 管理 Azure Databricks 預覽。
學習如何使用 Lakeflow 管線建立一個 medallion 管線,從端到端處理非結構化文件。 此範例使用 samples.sec.contracts 範例資料集,即一組由美國證券交易委員會(SEC)提交的法律協議,以PDF格式儲存在Unity Catalog卷中。
管線會用 AutoLoader 將 PDF 作為外部 FILE 參考資料擷取,使用 AI 函式解析每份文件,將其分類為協議類型,並為每種類型擷取結構化欄位。
關於類型參考,請參見 FILE 類型。
在本教學課程中,您將:
- 用 AutoLoader 逐步從某卷中擷取合約 PDF,作為外部
FILE參考。 - 將每份文件解析為功能
ai_parse_document,並依功能分類ai_classify。 - 為每個與函數的
ai_extract協議類型擷取結構化欄位。
結果是一條類似獎章的管線:青銅(原始外部 FILE 參考)、銀(解析與機密文件)和金(依協議類型提取欄位)。 如需詳細資訊,請參閱獎章湖屋架構的概念。 青銅層是一個 串流資料表 ,會逐步擷取檔案,而銀層和金層則是 實體化的視圖 ,只有當輸入改變時才會重新計算。
要求
您必須滿足下列需求,才能完成本教學課程:
- 登入啟用Unity Catalog的Azure Databricks工作區。
- 擁有在結構中建立資料表和建立管線的權限。
- 使用預覽頻道。
samples.sec.contracts該資料集預設在所有工作區都能取得,因此不需要額外的設定。 因為檔案已經存在於 Unity 目錄卷中,這個教學會將它們儲存為 FILE EXTERNAL 參考資料,而不會複製其內容。 要將管線調整到你自己的 PDF,請將來源路徑指向包含你檔案的卷。 關於其他擷取選項,請參見「 檔案擷取」作為檔案類型。
建立檔案處理流程
流程文件處理分為三個階段。
步驟 1. 青銅:將原始 PDF 作為外部檔案參考
使用自動載入器逐步閱讀該卷的合約 PDF。 讀取檔案 format => 'file' 時,會擷取每個檔案的參考資料和元資料,而不會實體化其位元組。 將欄位宣告為 FILE EXTERNAL 參考,讓每個檔案在位,但不會複製其內容。
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/")
)
-
適用於大型檔案:大型 PDF 會留在磁碟區,而表格列只儲存一個輕量級
FILE參考(uri,size,content_type,checksum)。 與 型BINARY別比較,該型別將列中的位元組內嵌。 -
增量處理:串流資料表在新檔案抵達原始碼時逐步擷取,不重新處理現有檔案。
samples.sec.contracts此範例中的資料集是靜態的,但若有即時來源,每次管線更新都會有新檔案被擷取。 若要同時傳播來源變更與刪除,請用AUTO CDCS 來接收變更訂閱源。 請參見 「與 AUTO CDC 的更新與刪除」。
步驟 2。 銀色:解析與分類文件
將每個 FILE 函式傳遞到 ai_parse_document 函式 ,將原始 PDF 轉換成包含文件元素、版面中繼資料和文字的結構化 VARIANT 文件。 由於 ai_parse_document 接受欄位 FILE ,它會直接從儲存讀取文件,且不會將位元組載入叢集記憶體。
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
將解析步驟定義為對 raw_contracts 串流表的具體視圖,可讓計算過程更為精進。 每次管線更新只會執行 ai_parse_document 自上次更新以來新增的檔案,而非整個資料表。 因為 ai_parse_document 這是最昂貴的步驟,這樣可以避免重新解析你已經處理過的文件。 實體化檢視的增量刷新需要無伺服器運算;在無伺服器模式下運行管線。 參見 Spark 宣告式管線。
接著,將解析後的輸出傳給 ai_classify 函式 ,讓每份文件被指派為五種協議類型之一。 有解析錯誤的文件會在分類前被過濾掉。 此範例釘 ai_classify 在版本 2.1,該版本以每個標籤物件回傳分類,因此從鍵中讀取標籤 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""")
)
提示
為提升分類準確度,請新增標籤描述及instructions選項。ai_classify 請參閱 ai_classify 函式。
步驟 3。 黃金:依協議類型提取欄位
每種協議類型都有其相關的欄位集合。 將機密文件過濾成一種類型,將解析後的內容傳入ai_extract你想要的欄位結構,然後將回應平整成有型的欄位。 此範例是 ai_extract 針對版本 2.1,每個擷取欄位都是一個物件,因此讀取其 value 鍵。
以下範例構建顧問協議的黃金表:
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")
)
有了這些語句,你就擁有一個完全增量的流程:當新的合約 PDF 進入卷中時,Auto Loader 會將它們作為外部 FILE 參考資料匯入, ai_parse_document 並 ai_classify 路由每份文件,而 consulting_agreements 金色實體化檢視則顯示已擷取的欄位。
自己探索吧
該管線將文件分類為五種協議類型,但僅 consulting_agreement提取欄位。 要擴充,對剩餘的每個類型重複金色步驟,並調整 contract_type 過濾器和 ai_extract 結構以匹配該類型相關的欄位。 例如:
-
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
其他資源
-
FILE類型 - 以檔案類型來導入檔案
- FILE 函式快速啟動
- 了解更多關於自動裝填機的資訊。 請參閱 什麼是自動載入器?。