Nota:
El acceso a esta página requiere autorización. Puede intentar iniciar sesión o cambiar directorios.
El acceso a esta página requiere autorización. Puede intentar cambiar los directorios.
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:
- Ingerir de forma incremental PDFs de contratos de un volumen como referencias gestionadas
FILEcon Auto Loader. - Analiza cada documento con
ai_parse_documentfunción y clasifícalo conai_classifyfunción. - Extraer campos estructurados para cada tipo de acuerdo con
ai_extractfunción.
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
FILEpara 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
FileSpacede 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
FileSpacede la tabla, mientras que la fila de la tabla almacena solo una referencia ligeraFILE(uri,size,content_type,checksum). Compáralo con elBINARYtipo, que alinea los bytes de la fila. -
Ciclo de vida gestionado de archivos: Unity Catalog copia cada archivo ingerido en el
FileSpacede 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.contractsconjunto 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 conAUTO 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
Python
Cuaderno Python de la cadena de procesamiento de archivos
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
-
FILEtipo - Importar archivos como tipo FILE
- Inicio rápido de funciones FILE
- Más información sobre Auto Loader. Consulte ¿Qué es Auto Loader?.