Pular para o conteúdo principal

Preenchimento retroativo de dados históricos com oleoduto

Na engenharia de dados, backfilling refere-se ao processo de processamento retroativo de dados históricos por meio de um pipeline de dados que foi projetado para processar dados atuais ou de transmissão.

Normalmente, esse é um fluxo separado que envia dados para suas tabelas existentes. A ilustração a seguir mostra um fluxo de preenchimento enviando dados históricos para as tabelas de bronze no seu pipeline.

Fluxo de backfill adicionando dados históricos a um fluxo de trabalho existente

Alguns cenários que podem exigir um preenchimento:

  • Processe dados históricos de um sistema legado para ensinar um modelo machine learning (ML) ou construir um painel de análise de tendências históricas.
  • Reprocesse um subconjunto de dados devido a um problema de qualidade de dados com fonte de dados upstream.
  • Seus requisitos de negócios mudaram e você precisa preencher dados para um período de tempo diferente que não foi coberto pelo pipeline inicial.
  • Sua lógica de negócios mudou e você precisa reprocessar dados históricos e atuais.

O preenchimento retroativo em Lakeflow pipelines é compatível com um fluxo de acréscimo especializado que usa a opção ONCE. Consulte append_flow ou CREATE FLOW (pipelines) para obter mais informações sobre a opção ONCE.

Considerações ao preencher dados históricos em uma tabela de transmissão

  • Normalmente, anexe os dados à tabela de transmissão de bronze. As camadas de prata e ouro a jusante coletarão os novos dados da camada de bronze.
  • Garanta que seu pipeline possa lidar com dados duplicados com eficiência caso os mesmos dados sejam anexados várias vezes.
  • Garanta que o esquema histórico de dados seja compatível com o esquema de dados atual.
  • Considere o tamanho do volume de dados e o tempo de processamento necessário SLA) e configure adequadamente os tamanhos cluster e dos lotes.

Exemplo: Adicionar um aterro a um pipeline existente

Neste exemplo, digamos que você tenha um pipeline que ingere dados brutos de registro de eventos de uma fonte de armazenamento cloud , a partir de 1º de janeiro de 2025. Mais tarde, você percebe que deseja preencher os três anos anteriores de histórico de dados para casos de uso de relatórios e análises posteriores. Todos os dados estão em um local, particionados por ano, mês e dia, no formato JSON.

pipelineinicial

Este é o código pipeline inicial que ingere incrementalmente os dados brutos de registro de eventos do armazenamento cloud .

Python
from pyspark import pipelines as dp

source_root_path = spark.conf.get("registration_events_source_root_path")
begin_year = spark.conf.get("begin_year")
incremental_load_path = f"{source_root_path}/*/*/*"

# create a streaming table and the default flow to ingest streaming events
@dp.table(name="registration_events_raw", comment="Raw registration events")
def ingest():
return (
spark
.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.inferColumnTypes", "true")
.option("cloudFiles.maxFilesPerTrigger", 100)
.option("cloudFiles.schemaEvolutionMode", "addNewColumns")
.option("modifiedAfter", "2025-01-01T00:00:00.000+00:00")
.load(incremental_load_path)
.where(f"year(timestamp) >= {begin_year}") # safeguard to not process data before begin_year
)

Aqui usamos a opção modifiedAfter Auto Loader para garantir que não estamos processando todos os dados do caminho de armazenamento cloud . O processamento incremental é interrompido nesse limite.

dica

Outras fontes de dados, como Kafka, Kinesis e Azure Event Hubs, têm opções de leitor equivalentes para atingir o mesmo comportamento.

Dados de preenchimento dos últimos 3 anos

Agora você deseja adicionar um ou mais fluxos para preencher dados anteriores. Neste exemplo, tomemos os seguintes passos:

  • Use o fluxo append once . Isso executa um preenchimento único sem continuar a execução após esse primeiro preenchimento. O código permanece no seu pipeline e, se o pipeline for totalmente atualizado, o backfill será reexecutado.
  • Crie três fluxos de preenchimento, um para cada ano (nesse caso, os dados são divididos por ano no caminho). No Python, parametrizamos a criação dos fluxos, mas no SQL repetimos o código três vezes, uma para cada fluxo.

Se você estiver trabalhando em seu próprio projeto e não estiver usando compute serverless , talvez seja necessário atualizar o número máximo de trabalhadores para o pipeline. Aumentar o número máximo de trabalhadores garante que você tenha o recurso para processar o histórico de dados enquanto continua a processar os dados de transmissão atuais dentro do SLA esperado.

dica

Se você usar compute serverless com dimensionamento automático aprimorado (o default), seu cluster aumentará automaticamente de tamanho quando sua carga aumentar.

Python
from pyspark import pipelines as dp

source_root_path = spark.conf.get("registration_events_source_root_path")
begin_year = spark.conf.get("begin_year")
backfill_years = spark.conf.get("backfill_years") # e.g. "2024,2023,2022"
incremental_load_path = f"{source_root_path}/*/*/*"

# meta programming to create append once flow for a given year (called later)
def setup_backfill_flow(year):
backfill_path = f"{source_root_path}/year={year}/*/*"
@dp.append_flow(
target="registration_events_raw",
once=True,
name=f"flow_registration_events_raw_backfill_{year}",
comment=f"Backfill {year} Raw registration events")
def backfill():
return (
spark
.read
.format("json")
.option("inferSchema", "true")
.load(backfill_path)
)

# create the streaming table
dp.create_streaming_table(name="registration_events_raw", comment="Raw registration events")

