Aller au contenu principal

Récupérer un pipeline à partir d'une défaillance du point de contrôle de streaming

Récupérer un LakeFlow Pipelines lorsqu'un point de contrôle de streaming devient invalide ou corrompu, en utilisant un refresh complet, une sauvegarde et un remplissage rétroactif, ou une Reset sélective du point de contrôle.

Qu'est-ce qu'un point de contrôle de streaming ?

Dans Apache Spark Structured Streaming, un point de contrôle est un mécanisme utilisé pour persister l'état d'une query de streaming. Cet état inclut :

  • **Informations sur la progression** : Quels décalages de la source ont été traités.
  • État intermédiaire : Données qui doivent être conservées à travers les micro-batchs pour les opérations avec état (par exemple, les agrégations, mapGroupsWithState).
  • Métadonnées : Information sur l'exécution de la query de streaming.

Les points de contrôle sont essentiels pour garantir la tolérance aux pannes et la cohérence des données dans les applications de streaming :

  • Tolérance aux pannes : si une application de streaming échoue (par exemple, en raison d'une défaillance de nœud, d'un plantage d'application), le point de contrôle permet à l'application de redémarrer à partir du dernier état de point de contrôle réussi au lieu de retraiter toutes les données depuis le début. Ceci empêche la perte de données et assure le traitement incrémentiel.
  • Traitement « exactly-once » : Pour de nombreuses sources de streaming, les points de contrôle, conjointement avec les puits idempotents, permettent des garanties de traitement « exactly-once » selon lesquelles chaque enregistrement est traité exactement une fois, même en cas de défaillance, ce qui empêche les doublons ou les omissions.
  • Gestion d'état : Pour les transformations avec état, les points de contrôle conservent l'état interne de ces opérations, permettant à la query de streaming de continuer à traiter correctement de nouvelles données basées sur l'état historique accumulé.

Points de contrôle des pipelines

Les pipelines s'appuient sur Structured Streaming et abstraient une grande partie de la gestion des checkpoints sous-jacente, offrant une approche déclarative. Lorsque vous définissez une table de streaming dans votre pipeline, il existe un état de point de contrôle pour chaque flux écrivant dans la table de streaming. Ces emplacements de checkpoint sont internes au pipeline et ne sont pas accessibles aux utilisateurs.

Vous n'avez généralement pas besoin de gérer ou de comprendre les points de contrôle sous-jacents pour les tables de streaming, sauf dans les cas suivants :

  • **Retour en arrière et relecture** : Si vous souhaitez retraiter les données à partir d'un point spécifique dans le temps tout en conservant l'état actuel de la table, vous devez Reset le point de contrôle de la table de streaming.
  • **Récupération après un échec ou une corruption de point de contrôle** : Si une query écrivant dans la table de streaming a échoué en raison d'erreurs liées aux points de contrôle, cela provoque une défaillance grave, et la query ne peut pas progresser davantage. Il existe trois approches que vous pouvez utiliser pour récupérer de ce type de défaillance :
    • Full table refresh : This Resets the table and wipes out the existing data.
    • Full table refresh with backup and backfill : Vous effectuez une sauvegarde de la table avant d'effectuer un full table refresh et de rattraper les anciennes données, mais cette opération est très coûteuse et ne devrait être utilisée qu'en dernier recours.
    • Reset checkpoint et continuation incrémentielle : si vous ne pouvez pas vous permettre de perdre des données existantes, vous devez effectuer une réinitialisation sélective du point de contrôle pour les flux de streaming concernés.

Exemple : échec du pipeline en raison d'un changement de code

Imaginez un scénario dans lequel vous avez un pipeline qui traite un flux de données modifiées, ainsi que l'instantané initial de la table à partir d'un système de stockage cloud, tel qu'Amazon S3, et écrit dans une table de streaming SCD-1.

Le pipeline comprend deux flux de streaming :

  • customers_incremental_flow: lit de manière incrémentielle le flux CDC de la table source customer, filtre les enregistrements dupliqués et les intègre dans la table cible.
  • customers_snapshot_flow: Lit une seule fois l'instantané initial de la table source customers et met à jour les enregistrements dans la table cible.

Exemple de pipelines CDC pour la récupération après une défaillance du point de contrôle.

Python
@dp.temporary_view(name="customers_incremental_view")
def query():
return (
spark.readStream.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.inferColumnTypes", "true")
.option("cloudFiles.includeExistingFiles", "true")
.load(customers_incremental_path)
.dropDuplicates(["customer_id"])
)

@dp.temporary_view(name="customers_snapshot_view")
def full_orders_snapshot():
return (
spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.includeExistingFiles", "true")
.option("cloudFiles.inferColumnTypes", "true")
.load(customers_snapshot_path)
.select("*")
)

dp.create_streaming_table("customers")

dp.create_auto_cdc_flow(
flow_name = "customers_incremental_flow",
target = "customers",
source = "customers_incremental_view",
keys = ["customer_id"],
sequence_by = col("sequenceNum"),
apply_as_deletes = expr("operation = 'DELETE'"),
apply_as_truncates = expr("operation = 'TRUNCATE'"),
except_column_list = ["operation", "sequenceNum"],
stored_as_scd_type = 1
)
dp.create_auto_cdc_flow(
flow_name = "customers_snapshot_flow",
target = "customers",
source = "customers_snapshot_view",
keys = ["customer_id"],
sequence_by = lit(0),
stored_as_scd_type = 1,
once = True
)

Après le déploiement de ce pipeline, il s'exécute avec succès et commence à traiter le flux de données de modification et l'instantané initial.

Plus tard, vous réalisez que la logique de déduplication dans la query customers_incremental_view est redondante et provoque un goulot d'étranglement des performances. Vous supprimez le dropDuplicates() pour améliorer les performances :

Python
@dp.temporary_view(name="customers_raw_view")
def query():
return (
spark.readStream.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.inferColumnTypes", "true")
.option("cloudFiles.includeExistingFiles", "true")
.load()
# .dropDuplicates()
)

Après avoir supprimé l'API dropDuplicates() et redéployé le pipeline, la mise à jour échoue avec l'erreur suivante :

Streaming stateful operator name does not match with the operator in state metadata.
This is likely to happen when a user adds/removes/changes stateful operators of existing streaming query.
Stateful operators in the metadata: [(OperatorId: 0 -> OperatorName: dedupe)];
Stateful operators in current batch: []. SQLSTATE: 42K03 SQLSTATE: XXKST

Cette erreur indique que la modification n'est pas autorisée en raison d'une incompatibilité entre l'état du checkpoint et la définition de query actuelle, empêchant le pipeline de progresser davantage.

Des échecs liés aux points de contrôle peuvent survenir pour diverses raisons au-delà de la simple suppression de l'API dropDuplicates. Les scénarios courants incluent :

  • Ajout ou suppression d'opérateurs avec état (par exemple, en introduisant ou en supprimant dropDuplicates() ou des agrégations) dans une query de streaming existante.
  • Ajout, suppression ou combinaison de sources de streaming dans une query précédemment à point de contrôle (par exemple, l'union d'une query de streaming existante avec une nouvelle, ou l'ajout/la suppression de sources à partir d'une opération d'union existante).
  • Modification du schéma d'état des opérations de streaming avec état (telles que la modification des colonnes utilisées pour la déduplication ou l'agrégation).

Pour obtenir une liste complète des modifications prises en charge et non prises en charge, reportez-vous au Guide Structured Streaming de Spark et aux Types de modifications dans les queries Structured Streaming.

Options de récupération

Il existe trois stratégies de récupération, en fonction de vos exigences en matière de durabilité des données et de vos contraintes de ressources :

Méthodes

Complexité

Coût

Perte de données potentielle

Duplication potentielle des données

Nécessite un instantané initial.

Full table Reset

Full table refresh

Faible

Medium

Oui (Si aucune capture initiale n'est disponible ou si les fichiers bruts ont été supprimés à la source.)

Non (Pour appliquer les modifications à la table cible.)

Oui

Oui

Refresh complet de la table avec sauvegarde et remplissage

Medium

Haute

Non

Non (Pour les sinks idempotents. Par exemple, auto CDC.)

Non

Non

Reset le point de contrôle de la table

Moyenne-Élevée (Moyenne pour les sources en mode ajout seul qui fournissent des décalages immuables.)

Faible

Non (Nécessite une attention particulière.)

Non (pour les rédacteurs idempotents. Par exemple, CDC auto vers la table cible uniquement.)

Non

Non

Méthodes

Complexité

Coût

Perte de données potentielle

Duplication potentielle des données

Nécessite un instantané initial.

Full table Reset

Full table refresh

Faible

Medium

Oui (Si aucune capture initiale n'est disponible ou si les fichiers bruts ont été supprimés à la source.)

Non (Pour appliquer les modifications à la table cible.)

Oui

Oui

Refresh complet de la table avec sauvegarde et remplissage

Medium

Haute

Non

Non (Pour les sinks idempotents. Par exemple, auto CDC.)

Non

Non

Reset le point de contrôle de la table

Moyenne-Élevée (Moyenne pour les sources en mode ajout seul qui fournissent des décalages immuables.)

Faible

Non (Nécessite une attention particulière.)

Non (pour les rédacteurs idempotents. Par exemple, CDC auto vers la table cible uniquement.)

Non

Non

La complexité moyenne à élevée dépend du type de source de streaming et de la complexité de la query.

Recommandations

  • Utilisez un full table refresh si vous ne voulez pas gérer la complexité d'un checkpoint Reset et que vous pouvez recalculer l'ensemble de la table. Une full refresh vous donne également la possibilité d’apporter des modifications au code.
  • Utilisez un refresh complet de la table avec sauvegarde et remplissage rétroactif si vous ne voulez pas gérer la complexité d'un reset de point de contrôle, et que le coût supplémentaire lié à la sauvegarde et au remplissage rétroactif des données historiques ne vous dérange pas.
  • Utilisez le point de contrôle de Reset de table si vous devez préserver les données existantes dans la table et continuer à traiter les nouvelles données de manière incrémentale. Cependant, cette approche nécessite une manipulation minutieuse du reset du checkpoint afin de s'assurer que les données existantes dans la table ne sont pas perdues et que le pipeline peut continuer à traiter de nouvelles données.

Reset le point de contrôle et continuer de manière incrémentale

Pour Reset le point de contrôle et continuer le traitement de manière incrémentielle, suivez ces étapes :

  1. Arrêtez le pipeline : Assurez-vous que le pipeline n'a pas de mises à jour actives en cours d'exécution.

  2. Déterminez la position de départ du nouveau point de contrôle : identifiez le dernier décalage ou timestamp réussi à partir duquel vous souhaitez continuer le traitement. Il s'agit généralement du dernier offset traité avec succès avant que l'échec ne se produise.

    Étant donné que vous lisez les fichiers JSON à l'aide d'Auto Loader dans l'exemple précédent, vous pouvez utiliser l'option modifiedAfter pour définir un Timestamp pour le moment où Auto Loader start à traiter de nouveaux fichiers.

    Pour les sources Kafka, vous pouvez utiliser l'option startingOffsets afin de spécifier les offsets à partir desquels la query streaming doit start le traitement des nouvelles données.

    Pour les sources Delta Lake, vous pouvez utiliser l'option startingVersion pour spécifier la version à partir de laquelle la query de streaming doit start le traitement des nouvelles données.

  3. Apporter des modifications au code : vous pouvez modifier la query de streaming pour supprimer l'API dropDuplicates() ou apporter d'autres modifications. Vérifiez également que vous avez ajouté l’option modifiedAfter au chemin de lecture de l’Auto Loader.

    Python
    @dp.temporary_view(name="customers_incremental_view")
    def query():
    return (
    spark.readStream.format("cloudFiles")
    .option("cloudFiles.format", "json")
    .option("cloudFiles.inferColumnTypes", "true")
    .option("cloudFiles.includeExistingFiles", "true")
    .option("modifiedAfter", "2025-04-09T06:15:00")
    .load(customers_incremental_path)
    # .dropDuplicates(["customer_id"])
    )
remarque

La fourniture d'un modifiedAfter timestamp incorrect peut entraîner une perte ou une duplication de données. Vérifiez que le timestamp est correctement défini pour éviter de traiter à nouveau les anciennes données ou de manquer les nouvelles données.

Si votre query comporte une jointure stream-stream ou une union stream-stream, vous devez appliquer la stratégie ci-dessus pour toutes les sources de streaming participantes. Par exemple :

Python
cdc_1 = spark.readStream.format("cloudFiles")...
cdc_2 = spark.readStream.format("cloudFiles")...
cdc_source = cdc_1..union(cdc_2)
  1. Identifiez les noms de flux associés à la table de streaming pour laquelle vous souhaitez Reset le point de contrôle. Vous devez transmettre chaque nom de flux dans un format catalog.schema.flow_name entièrement qualifié. Dans l'exemple, le nom de flux entièrement qualifié est my_catalog.my_schema.customers_incremental_flow. Le flux utilise un paramètre flow_name explicite de customers_incremental_flow. Si vous ne définissez pas de nom de flux explicite, le nom de flux default est le nom de table cible entièrement qualifié au format catalog.schema.table (par exemple, my_catalog.my_schema.customers). Vous pouvez trouver les noms de flux dans le code du pipeline, l'interface utilisateur du pipeline ou les logs d'événements du pipeline.

  2. Reset le point de contrôle : créez un Notebook Python et associez-le à un cluster Databricks.

    Vous aurez besoin des informations suivantes pour pouvoir Reset le checkpoint :

    • URL du workspace Databricks
    • ID du pipeline
    • Nom(s) du ou des flux pour lesquels vous réinitialisez le point de contrôle.
    Python
    import requests
    import json

    # Define your Databricks instance and pipeline ID
    databricks_instance = "<DATABRICKS_URL>"
    token = dbutils.notebook.entry_point.getDbutils().notebook().getContext().apiToken().get()
    pipeline_id = "<YOUR_PIPELINE_ID>"
    flows_to_reset = ["<YOUR_FLOW_NAME>"]
    # Set up the API endpoint
    endpoint = f"{databricks_instance}/api/2.0/pipelines/{pipeline_id}/updates"


    # Set up the request headers
    headers = {
    "Authorization": f"Bearer {token}",
    "Content-Type": "application/json"
    }

    # Define the payload
    payload = {
    "reset_checkpoint_selection": flows_to_reset
    }

    # Make the POST request
    response = requests.post(endpoint, headers=headers, data=json.dumps(payload))

    # Check the response
    if response.status_code == 200:
    print("Pipeline update started successfully.")
    else:
    print(f"Error: {response.status_code}, {response.text}")
  3. Exécuter le pipeline : le pipeline start à traiter de nouvelles données à partir de la position de départ spécifiée avec un nouveau point de contrôle, en préservant les données de table existantes tout en poursuivant le traitement incrémentiel.

Bonnes pratiques

  • Évitez d'utiliser les fonctionnalités d'aperçu privé en production.
  • Testez vos modifications avant de les appliquer dans votre environnement de production.
    • Créez un pipeline de test, idéalement dans un environnement inférieur. Si cela n'est pas possible, essayez d'utiliser un autre catalogue et un autre schéma pour votre test.
    • Reproduisez l'erreur.
    • Appliquer les modifications.
    • Validez les résultats et prenez une décision sur le lancement ou non.
    • Déployez les modifications sur vos pipelines de production.