Tutorial: Construye una cadena de procesamiento de archivos con el tipo ARCHIVO

Importante

Esta característica se encuentra en su versión beta. Los administradores del área de trabajo pueden controlar el acceso a esta característica desde la página Vistas previas . Consulte Administrar versiones preliminares de Azure Databricks.

Aprende a construir una pipeline medallion con Lakeflow pipeline que procese documentos no estructurados de principio a fin. Este ejemplo utiliza el samples.sec.contracts conjunto de datos de ejemplo, una colección de acuerdos legales presentados ante la SEC almacenados en PDFs en un volumen del Unity Catalog.

La tubería ingiere los PDFs como referencias gestionadas FILE con Auto Loader, analiza cada documento con funciones de IA, lo clasifica en un tipo de acuerdo y extrae campos estructurados para cada tipo.

Para la referencia tipográfica, véase FILE tipo.

En este tutorial, aprenderá lo siguiente:

El resultado es una tubería tipo medallón: bronce (referencias gestionadas FILE en bruto), plata (documentos analizados y clasificados) y oro (campos extraídos por tipo de acuerdo). Consulte ¿Qué es la arquitectura de medallion lakehouse? para obtener más información. La capa de bronce es una tabla de flujo que ingiere archivos de forma incremental, y las capas de plata y oro son vistas materializadas que solo se recalculan cuando cambian sus entradas.

Requisitos

Para completar este tutorial, debe cumplir los siguientes requisitos:

  • Inicia sesión en un espacio de trabajo de Azure Databricks con Unity Catalog activado.
  • Ten habilitado el FILE tipo para tu espacio de trabajo. Los administradores del espacio de trabajo pueden activarlo desde la página de Previsualizaciones . Consulte Administrar versiones preliminares de Azure Databricks.
  • Tener permisos para crear tablas en un esquema y para crear una canalización.
  • Ten un volumen del Catálogo de Unity donde puedas escribir. Declaras este volumen como el de FileSpacela tabla de bronce , y Unity Catalog copia los archivos ingeridos en él como almacenamiento gestionado.
  • Usa el canal de Vista Previa.

El samples.sec.contracts conjunto de datos está disponible por defecto en todos los espacios de trabajo. Este tutorial almacena los PDFs ingeridos como FILE MANAGED referencias: Unity Catalog copia cada archivo en el volumen que declaras como de la FileSpace tabla y lo gestiona junto con la tabla, así que eliminar filas hace que los archivos referenciados sean elegibles para recogida de basura y la tabla y sus archivos permanecen sincronizados. Para adaptar la canalización a tus propios PDFs, apunta la ruta fuente a un volumen que contenga tus archivos. Para otras opciones de ingestión, consulta archivos Ingest como el tipo de archivo.

Crear la cadena de procesamiento de archivos

El oleoducto procesa documentos en tres etapas.

Paso 1. Bronce: ingirir PDFs en bruto como referencias gestionadas de ARCHIVO

Usa Auto Loader para leer de forma incremental los PDFs de los contratos del volumen. Leer archivos con format => 'file' captura una referencia y metadatos para cada archivo sin materializar sus bytes. Declarar la columna como FILE MANAGED copia cada archivo en la tabla FileSpace, el volumen que configuras con la databricks.filespace-preview propiedad de tabla, así que Unity Catalog gestiona los archivos con la tabla.

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/")
  )
  • Funciona para archivos grandes: un PDF grande está en la tabla FileSpace, mientras que la fila de la tabla almacena solo una referencia ligera FILE (uri, size, content_type, checksum). Compáralo con el BINARY tipo, que alinea los bytes de la fila.
  • Ciclo de vida gestionado de archivos: Unity Catalog copia cada archivo ingerido en las FileSpace tablas y lo gestiona junto con la tabla: eliminar filas hace que los archivos referenciados sean elegibles para recogida de basura, de modo que la tabla y sus archivos permanezcan sincronizados. Para más detalles, véase FILE MANAGED y FILE EXTERNAL.
  • Procesamiento incremental: la tabla de streaming ingiere de forma incremental nuevos archivos a medida que llegan al código-fuente, sin reprocesar los existentes. El samples.sec.contracts conjunto de datos en este ejemplo es estático, pero con un código fuente activo, se recogen nuevos archivos en cada actualización de la canalización. Para propagar también cambios y eliminaciones de fuente, ingiere el feed de cambios con AUTO CDC. Consulta Aplicar actualizaciones y eliminaciones con AUTO CDC.

Paso 2. Plata: analizar y clasificar documentos

Pasa cada FILE función para ai_parse_document convertir el PDF en bruto en una estructura VARIANT que contenga elementos del documento, metadatos de maquetación y texto. Al ai_parse_document aceptar una FILE columna, lee el documento directamente desde el almacenamiento y nunca carga los bytes en la memoria del clúster.

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

Nota:

Definir el paso de análisis como una vista materializada sobre la raw_contracts tabla de flujo incrementa el cálculo. Cada actualización de pipeline se ejecuta ai_parse_document solo en los archivos añadidos desde la última actualización, no en toda la tabla. Como ai_parse_document es el paso más caro, evita que se reparen los documentos que ya has procesado. La actualización incremental de vistas materializadas requiere computación sin servidor; Ejecuta la pipeline en serverless. Consulte Spark Declarative Pipelines.

A continuación, pasa la salida analizada a ai_classify la función para asignar a cada documento uno de cinco tipos de acuerdo. Los documentos con errores de análisis se filtran antes de la clasificación. Este ejemplo se fija ai_classify a la versión 2.1, que devuelve la clasificación como un objeto por etiqueta, por lo que lee la etiqueta desde la value clave.

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

Para mejorar la precisión de la clasificación, añade descripciones de etiquetas y una instructions opción a ai_classify. Consulte la ai_classify función.

Paso 3. Oro: extraer campos por tipo de acuerdo

Cada tipo de acuerdo tiene su propio conjunto de campos relevantes. Filtra los documentos clasificados a un solo tipo, pasa el contenido analizado para ai_extract que funcione con un esquema de los campos que quieres y luego aplana la respuesta en columnas tipadas. Este ejemplo se relaciona ai_extract con la versión 2.1, en la que cada campo extraído es un objeto, así que lee su value clave.

El siguiente ejemplo construye la tabla de oro para acuerdos de consultoría:

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

Con estas sentencias, tienes una pipeline totalmente incremental: a medida que llegan nuevos PDFs de contratos en el volumen, Auto Loader los ingiere como referencias gestionadas FILE y ai_classify enruta cada documento, y la consulting_agreements vista dorada materializada muestra los campos ai_parse_document extraídos.

Cuadernos de ejemplo

Los siguientes cuadernos contienen la cadena completa de este tutorial. Estos cuadernos son código fuente de pipeline, no cuadernos ejecutables. Importa el cuaderno para tu lenguaje y luego especifica su ruta en el campo Código fuente cuando configures la canalización. Consulte Configuración de canalizaciones.

SQL

Cuaderno SQL de la cadena de procesamiento de archivos

Obtención del cuaderno

Python

Cuaderno Python de la cadena de procesamiento de archivos

Obtención del cuaderno

Explora por tu cuenta

El oleoducto clasifica documentos en cinco tipos de acuerdo, pero extrae campos solo consulting_agreementpara . Para extenderlo, repite el paso dorado para cada tipo restante, cambiando el contract_type filtro y el ai_extract esquema para que coincidan con los campos relevantes para ese tipo. Por ejemplo:

  • 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

Recursos adicionales