# append the original incremental, streaming flow
@dp.append_flow(
target="registration_events_raw",
name="flow_registration_events_raw_incremental",
comment="Raw registration events")
def ingest():
return (
spark
.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.inferColumnTypes", "true")
.option("cloudFiles.maxFilesPerTrigger", 100)
.option("cloudFiles.schemaEvolutionMode", "addNewColumns")
.option("modifiedAfter", "2024-12-31T23:59:59.999+00:00")
.load(incremental_load_path)
.where(f"year(timestamp) >= {begin_year}")
)

# parallelize one time multi years backfill for faster processing
# split backfill_years into array
for year in backfill_years.split(","):
setup_backfill_flow(year) # call the previously defined append_flow for each year

Esta implementação destaca vários padrões importantes.

Separação de preocupações

  • O processamento incremental é independente das operações de preenchimento.
  • Cada fluxo tem suas próprias configurações e configurações de otimização.
  • Há uma distinção clara entre operações incrementais e de preenchimento.

Execução controlada

  • Usar a opção ONCE garante que cada preenchimento seja executado exatamente uma vez.
  • O fluxo de preenchimento permanece no gráfico pipeline , mas se torna paradoxo após sua conclusão. Ele está pronto para uso na refresh completa, automaticamente.
  • Há uma trilha de auditoria clara das operações de aterramento na definição do pipeline.

Otimização de processamento

  • Você pode dividir o preenchimento grande em vários preenchimentos menores para um processamento mais rápido ou para controlar o processamento.
  • Usando o dimensionamento automático aprimorado, escale dinamicamente o tamanho cluster com base na carga atual cluster .

evolução do esquema

  • Usar schemaEvolutionMode="addNewColumns" manipula alterações de esquema com elegância.
  • Você tem inferência de esquema consistente em dados históricos e atuais.
  • Há um tratamento seguro de novas colunas em dados mais recentes.

Exemplo: preencher um destino SCD durante uma migração

Um cenário comum de migração é uma tabela de dimensão de alteração lenta (SCD) que já existe em um sistema legado com anos de histórico acumulado, mas cujo registro de alterações original não está mais disponível. Como os eventos de mudança originais desapareceram, você repete o histórico da tabela antiga no novo alvo AUTO CDC uma vez e anexa um novo feed do CDC daqui para frente. Para obter mais informações sobre AUTO CDC e tipos SCD, consulte As APIs AUTO CDC: Simplifique a captura de dados de alterações (CDC) com pipelines.

O padrão é um fluxo AUTO CDC único para a mesma tabela de transmissão que o fluxo AUTO CDC contínuo tem como destino. Um destino AUTO CDC aceita apenas fluxos AUTO CDC, portanto, a semente também deve ser um fluxo AUTO CDC. Um fluxo de acréscimo INSERT INTO ONCE simples para a mesma tabela falha na validação:

  1. Crie a tabela de transmissão de destino na qual seu fluxo AUTO CDC grava.
  2. Preencha o histórico legado uma vez com um fluxo AUTO CDC ONCE que lê a tabela SCD legada como uma transmissão, sequenciada pela coluna de início de validade legada. Reproduza as linhas legadas como eventos de mudança em vez de formatá-las você mesmo. AUTO CDC cria as colunas de história __START_AT e __END_AT para um destino SCD Tipo 2, portanto, não grave nessas colunas diretamente.
  3. Anexe o fluxo AUTO CDC em andamento lendo o feed de dados alterados mais recente. AUTO CDC resolve a ordenação por key, portanto, a transição deve ser válida para cada key de negócio individualmente: a primeira alteração em tempo real de cada key deve ser sequenciada após a última alteração semeada para essa mesma key. Um valor de sequência que seja apenas posterior ao máximo legado global ainda pode estar obsoleto para uma key individual, e a primeira alteração em tempo real dessa key é então ignorada ou ordenada incorretamente.

O código a seguir cria uma tabela de transmissão que usa os passos acima:

SQL
CREATE OR REFRESH STREAMING TABLE customers_history;

-- One-time seed: replay the legacy history as change events
CREATE FLOW customers_history_seed
AS AUTO CDC ONCE INTO customers_history
FROM stream(legacy.customers_scd2)
KEYS (customer_id)
SEQUENCE BY valid_from
STORED AS SCD TYPE 2;

-- Ongoing live CDC into the same target
CREATE FLOW customers_history_cdc
AS AUTO CDC INTO customers_history
FROM stream(customers_cdc_bronze)
KEYS (customer_id)
SEQUENCE BY change_timestamp
STORED AS SCD TYPE 2;

Ambos os fluxos devem concordar quanto às suas keys, ao seu tipo de SCD e ao tipo de dados da sua coluna de sequenciamento. No exemplo anterior, ambos os fluxos são sequenciados por um Timestamp, que utiliza um único tempo de transição para separar a história semeada do feed ativo. Se a tabela legada for sequenciada por um valor de um tipo diferente do feed ativo, converta um deles para que os tipos correspondam.

O mesmo formato funciona para um destino SCD tipo 1: altere STORED AS SCD TYPE 2 para STORED AS SCD TYPE 1 em ambos os fluxos, e o destino mantém apenas a linha atual por key. Antes de confiar em qualquer um dos formatos, valide em uma amostra de chaves se a primeira alteração ativa para uma chave semeada produz exatamente uma nova versão e fecha corretamente a anterior. Uma lacuna de sequenciamento por key geralmente aparece nessa etapa.

Recurso adicional