Tutorial: criar um pipeline de processamento de arquivos com o tipo FILE
Beta
Este recurso está em Beta. Os administradores do Workspace podem controlar o acesso a este recurso na página Pré-visualizações . Consulte Gerenciar prévias do Databricks.
Saiba como criar um pipeline medalhão com o Lakeflow pipeline que processa documentos não estruturados de ponta a ponta. Este exemplo usa o dataset de exemplo samples.sec.contracts, uma coleção de contratos legais arquivados na 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 AI, classifica-o em um tipo de contrato e extrai campos estruturados para cada tipo.
Para a referência de tipo, consulte tipoFILE.
Neste tutorial, você:
- Ingira incrementalmente PDFs de contrato de um volume como referências externas de
FILEcom o Auto Loader. - Analise cada documento com a função
ai_parse_documente classifique-o com a funçãoai_classify. - Extraia campos estruturados para cada tipo de contrato com a função
ai_extract.
O resultado é um pipeline no estilo medalhão: bronze (referências externas brutas FILE), silver (documentos analisados e classificados) e ouro (campos extraídos por tipo de contrato). Consulte Qual é a arquitetura medalhão do lakehouse? para obter mais informações. A camada bronze é uma tabela de transmissão que ingere arquivos incrementalmente, e as camadas silver e ouro são views materializadas que recomputam apenas quando suas entradas são alteradas.
[[ ## completed ##]]
Requisitos
Para completar este tutorial, você deve atender aos seguintes requisitos:
- Esteja logado em um workspace do Databricks com o Unity Catalog habilitado.
- Ter permissões para criar tabelas em um esquema e para criar um pipeline.
- Use o canal de pré-visualização.
O dataset samples.sec.contracts está disponível em todos os Workspace por default, portanto, nenhuma configuração adicional é necessária. Como os arquivos já residem em um volume do Unity Catalog, este tutorial os armazena como referências FILE EXTERNAL sem copiar seus conteúdos. 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, consulte Ingerir arquivos como o tipo FILE.
Criar o pipeline de processamento de arquivos
O pipeline processa documentos em três estágios.
O passo 1. Bronze: ingerir PDFs brutos como referências de ARQUIVO externas
Use o Auto Loader para ler incrementalmente os PDFs de contrato do volume. A leitura de 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
- Python
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');
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
FILEleve (uri,size,content_type,checksum). Compare isso com o tipoBINARY, que coloca os bytes inline na linha. - Processamento incremental : a tabela de transmissão ingere incrementalmente novos arquivos à medida que chegam à origem, sem reprocessar os existentes. O dataset
samples.sec.contractsneste exemplo é estático, mas com uma origem ativa, novos arquivos são capturados a cada atualização do pipeline. Para também propagar alterações e exclusões de origem, ingira o feed de alterações comAUTO CDC. Consulte Aplicar atualizações e exclusões com AUTO CDC.
O passo 2. Silver: analisar e classificar documentos
Passe cada FILE para a funçãoai_parse_document para converter o PDF bruto em um VARIANT estruturado contendo elementos do documento, metadados de disposição e texto. Como ai_parse_document aceita uma coluna FILE, ele lê o documento diretamente do armazenamento e nunca carrega os bytes na memória do cluster.
- SQL
- Python
CREATE OR REFRESH MATERIALIZED VIEW parsed_contracts AS
SELECT
path,
ai_parse_document(file) AS parsed
FROM raw_contracts;
@dp.materialized_view(name="parsed_contracts")
def parsed_contracts():
return (
spark.read.table("raw_contracts")
.selectExpr("path", "ai_parse_document(file) AS parsed")
)
Definir o passo de análise como uma view materializada sobre a tabela de transmissão raw_contracts torna a computação incremental. Cada atualização de pipeline realiza a execução de ai_parse_document apenas nos arquivos adicionados desde a última atualização, não na tabela inteira. Como ai_parse_document é o passo mais caro, isso evita reanalisar documentos que você já processou. O refresh incremental de views materializadas requer compute Serverless; faça a execução do pipeline em Serverless. See Spark Declarative Pipelines.
Em seguida, passe a saída analisada para a funçãoai_classify para atribuir a cada documento um dos cinco tipos de contrato. Documentos com erros de análise são filtrados antes da classificação. Este exemplo faz o pin de ai_classify para a versão 2.1, que retorna a classificação como um objeto por rótulo, portanto, leia o rótulo da key value.
[[ ## completed ##]]
- SQL
- Python
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);
@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""")
)
Para melhorar a precisão da classificação, adicione descrições de rótulos e uma opção instructions a ai_classify. Consulte a função ai_classify.
O passo 3. Ouro: extrair campos por tipo de acordo
Cada tipo de contrato tem seu próprio conjunto de campos relevantes. Filtre os documentos classificados para um tipo, passe o conteúdo analisado para a funçãoai_extract com um esquema dos campos desejados e, em seguida, achate a resposta em colunas tipadas. Este exemplo fixa ai_extract na versão 2.1, na qual cada campo extraído é um objeto, portanto, leia sua key value.
O exemplo a seguir cria a tabela ouro para contratos de consultoria:
- SQL
- Python
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;
@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 instruções, você tem um pipeline totalmente incremental: à medida que novos PDFs de contrato chegam ao volume, o Auto Loader os ingere como referências FILE externas, ai_parse_document e ai_classify roteiam cada documento, e a view materializada ouro consulting_agreements exibe os campos extraídos.
[[ ## completed ##]]
Explore por conta própria
O pipeline classifica documentos em cinco tipos de acordo, mas extrai campos apenas para consulting_agreement. Para estendê-lo, repita o passo ouro para cada tipo restante, alterando o filtro contract_type e o esquema ai_extract para corresponder aos campos relevantes para esse tipo. Por exemplo:
affiliate_agreement:party_1_name,party_2_name,commission_rate,payment_frequencymarketing_agreement:party_1_name,party_2_name,effective_date,territoryhosting_agreement:provider_name,customer_name,effective_date,term_lengthescrow_agreement:owner_name,licensee_name,escrow_agent_name,software_name
Recursos adicionais
FILETipo- Ingerir arquivos como o tipo FILE
- Início rápido das funções FILE
- Saiba mais sobre o Auto Loader. Consulte O que é o Auto Loader?.