教程:构建带有 FILE 类型的文件处理流水线

Important

此功能在 Beta 版中。 工作区管理员可以从 预览 页控制对此功能的访问。 请参阅 Manage Azure Databricks 预览版

了解如何用Lakeflow管道构建一个Medallion管道,能够端到端处理非结构化文档。 本示例使用 samples.sec.contracts 示例数据集,即以PDF形式存储在Unity目录卷中的SEC提交法律协议集合。

流水线通过自动加载器将PDF作为外部 FILE 引用导入,使用AI函数解析每个文档,将其分类为协议类型,并为每种类型提取结构化字段。

关于类型参考,请参见 FILE 类型

在本教程中,你将:

  • 通过 Auto Loader 逐步从卷中导入合同 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参考(urisizecontent_typechecksum, )。 与类型 BINARY 对比,后者将行中的字节内嵌。
  • 增量处理:流式表在新文件到达源时逐步导入,不重新处理已有文件。 samples.sec.contracts本例中的数据集是静态的,但对于实时源,每次管道更新都会获取新文件。 为了传播源变更和删除,请用 摄取变 AUTO CDC更源。 参见 “应用自动疾病控制中心(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 Declarative Pipelines

接着,将解析输出传递给 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进入卷时,自动加载器会将其作为外部 FILE 引用接收, ai_parse_documentai_classify 路由每份文档,金 consulting_agreements 色实体视图会显示提取的字段。

自己探索吧

该管道将文件分类为五种协议类型,但仅 consulting_agreement提取字段。 要扩展它,对剩余的每个类型重复金步,修改 contract_type 过滤器和 ai_extract 模式以匹配该类型相关的字段。 例如:

  • affiliate_agreementparty_1_nameparty_2_namecommission_ratepayment_frequency
  • marketing_agreementparty_1_nameparty_2_nameeffective_dateterritory
  • hosting_agreementprovider_namecustomer_nameeffective_dateterm_length
  • escrow_agreementowner_namelicensee_nameescrow_agent_namesoftware_name

其他资源