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

Importante

Esta característica se encuentra en su versión beta. Para usarlo, un administrador de espacio de trabajo debe activar Tipo de archivo desde la página de Previsualizaciones . Consulte Administrar versiones preliminares de Azure Databricks.

Aprende cómo crear una canalización de medallón con Lakeflow Pipeline para procesar 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 canalización ingiere los PDFs como referencias gestionadas FILE con Auto Loader, analiza cada documento con funciones de IA, lo clasifica en un tipo de contrato 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 almacén de lago de datos de medallón? 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.
  • Activa el tipo FILE 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 FileSpace de la tabla de bronce, y Unity Catalog copia los archivos ingeridos en él como almacenamiento gestionado.
  • Uso del canal de versión preliminar

El samples.sec.contracts conjunto de datos está disponible por defecto en todos los espacios de trabajo. Este tutorial almacena los PDFs ingeridos como referencias FILE MANAGED: Unity Catalog copia cada archivo en el volumen que declaras como el FileSpace de la tabla y lo gestiona junto con la tabla, así que eliminar filas hace que los archivos referenciados se puedan recoger mediante el recolector 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 ver otras opciones de importación, consulte Ingerir archivos como el tipo FILE.

Crear la cadena de procesamiento de archivos

El oleoducto procesa documentos en tres etapas.

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

Usa Auto Loader para leer de forma incremental los archivos PDF de contratos desde el volumen de almacenamiento. 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 el FileSpace de la tabla, el volumen que configuras con la propiedad de tabla databricks.filespace-preview, 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 el FileSpace de la tabla, 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 el FileSpace de la tabla 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 de este ejemplo es estático, pero con un origen activo, se detectan 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 a la función ai_parse_document para convertir el PDF sin procesar en VARIANT estructurado que contiene 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 tabla de flujo raw_contracts incrementa el cálculo. Cada actualización de la canalización ejecuta ai_parse_document solo en los archivos que se han añadido desde la última actualización, no en toda la tabla. Como ai_parse_document es el paso más caro, esto evita tener que volver a analizar los documentos que ya has procesado. La actualización incremental de las vistas materializadas requiere cómputo sin servidor; ejecute la canalización en un entorno sin servidor. Consulte Spark Declarative Pipelines.

A continuación, pasa el resultado del análisis a la función ai_classify para asignar a cada documento uno de los cinco tipos de acuerdo. Los documentos con errores de análisis se filtran antes de la clasificación. Este ejemplo fija ai_classify en la versión 2.1, que devuelve la clasificación como un objeto por etiqueta, así que lee la etiqueta de la clave 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

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 único tipo, pasa el contenido procesado a la función ai_extract con un esquema de los campos que desees y, a continuación, convierte la respuesta en columnas tipadas. Este ejemplo fija ai_extract en la versión 2.1, en la que cada campo extraído es un objeto, por lo que lee su clave value.

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_parse_document y ai_classify enrutan cada documento, y la consulting_agreements vista dorada materializada muestra los campos extraídos.

Cuadernos de ejemplo

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

SQL

Cuaderno SQL de la cadena de procesamiento de archivos

Obtener el portátil

Python

Cuaderno Python de la cadena de procesamiento de archivos

Obtener el portátil

Explora por tu cuenta

La canalización clasifica los documentos en cinco tipos de acuerdos, pero extrae campos solo de consulting_agreement. 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