Учебник: Постройте конвейер обработки файлов с типом FILE

Important

Эта функция доступна в бета-версии. Администраторы рабочей области могут управлять доступом к этой функции на странице "Предварительные версии ". См. статью "Управление предварительными версиями Azure Databricks".

Узнайте, как построить модульный пайплайн с помощью Lakeflow, который обрабатывает неструктурированные документы от до конца. В этом примере используется samples.sec.contracts примерный набор данных — коллекция юридических соглашений, поданных SEC, хранящихся в виде PDF в томе Unity Catalog.

Конвейер принимает PDF-файлы как управляемые FILE ссылки с помощью Auto Loader, анализирует каждый документ с помощью AI-функций, классифицирует его по типу соглашения и извлекает структурированные поля для каждого типа.

Для справочника типа см.FILE тип.

При работе с этим руководством вы сделаете следующее:

В результате получается трубопровод в виде медальона: бронза (необработанные управляемые FILE ссылки), серебро (анализированные и засекреченные документы) и золото (извлекаемые поля по типу соглашения). Дополнительные сведения см. в статье об архитектуре medallion lakehouse. Бронзовый слой — это потоковая таблица , которая постепенно поглощает файлы, а серебряные и золотые слои — это материализованные виды , которые пересчитываются только при изменении входных данных.

Требования

Чтобы завершить работу с этим руководством, необходимо выполнить следующие требования:

  • Войдите в рабочее пространство Azure Databricks с включённым каталогом Unity.
  • Включите FILE тип для вашего рабочего пространства. Администраторы рабочего пространства могут включить его на странице предпросмотра . См. статью "Управление предварительными версиями Azure Databricks".
  • Имейте права на создание таблиц в схеме и для создания конвейера.
  • Имейте том Unity Catalog, в который можно писать. Вы объявляете этот том как бронзовые таблицы FileSpace, и Unity Catalog копирует погружённые файлы в него как управляемое хранилище.
  • Используйте канал предварительного просмотра.

Набор samples.sec.contracts данных доступен во всех рабочих пространствах по умолчанию. В этом учебном руководстве сохраняются введённые PDF-файлы как FILE MANAGED ссылки: Unity Catalog копирует каждый файл в том, который вы объявляете как таблицу FileSpace , и управляет им вместе с таблицей, поэтому удаление строк делает ссылки пригодными для сборки мусора, а таблица и её файлы остаются синхронизированными. Чтобы адаптировать конвейер к вашим собственным PDF-файлам, укажите исходный путь на том, который содержит ваши файлы. Для других вариантов загрузки см. Ingest files как тип FILE.

Создать конвейер обработки файлов

Конвейер обрабатывает документы в три этапа.

Этап 1. Бронза: вводите сырые PDF-файлы как управляемые ссылки на FILE

Используйте Auto Loader, чтобы постепенно читать PDF-файлы контракта из тома. Чтение файлов с format => 'file' файлами фиксирует ссылку и метаданные для каждого файла, не материализуя его байты. Объявление столбца как FILE MANAGED копирует каждый файл в таблицу FileSpace, том, который вы задаёте с databricks.filespace-preview помощью свойства таблицы, так что Unity Catalog управляет файлами с этой таблицей.

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/")
  )
  • Работает для больших файлов: большой PDF находится в таблице FileSpace, а в строке таблицы хранится только лёгкая FILE ссылка (uri, size, content_typechecksum, ). Сравните это с BINARY типом, который выстраивает байты в строке.
  • Управляемый жизненный цикл файла: Unity Catalog копирует каждый погружённый файл в таблицы FileSpace и управляет им вместе с таблицей: удаление строк делает ссылки на ссылки подходящими для сборки мусора, чтобы таблица и её файлы оставались синхронизированы. Для подробностей см. FILE MANAGED и FILE EXTERNAL.
  • Инкрементальная обработка: таблица потоковой передачи постепенно принимает новые файлы по мере их поступления в исходник, не перерабатывая существующие. 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")
  )

Примечание.

Определение шага разбора как материализированного представления над таблицей 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""")
  )

Tip

Для повышения точности классификации добавьте описания ярлыков и опцию 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_documentai_classify каждый документ, и consulting_agreements золотой материализованный вид появляется на поверхности извлечённых полей.

Примеры записных книжек

Следующие блокноты содержат полный конвейер из этого учебника. Эти блокноты — это исходный код конвейера, а не управляемые блокноты. Импортируйте блокнот для вашего языка, затем укажите его путь в поле исходного кода при настройке конвейера. См. раздел "Настройка конвейеров".

SQL

Файловый конвейер SQL notebook

Получите ноутбук

Python

Блокнот для обработки файлов на Python

Получите ноутбук

Исследуйте самостоятельно

Конвейер классифицирует документы на пять типов соглашений, но извлекает поля только consulting_agreementдля . Чтобы расширить его, повторите шаг золота для каждого оставшегося типа, изменяя contract_type фильтр и ai_extract схему так, чтобы они соответствовали полям, релевантным этому типу. Рассмотрим пример.

  • affiliate_agreement: party_1_name, , party_2_namecommission_ratepayment_frequency
  • marketing_agreement: party_1_name, , party_2_nameeffective_dateterritory
  • hosting_agreement: provider_name, , customer_nameeffective_dateterm_length
  • escrow_agreement: owner_name, , licensee_nameescrow_agent_namesoftware_name

Дополнительные ресурсы