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

Modes d’échec et compromis :

Opération de Dataflow pour les charges de travail de streaming

Déploiement, modèles et stratégies de mise à niveau

undefined

undefined

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 :

  1. 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.
  2. 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.
  3. 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.
  4. 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).
  5. 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.
  6. 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.
  7. 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.
  8. 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.
  9. 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 →

Parcourir Google →

Related guides

Accès tout-en-un

Un seul abonnement. Chaque examen.

Chaque plan débloque la recherche de réponses illimitée, les tests pratiques, les explications IA et la bibliothèque complète de ressources — en plus de 20 langues.

Mensuel
24.87
Just €0.83/day
Tout inclus :
  • Recherche de réponses illimitée
  • Tests pratiques illimités
  • Explications basées sur l'IA
  • Bibliothèque complète de ressources
  • Plus de 20 langues
  • Mises à jour de contenu hebdomadaires
  • Récompenses et parrainages
  • Support prioritaire
Commencer l'essai gratuit

Aucune carte de crédit requise*

Meilleur rapport qualité/prix
12 mois
179.87
Just €0.49/daySave 40%
Tout inclus :
  • Recherche de réponses illimitée
  • Tests pratiques illimités
  • Explications basées sur l'IA
  • Bibliothèque complète de ressources
  • Plus de 20 langues
  • Mises à jour de contenu hebdomadaires
  • Récompenses et parrainages
  • Support prioritaire
Commencer l'essai gratuit

Aucune carte de crédit requise*

✓ Plan gratuit inclus · ✓ Annulez à tout moment · ✓ Tous les plans débloquent le produit complet