Remplissage rétroactif des données historiques avec des pipelines
En Data Engineering, le remplissage rétrospectif désigne le processus de traitement rétroactif des données historiques par le biais d'un pipeline de données conçu pour traiter les données actuelles ou en streaming.
Généralement, il s’agit d’un flux distinct qui envoie des données dans vos tables existantes. L'illustration suivante montre un flux de remplissage historique envoyant des données historiques aux tables bronze de votre pipeline.

Voici quelques scénarios qui pourraient nécessiter un remplissage :
- Traiter les données historiques d'un système hérité pour entraîner un modèle de Machine Learning (ML) ou créer un tableau de bord d'analyse des tendances historiques.
- Retraiter un sous-ensemble de données en raison d'un problème de qualité des données avec des sources de données en amont.
- Vos besoins commerciaux ont changé et vous devez effectuer un remplissage rétroactif des données pour une période différente qui n'était pas couverte par le pipeline initial.
- La logique de votre entreprise a changé et vous devez retraiter les données historiques et actuelles.
Un rattrapage dans les LakeFlow Pipelines est pris en charge par un flux d’ajout spécialisé qui utilise l’option ONCE. Consultez append_flow ou CREATE FLOW (pipelines) pour plus d’informations sur l’option ONCE.
Considérations lors du remplissage des données historiques dans une table de streaming
- Généralement, ajoutez les données à la table de streaming bronze. Les couches Silver et Gold en aval reprendront les nouvelles données de la couche Bronze.
- Assurez-vous que votre pipeline peut gérer les données en double avec élégance si les mêmes données sont ajoutées plusieurs fois.
- Assurez-vous que le schéma des données historiques est compatible avec le schéma des données actuelles.
- Considérez la taille du volume de données et le SLA de temps de traitement requis, et configurez en conséquence les tailles de clusters et de batch.
Exemple : Ajouter un remplissage à un pipeline existant
Dans cet exemple, supposons que vous ayez un pipeline qui ingère les données brutes d'enregistrement d'événements à partir d'une source de stockage cloud, à partir du 01.01.2025. Vous réalisez ensuite que vous souhaitez remplir les trois années précédentes de données historiques pour les cas d'usage d'analyse et de rapports en aval. Toutes les données se trouvent dans un seul emplacement, partitionnées par année, mois et jour, au format JSON.
Pipeline initial
Voici le code de pipeline de départ qui ingère de manière incrémentale les données brutes d'enregistrement d'événements à partir du stockage cloud.
- Python
- SQL
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
)
-- 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
Ici, nous utilisons l'option modifiedAfter Auto Loader pour nous assurer que nous ne traitons pas toutes les données à partir du chemin de stockage cloud. Le traitement incrémentiel est interrompu à cette limite.
D'autres sources de données, telles que Kafka, Kinesis et Azure Event Hub, ont des options de lecture équivalentes pour obtenir le même comportement.
Remplir les données des 3 années précédentes.
Vous souhaitez maintenant ajouter un ou plusieurs flux pour un remplissage des données précédentes. Dans cet exemple, suivez les étapes suivantes :
- Utilisez le flux
append once. Cela effectue un rattrapage unique sans continuer à s'exécuter après ce premier rattrapage. Le code reste dans votre pipeline, et si le pipeline est entièrement actualisé, le rattrapage est réexécuté. - Créez trois flux de remplissage, un pour chaque année (dans ce cas, les données sont divisées par année dans le chemin). Pour Python, nous paramétrons la création des flux, mais en SQL nous répétons le code trois fois, une fois pour chaque flux.
Si vous travaillez sur votre propre projet et que vous n'utilisez pas de compute serverless, vous pouvez mettre à jour le nombre maximum de Workers pour le pipeline. L'augmentation du nombre maximal de Workers garantit que vous disposez des Ressources pour traiter les données historiques tout en continuant de traiter les données de streaming actuelles dans le respect du SLA attendu.
Si vous utilisez le Serverless compute avec l'autoscaling amélioré (le default), votre cluster augmente automatiquement en taille lorsque votre charge augmente.
- Python
- SQL
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
-- 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
);
Cette implémentation met en évidence plusieurs modèles importants.
Séparation des préoccupations
- Le traitement incrémentiel est indépendant des Opérations de rattrapage.
- Chaque flux dispose de sa propre configuration et de ses propres paramètres d'optimisation.
- Il existe une distinction claire entre les Opérations incrémentielles et les Opérations de remplissage.
Exécution contrôlée
- L'utilisation de l'option
ONCEgarantit que chaque remplissage rétroactif s'exécute exactement une seule fois. - Le flux de remplissage rétrospectif reste dans le graphe de pipeline, mais devient inactif une fois terminé. Il est prêt à être utilisé lors du full refresh, automatiquement.
- Il existe une piste d'audit claire des opérations de remplissage des données dans la définition du pipeline.
Optimisation du traitement
- Vous pouvez diviser le grand backfill en plusieurs backfills plus petits pour un traitement plus rapide, ou pour contrôler le traitement.
- L'autoscaling amélioré redimensionne dynamiquement la taille du cluster en fonction de la charge actuelle du cluster.
évolution des schémas
- L'utilisation de
schemaEvolutionMode="addNewColumns"gère les modifications de schéma avec élégance. - Vous disposez d'une inférence de schéma cohérente pour les données historiques et actuelles.
- Il y a une gestion sécurisée des nouvelles colonnes dans les données plus récentes.