Observação
O acesso a essa página exige autorização. Você pode tentar entrar ou alterar diretórios.
O acesso a essa página exige autorização. Você pode tentar alterar os diretórios.
Importante
Esse recurso está em Beta. Os administradores do workspace podem controlar o acesso a esse recurso na página Visualizações . Consulte Gerenciar visualizações do Azure Databricks.
Aprenda como construir um pipeline medallion com o pipeline Lakeflow que processa documentos não estruturados de ponta a ponta. Este exemplo utiliza o samples.sec.contracts conjunto de dados de exemplo, uma coleção de acordos legais arquivados pela SEC armazenados como PDFs em um volume do Unity Catalog.
O pipeline ingere os PDFs como referências externas FILE com o Auto Loader, analisa cada documento com funções de IA, classifica-o em um tipo de acordo e extrai campos estruturados para cada tipo.
Para a referência tipográfica, veja FILE tipo.
Neste tutorial, você irá:
- Ingerir progressivamente PDFs de contratos de um volume como referências externas
FILEcom o Auto Loader. - Analise cada documento com
ai_parse_documentfunção e classifique-o comai_classifyfunção. - Extraia campos estruturados para cada tipo de concordância com
ai_extractfunção.
O resultado é um oleoduto no estilo medalhão: bronze (referências externas FILE brutas), prata (documentos analisados e classificados) e ouro (campos extraídos por tipo de acordo). Veja O que é a arquitetura medalhão lakehouse? Para obter mais informações. A camada bronze é uma tabela de fluxo que ingere arquivos de forma incremental, e as camadas prata e ouro são vistas materializadas que só se recomputam quando suas entradas mudam.
Requirements
Para concluir este tutorial, você deve atender aos seguintes requisitos:
- Esteja logado em um workspace do Azure Databricks com o Unity Catalog ativado.
- Ter permissões para criar tabelas em um esquema e para criar um pipeline.
- Use o canal Prévia.
O samples.sec.contracts conjunto de dados está disponível em todos os espaços de trabalho por padrão, então não é necessária configuração adicional. Como os arquivos já estão em um volume do Catálogo Unity, este tutorial os armazena como FILE EXTERNAL referências sem copiar seu conteúdo. Para adaptar o pipeline aos seus próprios PDFs, aponte o caminho de origem para um volume que contenha seus arquivos. Para outras opções de ingestão, veja arquivos Ingest como o tipo FILE.
Criar o pipeline de processamento de arquivos
O pipeline processa documentos em três etapas.
Etapa 1. Bronze: ingerir PDFs brutos como referências externas de ARQUIVO
Use o Auto Loader para ler incrementalmente os PDFs dos contratos do volume. Ler arquivos com format => 'file' captura uma referência e metadados para cada arquivo sem materializar seus bytes. Declarar a coluna como FILE EXTERNAL faz referência a cada arquivo no local, sem copiar seu conteúdo.
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/")
)
-
Funciona para arquivos grandes: um PDF grande permanece no volume, enquanto a linha da tabela armazena apenas uma referência leve
FILE(uri,size,content_type,checksum). Compare isso com oBINARYtipo, que faz linhas nos bytes da linha. -
Processamento incremental: a tabela de streaming ingere progressivamente novos arquivos à medida que chegam à fonte, sem reprocessar os existentes. O
samples.sec.contractsconjunto de dados neste exemplo é estático, mas com uma fonte ativa, novos arquivos são captados a cada atualização do pipeline. Para também propagar alterações e deleções de fonte, ingire o feed de mudanças comAUTO CDC. Veja Aplicar atualizações e exclusões com AUTO CDC.
Etapa 2. Prata: analisar e classificar documentos
Passe cada FILEuma paraai_parse_document a função de converter o PDF bruto em uma estrutura VARIANT contendo elementos do documento, metadados de layout e texto. Como ai_parse_document aceita uma FILE coluna, ele lê o documento diretamente do armazenamento e nunca carrega os bytes na memória do cluster.
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")
)
Observação
Definir a etapa de análise como uma vista materializada sobre a raw_contracts tabela de fluxo incrementaliza o cálculo. Cada atualização do pipeline roda ai_parse_document apenas nos arquivos adicionados desde a última atualização, não em toda a tabela. Como ai_parse_document é a etapa mais cara, isso evita reparar documentos que você já processou. A atualização incremental das visualizações materializadas requer computação serverless; Execute o pipeline em serverless. Consulte Pipelines Declarativos do Spark.
Em seguida, passe a saída analisada para ai_classify a função para atribuir a cada documento um dos cinco tipos de acordo. Documentos com erros de análise são filtrados antes da classificação. Este exemplo se conecta ai_classify à versão 2.1, que retorna a classificação como um objeto por etiqueta, então leia o rótulo a partir da value chave.
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""")
)
Dica
Para melhorar a precisão da classificação, adicione descrições de rótulos e uma instructions opção a ai_classify. Veja a ai_classify função.
Etapa 3. Ouro: extrair campos por tipo de acordo
Cada tipo de acordo possui seu próprio conjunto de campos relevantes. Filtre os documentos classificados para um tipo, passe o conteúdo analisado para ai_extract funcionar com um esquema dos campos que você quer, e então achate a resposta em colunas digitadas. Este exemplo se conecta ai_extract à versão 2.1, na qual cada campo extraído é um objeto, então leia sua value chave.
O exemplo a seguir constrói a tabela ouro para acordos de consultoria:
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")
)
Com essas declarações, você tem um pipeline totalmente incremental: à medida que novos PDFs de contrato chegam ao volume, o Auto Loader os ingere como referências externas FILE , ai_parse_document roteia ai_classify cada documento, e a consulting_agreements visualização dourada materializada mostra os campos extraídos.
Explore por conta própria
O pipeline classifica documentos em cinco tipos de acordo, mas extrai campos apenas consulting_agreementpara . Para estendê-lo, repita o passo ouro para cada tipo restante, alterando o contract_type filtro e o ai_extract esquema para corresponder aos campos relevantes para aquele tipo. Por exemplo:
-
affiliate_agreement:party_1_name,party_2_name, ,commission_ratepayment_frequency -
marketing_agreement:party_1_name,party_2_name, ,effective_dateterritory -
hosting_agreement:provider_name,customer_name, ,effective_dateterm_length -
escrow_agreement:owner_name,licensee_name, ,escrow_agent_namesoftware_name
Recursos adicionais
-
FILEtipo - Ingeste arquivos como o tipo FILE
- Funções FILE quickstart
- Saiba mais sobre o Carregador Automático. Confira O que é o Carregador Automático?.