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.
Le flux de remplissage que vous utilisez dépend de la table cible et des données source. Pour une AUTO CDC cible SCD de type 1 à dimension à évolution lente (SCD) avec un instantané officiel, utilisez un AUTO CDC FROM SNAPSHOT flux ponctuel. Pour une migration SCD qui réexécute les modifications historiques, utilisez un flux AUTO CDC ponctuel.
Pour les rétrochargements en mode d’ajout seul : utilisez un flux d’ajout spécialisé avec l’option ONCE pour rétrocharger une table de streaming en mode d’ajout seul. 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
- En règle générale, ajoutez les données à la table de streaming bronze. Les couches silver et gold en aval récupèrent 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.
- Tenez compte du volume de données et de l'accord de niveau de service (SLA) de traitement requis, puis configurez le cluster et les tailles de batch en conséquence.
Exemple : Ajouter un remplissage rétroactif à un pipeline existant
Dans cet exemple, supposons que vous disposez d'un pipeline qui ingère des données d'enregistrement d'événements brutes à partir d'une source de stockage cloud, à compter du 01/01/2025. Vous réalisez ensuite que vous souhaitez effectuer un remplissage rétroactif des trois années précédentes de données historiques pour les cas d'usage de reporting et d'analyse en aval. Toutes les données se trouvent au même endroit, 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.
Ajouter un rattrapage à une table SCD de type 1 AUTO CDC
Utilisez un flux AUTO CDC FROM SNAPSHOT unique pour ajouter un instantané autoritatif à une cible SCD de type 1 qui reçoit également un flux de capture de données de changement (CDC) en continu. La version d’instantané et la colonne de séquençage CDC forment un domaine d’ordonnancement. Un événement CDC plus récent l’emporte sur un instantané plus ancien, tandis qu’un instantané plus récent l’emporte sur un événement CDC plus ancien.
Conditions requises
Avant d'ajouter le remplissage, assurez-vous que les flux remplissent les conditions requises suivantes :
- La cible utilise le SCD de type 1.
- La cible comporte exactement un flux
AUTO CDC FROM SNAPSHOTet un ou plusieurs fluxAUTO CDCau nom unique. - Tous les flux utilisent le même nombre de clés dans le même ordre. Les noms de clés du flux d’instantanés sont comparés sans distinction de majuscules et de minuscules avec les noms de clés
AUTO CDC. Plusieurs fluxAUTO CDCdoivent utiliser des noms de clés et des casses identiques. - La version d'instantané et chaque colonne de séquençage CDC ont exactement le même type de données.
- Le flux
AUTO CDC FROM SNAPSHOTne définit pas d’attentes. - Les flux
AUTO CDCn’utilisent pasIGNORE NULL UPDATES. En Python, ne définissez pasignore_null_updates,ignore_null_updates_column_listouignore_null_updates_except_column_list. - Le pipeline utilise le mode déclenché. Ce modèle ne prend pas en charge les pipelines continus.
Les deux types de flux peuvent utiliser l’interface de pipeline SQL ou Python. Vous pouvez combiner des flux SQL et Python dans la même cible.
L’instantané doit représenter l’état complet de la source à sa version. Si une clé cible est absente de l’instantané, AUTO CDC FROM SNAPSHOT traite cette absence comme une suppression à la version de l’instantané. Un événement CDC doté d’une version plus récente préserve ou restaure la clé.
Ajouter le rattrapage
Pour ajouter un remplissage par instantané ponctuel et continuer à traiter les événements CDC, procédez comme suit :
- Conservez la table cible existante et ses flux
AUTO CDCen cours dans la définition du pipeline. - Définissez l’instantané de référence et sa version. Pour un rappel Python, la première invocation doit renvoyer un instantané et une version. Renvoyez
Noneuniquement après qu’au moins un instantané a été traité. - Ajoutez un flux
AUTO CDC FROM SNAPSHOTaveconce=Trueen Python ouONCEen SQL. Pour un rattrapage SQL dans une cible existante, incluez une queryWITH VERSION. Un flux d'instantané SQL sansWITH VERSIONne prend en charge qu'un chargement initial dans une cible vide. - Exécutez une mise à jour de pipeline Trigger pour traiter le remplissage et les événements CDC continus.
L'exemple suivant commence par un pipeline Python existant qui traite de façon incrémentielle les modifications de customers_cdc. Supposez que vous avez déjà exécuté ce pipeline et alimenté la cible customers :
from pyspark import pipelines as dp
from pyspark.sql.functions import col, expr
@dp.view
def customers_cdc():
return (
spark.readStream.table("main.bronze.customers_cdc")
.withColumn("change_timestamp", col("change_timestamp").cast("timestamp"))
)
dp.create_streaming_table("customers")
dp.create_auto_cdc_flow(
name="customers_incremental_cdc",
target="customers",
source="customers_cdc",
keys=["customer_id"],
sequence_by=col("change_timestamp"),
apply_as_deletes=expr("operation = 'DELETE'"),
except_column_list=["operation", "change_timestamp"],
stored_as_scd_type=1,
)
Pour remplir cette cible existante avec l'état de customers_snapshot au 1er janvier 2025, ajoutez le code suivant à la même définition de pipeline. Conservez la table cible existante et le flux AUTO CDC :
from datetime import datetime, timezone
from typing import Optional, Tuple
from pyspark.sql import DataFrame
backfill_version = datetime(2025, 1, 1, tzinfo=timezone.utc)
def backfill_snapshot_and_version(
latest_snapshot_version: Optional[datetime],
) -> Optional[Tuple[DataFrame, datetime]]:
if latest_snapshot_version is None:
return (spark.read.table("main.legacy.customers_snapshot"), backfill_version)
return None
dp.create_auto_cdc_from_snapshot_flow(
target="customers",
source=backfill_snapshot_and_version,
keys=["customer_id"],
stored_as_scd_type=1,
once=True,
)
Le rappel doit renvoyer un instantané et une version lors de sa première invocation. S’il renvoie None avant qu’un instantané ne soit traité, la mise à jour du pipeline échoue. Une fois l’instantané traité, le fait de renvoyer None indique qu’aucun instantané supplémentaire n’est disponible.
La version d’instantané est un élément Python datetime, qui correspond au type Spark SQL TIMESTAMP. Le flux AUTO CDC existant convertit change_timestamp en TIMESTAMP afin que les deux types de séquençage correspondent exactement. L’exemple utilise Python pour les deux flux, mais vous pouvez définir l’un ou l’autre des flux en SQL et combiner des flux SQL et Python dans la même cible. Pour connaître la syntaxe SQL, y compris la query WITH VERSION requise pour une cible non vide, consultez CREATE FLOW (pipelines).
Une fois que le flux d'instantanés a effectué son commit avec succès, les mises à jour incrémentielles ultérieures l'ignorent pendant que le flux AUTO CDC continue de traiter les nouveaux événements.
Un refresh complet de la cible réexécute le flux d’instantané ponctuel. Conservez l’instantané disponible et assurez-vous qu’il représente toujours l’état souhaité avant d’effectuer un refresh complet.
Ce modèle de remplissage unifié ne prend pas en charge le SCD de type 2 ou les cibles bittemporelles.
Exemple : Remblai une cible SCD lors d’une migration
Un scénario de migration courant est une table à dimensions à évolution lente (SCD) qui existe déjà dans un système hérité avec des années d'historique accumulé, mais dont le flux de modification d'origine n'est plus disponible. Étant donné que les événements de modification d'origine ont disparu, vous rejouez plutôt l'historique de la table héritée dans la nouvelle cible AUTO CDC une seule fois, puis vous y associez un nouveau flux CDC par la suite. Pour en savoir plus sur AUTO CDC et les types SCD, consultez The AUTO CDC APIs: Simplify change data capture with pipelines.
Le modèle est un flux AUTO CDC unique vers la même table de streaming que celle ciblée par le flux AUTO CDC en cours. Une cible AUTO CDC n'accepte que les flux AUTO CDC, le seed doit donc également être un flux AUTO CDC. Un flux d'ajout INSERT INTO ONCE simple vers la même table échoue à la validation :
- Créez la table de streaming cible vers laquelle votre flux
AUTO CDCécrit. - Amorcez l'historique hérité une seule fois avec un flux
AUTO CDC ONCEqui lit la table SCD héritée sous forme de Stream, séquencé par la colonne de start de validité héritée. Rejouez les lignes héritées en tant qu'événements de modification plutôt que de les façonner vous-même.AUTO CDCgénère les colonnes d'historique__START_ATet__END_ATpour une cible SCD de type 2, veillez donc à ne pas écrire directement dans ces colonnes. - Attachez le flux
AUTO CDCen cours qui lit le flux de données de modification récent.AUTO CDCrésout le classement par clé, de sorte que la transition doit être valide pour chaque clé métier individuellement : le premier changement en direct de chaque clé doit être séquencé après le dernier changement initialisé pour cette même clé. Une valeur de séquence qui est simplement postérieure au maximum global hérité peut toujours être obsolète pour une clé individuelle, et le premier changement en direct de cette clé est alors ignoré ou classé de manière incorrecte.
Le code suivant crée une table de streaming qui suit les étapes ci-dessus :
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;
Les deux flux doivent s'accorder sur leurs clés, leur type de SCD et le type de données de leur colonne de séquençage. Dans l'exemple précédent, les deux flux sont séquencés par un timestamp, qui utilise une heure de basculement unique pour séparer l'historique amorcé du flux en direct. Si la table héritée est séquencée par une valeur d'un type différent de celui du flux en direct, effectuez un cast de l'une d'entre elles afin que les types correspondent.
La même forme fonctionne pour une cible SCD de type 1 : remplacez STORED AS SCD TYPE 2 par STORED AS SCD TYPE 1 sur les deux flux, et la cible ne conservera que la ligne actuelle par clé. Avant de vous appuyer sur l'une ou l'autre forme, validez sur un échantillon de clés que le premier changement en direct pour une clé amorcée produit exactement une nouvelle version et ferme correctement la précédente. Un écart de séquençage par clé apparaît généralement à cette étape.
Ressources supplémentaires
- Charger et traiter les données de manière incrémentielle avec LakeFlow Pipelines
- Utilisez les flux dans les LakeFlow Pipelines.
- Les APIs AUTO CDC : simplifiez la capture des modifications de données avec les pipeline
- append_flow
- create_auto_cdc_from_snapshot_flow
- create_auto_cdc_flow
- CRÉER UN FLUX (pipelines)