Reposición de datos históricos con canalizaciones

En la ingeniería de datos, el backfilling se refiere al proceso de procesar de manera retroactiva los datos históricos a través de una canalización de datos diseñada para procesar datos actuales o de streaming.

Normalmente, se trata de un flujo independiente que envía datos a las tablas existentes. En la ilustración siguiente se muestra un flujo de reposición que envía datos históricos a las tablas bronce de la canalización.

Flujo de reposición que agrega datos históricos a un flujo de trabajo existente

Algunos escenarios que pueden requerir un reposición:

  • Procese datos históricos de un sistema heredado para entrenar un modelo de aprendizaje automático (ML) o crear un panel de análisis histórico de tendencias.
  • Reprocesar un subconjunto de datos debido a un problema de calidad en las fuentes de datos ascendentes.
  • Los requisitos empresariales cambiaron y necesita rerrellenar los datos durante un período de tiempo diferente que no estaba cubierto por la canalización inicial.
  • La lógica de negocios ha cambiado y debe volver a procesar los datos históricos y actuales.

Se admite una reposición en las canalizaciones de Lakeflow con un flujo de apéndice especializado que usa la opción ONCE. Consulte append_flow o CREATE FLOW (canalizaciones) para obtener más información sobre la ONCE opción.

Consideraciones al reponer datos históricos en una tabla de transmisión

  • Normalmente, anexe los datos a la tabla de streaming bronze. Las capas de plata y oro de bajada recogerán los nuevos datos de la capa de bronce.
  • Asegúrese de que la canalización puede controlar los datos duplicados correctamente en caso de que los mismos datos se anexen varias veces.
  • Asegúrese de que el esquema de datos históricos es compatible con el esquema de datos actual.
  • Tenga en cuenta el tamaño del volumen de datos y el Acuerdo de Nivel de Servicio de tiempo de procesamiento necesario y, en consecuencia, configure los tamaños del clúster y del lote.

Ejemplo: Adición de un relleno a una tubería existente

En este ejemplo, supongamos que tiene una canalización que ingiere datos de registro de eventos sin procesar desde un origen de almacenamiento en la nube, a partir del 1 de enero de 2025. Más adelante se da cuenta de que desea reponer los datos históricos de los últimos tres años para los casos de uso de informes y análisis posteriores. Todos los datos están en una ubicación, particionada por año, mes y día, en formato JSON.

Canalización inicial

Este es el código de canalización inicial que ingiere incrementalmente los datos de registro de eventos sin procesar del almacenamiento en la nube.

Pitón

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
    )

SQL

-- create a streaming table and the default flow to ingest streaming events
CREATE OR REFRESH STREAMING LIVE TABLE registration_events_raw AS
SELECT * FROM read_files(
  "/Volumes/gc/demo/apps_raw/event_registration/*/*/*",
  format => "json",
  inferColumnTypes => true,
  maxFilesPerTrigger => 100,
  schemaEvolutionMode => "addNewColumns",
  modifiedAfter => "2024-12-31T23:59:59.999+00:00"
)
WHERE year(timestamp) >= '2025'; -- safeguard to not process data before begin_year

Aquí usamos la opción modifiedAfter Auto Loader para asegurarnos de que no estamos procesando todos los datos de la ubicación de almacenamiento en la nube. El procesamiento incremental se corta en ese límite.

Sugerencia

Otros orígenes de datos, como Kafka, Kinesis y Azure Event Hubs, tienen opciones de lector equivalentes para lograr el mismo comportamiento.

Reposición de datos de 3 años anteriores

Ahora quiere agregar uno o varios flujos para rellenar los datos anteriores. En este ejemplo, siga estos pasos:

  • Use l flujo append once. Esto realiza una reposición única sin continuar ejecutándose después de ese rellenado inicial. El código permanece en la canalización, y si esta se actualiza por completo, la reposición de datos se vuelve a ejecutar.
  • Cree tres flujos de reposición, uno para cada año (en este caso, los datos se dividen por año en la ruta). Para Python, parametrizamos la creación de los flujos, pero en SQL repetimos el código tres veces, una vez para cada flujo.

Si está trabajando en un proyecto propio y no usa computación sin servidor, es posible que quiera actualizar el número máximo de trabajadores de la canalización. Aumentar el número máximo de trabajadores garantiza que tiene los recursos para procesar los datos históricos mientras continúa procesando los datos de streaming actuales dentro del SLA (Acuerdo de Nivel de Servicio) esperado.

Sugerencia

Si usa el proceso sin servidor con el escalado automático mejorado (valor predeterminado), el clúster aumenta automáticamente el tamaño cuando aumenta la carga.

Pitón

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

SQL

-- create the streaming table
CREATE OR REFRESH STREAMING TABLE registration_events_raw;

-- append the original incremental, streaming flow
CREATE FLOW
  registration_events_raw_incremental
AS INSERT INTO
  registration_events_raw BY NAME
SELECT * FROM STREAM read_files(
  "/Volumes/gc/demo/apps_raw/event_registration/*/*/*",
  format => "json",
  inferColumnTypes => true,
  maxFilesPerTrigger => 100,
  schemaEvolutionMode => "addNewColumns",
  modifiedAfter => "2024-12-31T23:59:59.999+00:00"
)
WHERE year(timestamp) >= '2025';


