Google PDE: Traitement en flux continu avec Dataflow et Apache Beam — Guide d'étude
Fait partie du Google Professional Data Engineer — Guide d’étude. Entraînez-vous avec des réponses vérifiées dans le centre d’examens Google, ou passez des tests chronométrés sur ExamRoll.io.
Aperçu
Le traitement de flux sur Google Cloud est centré sur le modèle de programmation unifié d’Apache Beam, exécuté par le runner Dataflow. Beam fournit une abstraction logique — des pipelines de transforms sur des PCollections — qui découple votre code des détails d’exécution tels que le parallélisme, l’autoscaling et la tolérance aux pannes. En streaming, la justesse des résultats dépend de la sémantique temporelle (temps de l’événement vs temps de traitement), du fenêtrage (fixe, glissant, de session, global), des filigranes (watermarks), des déclencheurs (triggers) et de la gestion des données en retard. L’excellence opérationnelle sur Dataflow nécessite un dimensionnement adéquat des workers, une politique d’autoscaling appropriée, le bon moteur de streaming, des choix de shuffle pertinents, une conception de récepteur idempotent, une gestion des messages non livrables et une observabilité robuste.
Modèle Apache Beam et sémantique temporelle
Pipelines, transforms, PCollections, runners :
- Un pipeline Beam applique un graphe orienté acyclique de PTransforms à des PCollections (bornées ou non bornées).
- Les runners (Dataflow, Spark, Flink, Direct) exécutent le pipeline ; Dataflow fournit l’autoscaling géré, le checkpointing et la visibilité opérationnelle.
- Les transforms incluent des opérations élément par élément (ParDo), de regroupement et de combinaison (GroupByKey, Combine), des jointures (CoGroupByKey) et des E/S (IOs) (PubSubIO, BigQueryIO, FileIO).
Fenêtres :
- Fenêtres fixes : tranches de temps qui ne se chevauchent pas (ex. : fenêtres basculantes de 1 minute) pour les agrégats périodiques.
- Fenêtres glissantes : fenêtres qui se chevauchent pour des métriques glissantes lissées (ex. : fenêtres de 5 minutes glissant toutes les minutes).
- Fenêtres de session : fenêtres dynamiques qui se ferment après une période d’inactivité, idéales pour les sessions utilisateur ou les rafales de données d’appareils.
- Fenêtre globale : la vue par défaut non fenêtrée de l’ensemble du flux non borné ; souvent associée à des déclencheurs pour une matérialisation périodique.
Temps de l’événement vs temps de traitement :
- Temps de l’événement : moment où l’événement s’est produit à la source ; permet des agrégations logiquement cohérentes malgré des latences de transport variables.
- Temps de traitement : moment où l’événement est observé par le pipeline ; utile pour les déclencheurs opérationnels mais pas pour la justesse sémantique.
Filigranes (Watermarks) :
- Un filigrane (watermark) estime la complétude du temps de l’événement (la supposition du runner qu’il a vu tous les événements jusqu’à un temps T).
- Les filigranes peuvent avancer de manière irrégulière ou se bloquer en cas de contre-pression ou de retards à la source ; une donnée en retard est toute donnée arrivant avec un horodatage < filigrane.
Déclencheurs (Triggers) et retard :
- Par défaut : déclencheur AfterWatermark qui s’active lorsque le filigrane dépasse la fin de la fenêtre ; avec une tolérance au retard (allowed lateness) = 0, les données en retard sont abandonnées.
- Les déclenchements anticipés (basés sur le temps de traitement ou sur un décompte) donnent des résultats préliminaires à faible latence.
- Les déclenchements tardifs permettent des corrections lorsque des données en retard arrivent ; le mode d’accumulation régit si les volets (panes) accumulent les résultats ou ignorent la sortie précédente.
- Choisissez la tolérance au retard en fonction des exigences métier et des compromis entre stockage et calcul ; un retard plus important augmente la rétention d’état et le coût.
Traitement avec état (Stateful), timers, sessionisation, dédoublonnage :
- Les DoFns avec état (Stateful) conservent un état par clé (ex. : dernier événement vu, agrégats en cours) et définissent des timers pour émettre ou effacer l’état.
- La sessionisation s’exprime naturellement via les SessionWindows ; pour une logique personnalisée, utilisez un état par clé et des timers basés sur le temps de traitement ou le temps de l’événement.
- Dédoublonnage : utilisez un identifiant stable par événement et soit Distinct/Combine par fenêtre, soit un état par clé (ex. : filtre de Bloom ou un ensemble avec TTL). Arbitrez entre la mémoire et les faux positifs par rapport à une précision stricte.
Modes d’échec et compromis :
- L’utilisation de fenêtres basées sur le temps de traitement pour les métriques métier provoque des dérives en cas de pics ou de nouvelles tentatives ; préférez les fenêtres basées sur le temps de l’événement.
- Des fenêtres trop petites avec des déclenchements anticipés fréquents provoquent des émissions de volets (panes) excessives et une amplification d’écriture au niveau du récepteur.
- Une tolérance au retard illimitée peut faire gonfler l’état ; limitez toujours le TTL de l’état et configurez des timers pour effacer les clés dormantes.
Opération de Dataflow pour les charges de travail de streaming
Dimensionnement et autoscaling des workers :
- L’autoscaling horizontal ajoute/supprime des workers en fonction du backlog, du retard du watermark, du CPU et du débit ; définissez une valeur
maxWorkersraisonnable pour absorber les pics. - Choisissez les types de machines en fonction des goulots d’étranglement : limités par le CPU (plus de vCPUs), limités par la mémoire (types à haute mémoire), limités par le réseau (des VM plus grandes réduisent la surcharge du shuffle).
- Augmentez le disque de démarrage pour les shuffles intensifs ou les récepteurs basés sur des fichiers. Surveillez le
system laget lebacklog seconds.
- L’autoscaling horizontal ajoute/supprime des workers en fonction du backlog, du retard du watermark, du CPU et du débit ; définissez une valeur
Streaming Engine et shuffle :
- Streaming Engine externalise l’état et le shuffle vers le backend du service, améliorant l’élasticité, réduisant la pression sur la mémoire des workers et permettant des mises à jour plus rapides.
- Pour les étapes lourdes en traitement par lots ou les regroupements de clés massifs, utilisez Dataflow Shuffle pour décharger les E/S du shuffle des workers. Les deux réduisent les défaillances de workers surchargés (hot-worker) et le
disk thrash.
Contre-pression (backpressure), clés surchargées (hot keys) et asymétrie (skew) :
- Dataflow gère la contre-pression via le rééquilibrage dynamique du travail ; néanmoins, ajustez le contrôle de flux de la source (par ex., nombre de messages/octets en attente pour Pub/Sub) lorsque c’est applicable.
- Les clés surchargées (par ex., des ID populaires) créent des
stragglers(workers à la traîne). Atténuez ce problème avec le partitionnement de clés (key#N), la pré-agrégation partielle puis un nouveau re-clés, ou des approximations basées sur des esquisses (sketch-based). - L’asymétrie due à des enregistrements aberrants (charges utiles énormes) ou à des publicateurs en rafales peut nécessiter des partitions par publicateur, le traitement par lots (batching) ou la compression.
Intégration avec Pub/Sub :
- Utilisez les sujets Pub/Sub pour l’ingestion ; activez les attributs de message pour les métadonnées (par ex., deviceId, horodatage de l’événement).
- Ingestez avec
PubSubIO; extrayez les horodatages des événements depuis les attributs ou la charge utile, sinon utilisez le temps de publication comme solution de repli. - Les clés de tri (
ordering keys) assurent un ordre par clé ; Dataflow nécessite tout de même un comportement idempotent en aval en raison de la livraisonat-least-once(au moins une fois).
Patrons pour le streaming vers BigQuery :
- Préférez
BigQueryIOavec la Storage Write API pour un débit élevé, une faible latence et une sémantiqueexactly-once(exactement une fois) au sein d’un flux via les décalages de flux (stream offsets) et les nouvelles tentatives automatiques. - Pour les pipelines simples à faible débit, les insertions en streaming sont acceptables ; définissez
insertIdpour dédupliquer les nouvelles tentatives du client. - Les requêtes sur les tampons de streaming sont
eventually consistent; pour les analyses critiques en termes de temps, exécutez la requête après un délai de tampon (par ex., attendez ~2x la latence de disponibilité observée), ou matérialisez via des fenêtres de micro-lots et le modecommittedde la Storage Write API.
- Préférez
Effets
exactly-once, idempotence, rejeu et récepteurs (sinks) :- Beam garantit un traitement
at-least-once; leexactly-oncedoit être atteint au niveau du récepteur en utilisant des écritures idempotentes, des transactions ou des clés de déduplication. - BigQuery : utilisez les flux par défaut (
default streams) ou les fluxcommittedde la Storage Write API pour unexactly-onceau sein d’un flux ; avec les insertions en streaming, définissez uninsertIdstable. - Fichiers : écrivez des fichiers temporaires avec des noms uniques, finalisez à la complétion de la fenêtre et assurez des renommages atomiques ; évitez l’écrasement pour prévenir les doublons partiels.
- Bases de données externes : utilisez des
upsertsbasés sur un ID stable ou implémentez des fenêtres de déduplication. - Concevez pour le rejeu : maintenez des transformations déterministes ; assurez-vous que les récepteurs dédupliquent lors des nouvelles tentatives.
- Beam garantit un traitement
Gestion des messages non livrables (dead-letter), routage des erreurs, observabilité :
- Encapsulez l’analyse/l’enrichissement à risque dans un
try/catchau sein d’unParDoet émettez les échecs vers unePCollectiondead-letter via unTupleTag; incluez la charge utile, le code d’erreur et le contexte. - Routez les DLQ (Dead-Letter Queues) vers BigQuery ou Cloud Storage pour analyse ; envisagez un sujet Pub/Sub distinct pour le retraitement.
- Observabilité : utilisez les métriques de job Dataflow (
watermark lag,system lag, débit), les compteurs personnalisés, les métriques de distribution et les journaux par étape dans Cloud Logging. Créez des alertes sur le retard (lag) et les taux d’erreur dans Cloud Monitoring. Utilisez Error Reporting pour agréger les exceptions.
- Encapsulez l’analyse/l’enrichissement à risque dans un
Patrons d’optimisation des performances :
- Lisez efficacement : pour les sources BigQuery, préférez la Storage Read API ou les lectures basées sur des requêtes qui ne sélectionnent que les champs et les filtres nécessaires.
Combine lifting: utilisez descombinerspour réduire le volume du shuffle avantGroupByKey.- Entrées secondaires (
side inputs) : mettez en cache les petites données de référence en mémoire ; surveillez lefanout(distribution) et la cadence de mise à jour. - Sérialisation : utilisez des schémas compacts (Avro/Proto) et évitez l’analyse JSON excessive sur les chemins critiques (
hot paths).
Déploiement, modèles et stratégies de mise à niveau
Flex Templates :
- Empaquetez les pipelines dans des modèles conteneurisés et paramétrés pour des déploiements reproductibles. Les Flex Templates prennent en charge les dépendances personnalisées, les images GPU et l’isolation de l’environnement.
- Externalisez les paramètres d’exécution (par ex., abonnement d’entrée, table de sortie, collecteur de lettres mortes, maxWorkers) pour permettre des déploiements spécifiques à chaque environnement.
Mises à jour de pipeline et compatibilité :
- Dataflow prend en charge la mise à jour sur place pour de nombreux pipelines de streaming si les noms des transformations, les spécifications d’état et les types de sortie restent compatibles. Utilisez des noms de PTransform stables.
- Pour les changements de graphe ou d’état incompatibles, effectuez un basculement contrôlé : démarrez la nouvelle tâche, puis drainez l’ancienne tâche pour terminer le travail en cours et arrêter de lire de nouveaux éléments.
Drainage et snapshots :
- Le drainage termine le traitement de manière propre, écrit les sorties restantes et s’arrête ; coordonnez-le avec la rétention ou les snapshots Pub/Sub pour éviter les pertes de données.
- Pour assurer la continuité, vous pouvez créer un snapshot Pub/Sub, démarrer le nouveau pipeline en reprenant la lecture à partir du snapshot ou d’un horodatage approprié, vérifier la sortie, puis drainer l’ancienne tâche.
Exemples de configuration :
- Exemple de fenêtrage avec déclencheurs précoces/tardifs et accumulation :
undefined
- Exemple de BigQueryIO avec la Storage Write API :
undefined
- Pièges courants :
- Écrire vers des collecteurs basés sur des fichiers en streaming sans écritures fenêtrées peut bloquer la finalisation ; activez les écritures fenêtrées et les déclencheurs.
- Croissance illimitée : oublier de limiter l’état ou la latence autorisée peut provoquer des fuites de mémoire et des échecs de mise à l’échelle.
- Horodatages manquants : ne pas assigner d’horodatages d’événement force le pipeline à utiliser par défaut le temps de traitement et lui fait perdre en exactitude en cas de retards variables.
Scénario de problème pratique
NovaTrack Inc. ingère la télémétrie IoT mondiale de 50 000 capteurs de température et doit fournir des agrégats à la minute, persister les données brutes et alimenter un tableau de bord en temps réel. Des messages mal formés occasionnels et une livraison dans le désordre sont attendus. La solution doit se mettre à l’échelle automatiquement, exposer les enregistrements incorrects pour inspection et prendre en charge les mises à niveau sans interruption de service.
Approche :
Ingestion et sémantique temporelle
- Créez un sujet Pub/Sub régional et des publicateurs par région avec les attributs deviceId et eventTs (RFC3339). Activez les clés d’ordonnancement (ordering keys) par deviceId lorsque c’est possible.
- Justification : Pub/Sub fournit une porte d’entrée (ingress) durable et élastique avec une livraison “au moins une fois”. Attacher les horodatages d’événement à la source préserve le véritable temps de l’événement ; l’ordonnancement par appareil réduit le réordonnancement intra-appareil sans créer de goulets d’étranglement centraux.
Pipeline de streaming Dataflow avec des fenêtres basées sur le temps de l’événement
- Lisez depuis un abonnement dédié via PubSubIO, en extrayant eventTs comme horodatage Beam, en vous rabattant sur publishTime s’il est manquant.
- Appliquez des FixedWindows de 1 minute avec un déclencheur précoce à 30 secondes et des déclenchements tardifs sur chaque élément en retard ; définissez une latence autorisée (allowed lateness) de 10 minutes et des panneaux cumulatifs (accumulating panes).
- Justification : Les fenêtres basées sur le temps de l’événement garantissent des agrégats à la minute précis ; les déclenchements précoces alimentent le tableau de bord avec une fraîcheur inférieure à la minute ; les déclenchements tardifs corrigent les agrégats à mesure que les données retardées arrivent. La limite de latence plafonne la taille de l’état et le coût.
Validation, enrichissement et routage des lettres mortes
- Implémentez un ParDo qui analyse le JSON, valide le schéma et les plages de valeurs, et enrichit les données avec de petites données de référence statiques via une entrée secondaire (side input) chargée depuis BigQuery au démarrage de la tâche.
- Utilisez des TupleTags pour émettre les enregistrements valides vers la sortie principale et les échecs vers une PCollection de lettres mortes (dead-letter) contenant la charge utile (payload), l’erreur, le deviceId et l’horodatage de l’analyse ; écrivez la DLQ (Dead-Letter Queue) dans une table BigQuery partitionnée.
- Justification : Les entrées secondaires (side inputs) conservent les données de référence en mémoire pour une faible latence. La capture des lettres mortes permet l’inspection et le retraitement ciblé des lignes incorrectes sans bloquer le flux principal.
Agrégation et atténuation des clés surchargées (hot keys)
- Cléifiez par deviceId et calculez la moyenne/min/max par minute avec des CombineFns. Pour les métriques régionales top-N, partitionnez (shard) par région#N pour éviter les clés surchargées, puis ré-agrégez.
- Justification : Les combineurs minimisent le volume de données brassées (shuffle) et le coût ; le partitionnement des clés (key sharding) empêche les goulets d’étranglement sur une seule clé lors de la centralisation régionale (fan-in).
Collecteurs et effets “exactly-once”
- Écrivez les événements bruts validés et les agrégats à la minute dans BigQuery en utilisant BigQueryIO avec la Storage Write API. Définissez un ID d’insertion stable basé sur deviceId + eventTs pour l’idempotence en cas de nouvelles tentatives personnalisées.
- Justification : La Storage Write API fournit une ingestion à haut débit et faible latence avec une sémantique “exactly-once” au sein d’un flux. Des ID stables garantissent le dédoublonnage en aval en cas de relectures.
Stratégie de cohérence du tableau de bord
- Le tableau de bord interroge les tables d’agrégats partitionnées avec une fenêtre de consultation de 2 minutes par rapport au watermark ou un délai fixe de 2x la latence de disponibilité observée pour les données en streaming.
- Justification : La visibilité du streaming dans BigQuery est cohérente à terme (eventually consistent) ; différer légèrement les lectures empêche de manquer des lignes en transit tout en conservant un comportement quasi temps réel.
Opérations : autoscaling et moteur de streaming
- Activez le Streaming Engine ; définissez maxWorkers en fonction du pic attendu (par ex., 3x la moyenne), sélectionnez un type de machine dimensionné pour l’analyse (parsing) et le chiffrement qui sont gourmands en CPU, et augmentez le disque de démarrage pour accommoder le brassage (shuffle) transitoire.
- Surveillez le retard du watermark (watermark lag), le backlog en secondes, le CPU et le débit par étape ; alertez en cas de retard soutenu et de pics de débit dans la DLQ.
- Justification : Le Streaming Engine externalise l’état et le brassage (shuffle) pour plus d’élasticité et des mises à niveau plus simples ; un dimensionnement correct (right-sizing) et une surveillance adéquate préviennent les violations silencieuses des SLO.
Déploiement et mises à niveau avec les Flex Templates
- Empaquetez le pipeline en tant que Flex Template avec des paramètres : abonnement d’entrée, tables de sortie, table DLQ, maxWorkers et région. Pour un changement incompatible, démarrez le nouveau pipeline ciblant le même sujet avec un nouvel abonnement, vérifiez les sorties, puis drainez l’ancienne tâche. Optionnellement, créez un snapshot Pub/Sub et faites en sorte que le nouvel abonnement reprenne la lecture à partir du snapshot pour garantir aucune perte de données.
- Justification : Les Flex Templates permettent des déploiements reproductibles et paramétrés. Un basculement de type bleu/vert (blue/green) vérifié avec drainage permet d’atteindre zéro perte de données et un temps d’arrêt minimal.
Retraitement et rattrapages par lots (backfills)
- Stockez des fichiers Avro compressés des événements bruts dans Cloud Storage via une sortie secondaire (side output) ; exécutez un pipeline Dataflow par lots (batch) pour effectuer un rattrapage (backfill) ou retraiter les données vers BigQuery lorsque les modèles ou les schémas chang
← Analytique BigQuery et ingénierie d’entrepôt de données · Tous les domaines · Messagerie →
Entraînez-vous sur ces questions → · Tests chronométrés sur ExamRoll.io →
Pass the whole exam — not just this question
You found this answer. Get every verified question and explanation in one place, and save hours of prep. Free to start.
Réussissez votre examen →