Tutorial: Construa um pipeline de processamento de ficheiros com o tipo FILE

Importante

Este recurso está em versão Beta. Os administradores do espaço de trabalho podem controlar o acesso a esse recurso na página Visualizações . Ver Gerir as pré-visualizações de Azure Databricks.

Aprenda a 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 num volume do Unity Catalog.

O pipeline ingere os PDFs como referências geridas FILE com o Auto Loader, analisa cada documento com funções de IA, classifica-o num tipo de acordo e extrai campos estruturados para cada tipo.

Para a referência tipográfica, veja FILE tipo.

Neste tutorial, você irá:

O resultado é um oleoduto ao estilo medalhão: bronze (referências geridas FILE brutas), prata (documentos analisados e classificados) e ouro (campos extraídos por tipo de acordo). Consulte O que é a arquitetura de medallion lakehouse? para obter mais informações. A camada bronze é uma tabela de fluxo que ingere ficheiros de forma incremental, e as camadas de prata e ouro são vistas materializadas que só se recalculam quando as suas entradas mudam.

Requisitos

Para concluir este tutorial, você deve atender aos seguintes requisitos:

  • Esteja iniciado num espaço de trabalho do Azure Databricks com o Unity Catalog ativado.
  • Tenha o FILE tipo ativado para o seu espaço de trabalho. Os administradores do espaço de trabalho podem ativar isso a partir da página de Pré-visualizações . Ver Gerir as pré-visualizações de Azure Databricks.
  • Ter permissões para criar tabelas num esquema e para criar um pipeline.
  • Tem um volume do Unity Catalog onde possas escrever. Declaras este volume como o da tabela FileSpacebronze , e o Unity Catalog copia os ficheiros ingeridos para ele como armazenamento gerido.
  • Use o canal de Pré-visualização.

O samples.sec.contracts conjunto de dados está disponível em todos os espaços de trabalho por defeito. Este tutorial armazena os PDFs ingeridos como FILE MANAGED referências: o Unity Catalog copia cada ficheiro para o volume que declaras como da tabela FileSpace e gere-o com a tabela, pelo que eliminar linhas torna os ficheiros referenciados elegíveis para recolha de lixo e a tabela e os seus ficheiros permanecem sincronizados. Para adaptar o pipeline aos seus próprios PDFs, aponte o caminho de origem para um volume que contenha os seus ficheiros. Para outras opções de ingestão, veja ficheiros Ingest como o tipo de ficheiro.

Criar o pipeline de processamento de ficheiros

O oleoduto processa documentos em três fases.

Passo 1. Bronze: ingerir PDFs brutos como referências de ficheiro geridas

Use o Auto Loader para ler incrementalmente os PDFs dos contratos do volume. Ler ficheiros com format => 'file' captura uma referência e metadados para cada ficheiro sem materializar os seus bytes. Declarar a coluna como FILE MANAGED copia cada ficheiro para a tabela FileSpace, o volume que defines com a databricks.filespace-preview propriedade tabela, para que o Unity Catalogue gere os ficheiros com a tabela.

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 ficheiros grandes: um PDF grande está nos arquivos da FileSpacetabela , enquanto a linha da tabela armazena apenas uma referência leve FILE (uri, size, content_type, checksum). Compare isto com o BINARY tipo, que faz linhas nos bytes da linha.
  • Ciclo de vida gerido dos ficheiros: O Unity Catalog copia cada ficheiro ingerido para as tabelas FileSpace e gere-o com a tabela: eliminar linhas torna os ficheiros referenciados elegíveis para recolha de lixo, para que a tabela e os seus ficheiros permaneçam sincronizados. Para mais detalhes, consulte FILE MANAGED e FILE EXTERNAL.
  • Processamento incremental: a tabela de streaming ingere progressivamente novos ficheiros à medida que chegam à fonte, sem reprocessar os já existentes. O samples.sec.contracts conjunto de dados neste exemplo é estático, mas com uma fonte ativa, novos ficheiros são captados em cada atualização do pipeline. Para também propagar alterações e eliminações da fonte, ingera o feed de alterações com AUTO CDC. Consulte Aplicar atualizações e eliminações com AUTO CDC.

Passo 2. Prata: analisar e classificar documentos

Passe cada FILE uma para ai_parse_document converter o PDF bruto numa estrutura VARIANT que contenha elementos do documento, metadados de layout e texto. Como ai_parse_document aceita uma FILE coluna, 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 o passo de análise como uma vista materializada sobre a raw_contracts tabela de fluxo incrementaliza o cálculo. Cada atualização do pipeline corre ai_parse_document apenas nos ficheiros adicionados desde a última atualização, não em toda a tabela. Como ai_parse_document é o passo mais dispendioso, isto evita reparar documentos que já processaste. A atualização incremental das vistas materializadas requer computação serverless; Executa o pipeline em serverless. Ver Pipelines Declarativos Spark.

De seguida, passa 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 fixa ai_classify a versão 2.1, que devolve a classificação como um objeto por etiqueta, por isso leia o rótulo 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 etiquetas e uma instructions opção a ai_classify. Consulte a função ai_classify.

Passo 3. Ouro: extrair campos por tipo de acordo

Cada tipo de acordo tem o seu próprio conjunto de campos relevantes. Filtra os documentos classificados para um tipo, passa o conteúdo analisado para ai_extract funcionar com um esquema dos campos que queres e depois achata a resposta em colunas digitadas. Este exemplo liga ai_extract à versão 2.1, em que cada campo extraído é um objeto, por isso leia a sua value chave.

O exemplo seguinte constrói a tabela dourada 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 estas declarações, tens um pipeline totalmente incremental: à medida que novos PDFs de contrato chegam ao volume, o Auto Loader ingere-os como referências geridas FILE , ai_parse_document encaminha ai_classify cada documento, e a consulting_agreements visualização dourada materializada destaca os campos extraídos.

Exemplos de cadernos

Os cadernos seguintes contêm o pipeline completo deste tutorial. Estes cadernos são código-fonte de pipeline, não cadernos executáveis. Importa o caderno para a tua linguagem e depois especifica o seu caminho no campo Código-Fonte quando configurares o pipeline. Veja Configurar pipelines.

SQL

Caderno SQL de pipeline de processamento de ficheiros

Obter bloco de notas

Python

Caderno Python de pipeline de processamento de ficheiros

Obter bloco de notas

Explora por tua conta

O pipeline classifica documentos em cinco tipos de concordância, mas extrai campos apenas consulting_agreementpara . Para prolongar, repita o passo ouro para cada tipo restante, alterando o contract_type filtro e o ai_extract esquema para corresponderem aos campos relevantes para esse tipo. Por exemplo:

  • 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 adicionais