-- one time backfill 2024
CREATE FLOW
  registration_events_raw_backfill_2024
AS INSERT INTO ONCE
  registration_events_raw BY NAME
SELECT * FROM read_files(
  "/Volumes/gc/demo/apps_raw/event_registration/year=2024/*/*",
  format => "json",
  inferColumnTypes => true
);

-- one time backfill 2023
CREATE FLOW
  registration_events_raw_backfill_2023
AS INSERT INTO ONCE
  registration_events_raw BY NAME
SELECT * FROM read_files(
  "/Volumes/gc/demo/apps_raw/event_registration/year=2023/*/*",
  format => "json",
  inferColumnTypes => true
);

-- one time backfill 2022
CREATE FLOW
  registration_events_raw_backfill_2022
AS INSERT INTO ONCE
  registration_events_raw BY NAME
SELECT * FROM read_files(
  "/Volumes/gc/demo/apps_raw/event_registration/year=2022/*/*",
  format => "json",
  inferColumnTypes => true
);

Esta implementación resalta varios patrones importantes.

Separación de responsabilidades

  • El procesamiento incremental es independiente de las operaciones de reposición.
  • Cada flujo tiene su propia configuración y configuración de optimización.
  • Hay una distinción clara entre las operaciones incrementales y de reposición.

Ejecución controlada

  • El uso de la ONCE opción garantiza que cada reposición se ejecute exactamente una vez.
  • El flujo de reposición permanece en el gráfico de canalización, pero después de completarse, se vuelve inactivo. Está listo para su uso en la actualización completa, automáticamente.
  • Hay una pista de auditoría clara de las operaciones de reposición en la definición de la canalización.

Optimización del procesamiento

  • Puede dividir las reposiciones grandes en varias reposiciones más pequeñas para acelerar el procesamiento o para tener más control sobre el procesamiento.
  • El escalado automático mejorado escala dinámicamente el tamaño del clúster en función de la carga del clúster actual.

Evolución del esquema

  • Usar schemaEvolutionMode="addNewColumns" maneja los cambios de esquema adecuadamente.
  • Usted dispone de una inferencia de esquema consistente en los datos históricos y actuales.
  • Hay un control seguro de las nuevas columnas en los datos más recientes.

Ejemplo: Reposición de un destino SCD durante una migración

Un escenario común de migración es una tabla de dimensiones que cambian lentamente (SCD) que ya existe en un sistema heredado con años de historia acumulada, pero cuya fuente de cambios original ya no está disponible. Como los eventos de cambio originales ya no existen, en su lugar reproduces una sola vez el historial de la tabla heredada en el nuevo destino AUTO CDC y, a continuación, conectas un nuevo flujo de CDC en adelante. Para obtener más información sobre AUTO CDC y los tipos SCD, consulta Las API de AUTO CDC: simplifican la captura de datos de cambios con canalizaciones.

El patrón es un flujo puntual AUTO CDC hacia la misma tabla de streaming a la que se dirige el flujo continuo AUTO CDC. Un AUTO CDC destino acepta solo flujos AUTO CDC, por lo que la semilla también debe ser un flujo AUTO CDC. Un flujo de adición simple INSERT INTO ONCE en la misma tabla no valida:

  1. Cree la tabla de streaming de destino en la que escribe su flujo AUTO CDC.
  2. Inicialice el historial heredado una sola vez mediante un flujo AUTO CDC ONCE que lee la tabla SCD heredada en transmisión, ordenada según la columna heredada de inicio de validez. Reproduzca las filas heredadas como eventos de cambio en lugar de darles forma usted mismo. AUTO CDC crea las columnas de historial de __START_AT y __END_AT para un destino SCD de tipo 2, así que no escriba esas columnas directamente.
  3. Adjunte el flujo continuo AUTO CDC que lee la fuente de cambios actualizada. AUTO CDC resuelve el orden para cada clave, por lo que la transición debe cumplirse para cada clave de negocio por separado: el primer cambio en producción de cada clave debe producirse después del último cambio generado a partir de datos semilla de esa misma clave. Un valor de secuencia que es simplemente posterior al máximo global de legado puede seguir siendo obsoleto para una clave individual, y el primer cambio en vivo de esa clave se ignora o ordena incorrectamente.

El siguiente código crea una tabla de streaming que sigue los pasos anteriores:

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 flujos deben coincidir en sus claves, su tipo SCD y el tipo de datos de su columna de secuenciación. En el ejemplo anterior, ambos flujos se secuencian según una marca de tiempo, que utiliza un único momento de corte para separar el historial inicial de la fuente en directo. Si las secuencias de la tabla heredada tienen un valor de un tipo diferente al de la transmisión en directo, emite una de ellas para que los tipos coincidan.

La misma forma funciona para un objetivo SCD Tipo 1: cambia STORED AS SCD TYPE 2 a STORED AS SCD TYPE 1 en ambos flujos, y el objetivo mantiene solo la fila actual por clave. Antes de basarse en cualquiera de las dos estructuras, valide con una muestra de claves que el primer cambio en producción de una clave inicializada produzca exactamente una nueva versión y cierre correctamente la versión anterior. Suele aparecer un desfase en la secuencia de teclas en ese paso.

Recursos adicionales