Google PDE: Messagerie, ingestion d'événements et services en temps réel — 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.
Vue d’ensemble
Les services de messagerie, d’ingestion d’événements et de temps réel sur Google Cloud s’articulent autour de Cloud Pub/Sub et Eventarc pour un transport découplé et durable ; Dataflow pour le traitement de flux avec état (stateful) ; et des récepteurs (sinks) tels que BigQuery, Cloud Storage et les bases de données opérationnelles. La conception pour une livraison au moins une fois (at-least-once), une consommation idempotente et une observabilité garantit des systèmes résilients qui s’adaptent de manière élastique tout en maintenant l’exactitude en cas de panne, de contre-pression (backpressure) et d’évolution des schémas.
Messagerie principale avec Pub/Sub
- Sujets et abonnements
- Les éditeurs (publishers) envoient des messages à un sujet (topic) ; les abonnés s’y attachent via des abonnements (plusieurs abonnés peuvent consommer indépendamment les mêmes messages).
- Types d’abonnement :
- Pull : les clients tirent explicitement les messages ; utilisez le streaming pull pour un débit maximal et moins d’allers-retours.
- Push : Pub/Sub livre les messages via HTTPS ; votre point de terminaison (endpoint) doit renvoyer un code 2xx pour accuser réception.
- Export vers BigQuery : un abonnement BigQuery livre les messages à une table BigQuery sans code ; idéal lorsque les charges utiles (payloads) correspondent au schéma déclaré et qu’une ingestion à faible latence dans les outils d’analyse est requise.
- Clés d’ordonnancement
- Activez l’ordonnancement des messages sur le sujet et l’abonnement pour recevoir les livraisons dans l’ordre par clé d’ordonnancement. Le débit par clé est sérialisé : un message en transit par clé peut bloquer les suivants ; utilisez de nombreuses clés (par exemple, hash(device_id)) pour passer à l’échelle.
- Fan-out et relecture (replay)
- Créez des abonnements distincts pour différents consommateurs afin d’isoler les charges de travail et la rétention.
- Utilisez la recherche (seek) ou un instantané (snapshot) pour relire à partir d’un horodatage ou d’un instantané pour la récupération et les rattrapages (backfills).
Compromis :
- L’ordonnancement réduit le parallélisme et le débit par clé ; désactivez-le sauf si c’est strictement nécessaire.
- Le mode Push simplifie le code client mais introduit des problématiques de mise à l’échelle du point de terminaison HTTP, de sécurité et de backoff ; le mode Pull offre plus de contrôle et de stabilité à haut débit.
Sémantiques de livraison, accusé de réception, rétention et lettres mortes
- Accusé de réception et délais
- Livraison au moins une fois (at-least-once) : des doublons peuvent survenir.
- Chaque livraison a un délai d’accusé de réception (ack deadline) (10 secondes par défaut). Prolongez-le (ModifyAckDeadline) lors du traitement de tâches longues ; le non-respect de ce délai est la cause la plus fréquente de livraisons Push en double.
- Un Nack (accusé de réception négatif) ou l’expiration du délai rend le message éligible à une nouvelle livraison.
- Rétention
- Les messages non acquittés sont conservés pendant le délai d’ack de l’abonnement et réessayés ; les messages acquittés peuvent être conservés jusqu’à la durée de rétention des messages du sujet pour être relus. Configurez la rétention pour couvrir votre durée maximale d’interruption plus le temps de récupération.
- Nouvelles tentatives (Retries)
- Pull : la nouvelle livraison a lieu après l’expiration du délai d’ack ; contrôlez la concurrence avec les limites de contrôle de flux (flow control).
- Push : backoff exponentiel ; seul un code HTTP 2xx est considéré comme un succès. Les codes 3xx/4xx/5xx déclenchent de nouvelles tentatives. Implémentez des gestionnaires (handlers) idempotents pour tolérer les répétitions.
- Sujets de lettres mortes (DLT)
- Configurez un sujet DL et un nombre maximal de tentatives de livraison par abonnement pour mettre en quarantaine les messages empoisonnés (poison messages).
- Surveillez le volume de la DLQ ; créez des flux de travail de tri et republiez sur le sujet principal après correction.
Exemple :
undefined
Résumé des sémantiques de livraison :
- Pub/Sub : au moins une fois (at-least-once), ordonnancement au mieux (best-effort) au sein d’une clé d’ordonnancement si activé.
- Récepteurs (Sinks) : les API d’insertion de BigQuery fournissent une atténuation des doublons (insertId ou offsets de flux de la Storage Write API), mais concevez tout de même des consommateurs et des rédacteurs (writers) idempotents.
Schémas, compatibilité et validation
- Schémas Pub/Sub
- Prise en charge native d’Avro et de Protocol Buffers avec des schémas stockés de manière centralisée.
- Paramètres de schéma au niveau du sujet : encodage (Avro ou Protobuf) et application (aucune, valider uniquement ou exiger).
- Le producteur publie des charges utiles encodées ; Pub/Sub valide par rapport au schéma actuel lorsque l’application est activée.
- Évolution et compatibilité
- Utilisez des changements rétrocompatibles (ajouter des champs optionnels, ajouter des champs avec des valeurs par défaut dans Avro, ne jamais réutiliser les balises dans Protobuf, éviter de supprimer ou de renommer des champs).
- Versionnez explicitement les schémas. Pour les changements non rétrocompatibles (breaking changes), publiez en double sur des sujets v1 et v2, ou ajoutez un champ de version et routez en conséquence.
- Contrats producteur-consommateur
- Les consommateurs doivent ignorer les champs inconnus et utiliser des valeurs par défaut pour les champs manquants.
- Testez la compatibilité des schémas sur tous les consommateurs avant la promotion ; validez dans les abonnements de pré-production (staging) avec la même application de schéma que la production.
Exemple court d’Avro (extrait) :
undefined
Intégration événementielle, Eventarc et interopérabilité avec Kafka
- Eventarc et CloudEvents
- Eventarc achemine les événements provenant des services Google Cloud (et de sources personnalisées via Pub/Sub) vers Cloud Run, GKE ou Workflows en utilisant la spécification CloudEvents. Des attributs tels que le type, la source et le sujet permettent un filtrage fin et une auditabilité.
- Utilisez des filtres d’attributs pour minimiser le fan-out et réduire la charge en aval.
- La livraison est de type « au moins une fois » (at-least-once) ; rendez les gestionnaires (handlers) idempotents et sans état (stateless) lorsque c’est possible.
- Exemple de déclencheur Eventarc :
gcloud eventarc triggers create gcs-finalize-to-run
–destination-run-service=ingestor
–event-filters=“type=google.cloud.storage.object.v1.finalized”
–event-filters=“bucket=my-data-bucket”
–service-account=eventarc-sa@PROJECT_ID.iam.gserviceaccount.com - Interopérabilité avec Kafka et migration gérée
- Les modèles Dataflow connectent Kafka <-> Pub/Sub pour une migration par phases. Mettez en miroir les topics en conservant les clés ; basculez d’abord les consommateurs, puis les producteurs, ou effectuez une double écriture (dual-write) pendant la transition.
- Pub/Sub Lite offre un streaming partitionné avec une capacité provisionnée, un routage basé sur les clés et un coût inférieur ; il est régional/zonal et convient aux charges de travail de type Kafka où une capacité prévisible et l’ordre par partition sont les principales préoccupations.
- Considérations sur la migration :
- Ordonnancement : mappez les clés Kafka avec les clés d’ordonnancement (ordering keys) de Pub/Sub ou les partitions de Lite.
- Décalages (Offsets) : transportez les décalages en tant qu’attributs de message pour le diagnostic ; les consommateurs ne peuvent pas se fier aux décalages Kafka après la migration.
- Livraison : acceptez le mode « au moins une fois » (at-least-once) ; imposez l’idempotence en aval.
- Schémas : migrez les définitions de Confluent Schema Registry vers les schémas Pub/Sub ou standardisez sur Protobuf/Avro avec des règles d’évolution compatibles.
Modèles d’ingestion en streaming, débit, mise à l’échelle, sécurité et opérations
- Modèles d’ingestion en temps réel
- Pub/Sub -> Dataflow -> BigQuery : utilisez le récepteur (sink) BigQuery Storage Write API pour un débit élevé et l’idempotence avec les décalages de flux (stream offsets) ; acheminez les échecs vers une table de lettres mortes pour inspection.
- Pub/Sub -> Dataflow -> Cloud Storage : archivez les événements bruts pour le retraitement ; utilisez des écritures fenêtrées et compressées pour équilibrer le coût et la latence.
- Pub/Sub -> banques de données opérationnelles : écrivez dans Bigtable pour les recherches à faible latence, Spanner pour les transactions fortement cohérentes, ou Cloud SQL/Firestore en fonction des besoins de la charge de travail. Assurez des opérations d’insertion/mise à jour (upserts) idempotentes basées sur un ID d’événement unique.
- Livraison au moins une fois, prévention des doublons et idempotence
- Transportez un
event_idet unevent_timeuniques dans chaque message ; imposez des UUID côté producteur. - Déduplication en streaming BigQuery : définissez
insertIdou utilisez la Storage Write API avec des flux ordonnés ; protégez tout de même les requêtes avec une logique de déduplication. - Exemple de déduplication au moment de la requête : WITH ranked AS ( SELECT t.*, ROW_NUMBER() OVER (PARTITION BY event_id ORDER BY event_time DESC) AS rn FROM dataset.events t ) SELECT * EXCEPT(rn) FROM ranked WHERE rn = 1;
- Pour les points de terminaison push, ne retournez un code 2xx qu’après un traitement réussi ; sinon, attendez-vous à une nouvelle livraison.
- Transportez un
- Débit des messages, quotas et mise à l’échelle
- Producteurs : regroupez les messages par lots et réutilisez les connexions ; parallélisez sur plusieurs clients. Utilisez de nombreuses clés d’ordonnancement pour mettre à l’échelle les charges de travail ordonnées.
- Abonnés : préférez l’extraction en streaming (streaming pull) avec contrôle de flux (nombre maximal d’octets/messages en attente). Dimensionnez les délais de confirmation (ack) en fonction du temps de traitement et prolongez-les si nécessaire.
- Surveillez et demandez des augmentations de quota pour le débit de publication et d’abonnement à mesure que les volumes augmentent ; prévoyez une marge de manœuvre (par exemple, 2x le pic attendu) pour absorber les rafales.
- Cohérence et disponibilité
- Le streaming BigQuery est à cohérence à terme (eventually consistent) pour la visibilité des requêtes ; pour les requêtes interactives qui doivent inclure les lignes streamées, attendez en fonction de la latence observée (par exemple, 2x le délai de disponibilité P50) ou concevez en utilisant des agrégations alignées sur les watermarks dans Dataflow et interrogez les résultats matérialisés.
- Sécurité
- IAM : accordez les rôles du moindre privilège (
pubsub.publisheraux producteurs sur le sujet ;pubsub.subscriberaux consommateurs sur l’abonnement). Utilisez des comptes de service dédiés pour chaque charge de travail. - Authentification push : configurez les abonnements push pour joindre des jetons OIDC d’un compte de service ; imposez la validation de l’audience sur le point de terminaison. Préférez les points de terminaison privés Cloud Run pour l’authentification et le TLS intégrés.
- Chiffrement : Pub/Sub chiffre les données en transit et au repos ; utilisez CMEK sur les sujets pour les clés gérées par le client. Appliquez les VPC Service Controls pour réduire le risque d’exfiltration de données. Utilisez le chiffrement côté client pour les champs de charge utile sensibles si nécessaire.
- IAM : accordez les rôles du moindre privilège (
- Diagnostic opérationnel du retard, de la nouvelle livraison et de l’échec de l’abonné
- Surveiller avec Cloud Monitoring :
subscription/num_undelivered_messagesetoldest_unacked_message_agepour le backlog (messages en attente).expired_ack_deadline_countpour détecter les confirmations (acks) manquées qui provoquent des doublons.publish_request_countetpull_request_countpour le débit.
- Enquêtez sur les événements manquants dans le tableau de bord en rejouant un jeu de données connu à travers le pipeline et en comparant les sorties étape par étape pour isoler la transformation ou le récepteur défectueux.
- Pour le streaming Dataflow :
- Utilisez l’autoscaling avec un
maxWorkersapproprié pour absorber la charge provenant de nombreuses sources. - Drainez les pipelines pour les mises à jour incompatibles afin de permettre au travail en cours de se terminer et d’éviter la perte de données.
- Utilisez l’autoscaling avec un
- Pour les notifications d’insertion BigQuery, acheminez les entrées d’audit de Cloud Logging via un récepteur (sink) filtré sur des tables spécifiques vers un sujet Pub/Sub pour les alertes.
- Surveiller avec Cloud Monitoring :
Scénario de problème pratique
Contoso Freight a besoin d’une plateforme d’événements mondiale et en temps réel pour ingérer 10 000 messages de télémétrie IoT par minute provenant de camions, enrichir les événements, alimenter des analyses interactives et déclencher des workflows lors de dépôts de fichiers par des partenaires externes. Certains CSV des partenaires contiennent des lignes malformées, et l’équipe d’analyse doit pouvoir inspecter les erreurs sans bloquer le flux.
- Créer la couche principale de messagerie et de schéma
- Action : Définir un schéma Avro pour la télémétrie et l’attacher à un sujet Pub/Sub
telemetryavec l’application du schéma (schema enforcement) définie surrequire. Activer l’ordonnancement des messages et publier avecordering_key = hash(device_id). - Justification : L’application du schéma au niveau du sujet rejette tôt les événements malformés. L’ordonnancement par appareil permet un traitement ordonné si nécessaire, tandis que le hachage répartit les clés pour préserver le débit.
- Provisionner des abonnements avec isolation et gestion des lettres mortes
- Action : Créer un abonnement en mode pull
telemetry-stream-subpour Dataflow avec un sujet de lettres mortestelemetry-dltetmax_delivery_attempts=10. Ajouter un abonnement BigQuerytelemetry-raw-bqpour stocker les événements bruts dans une table partitionnée par temps pour la traçabilité (lineage) et la relecture. - Justification : La DLQ (Dead-Letter Queue) isole les messages empoisonnés pour investigation. Un abonnement BigQuery distinct fournit un chemin d’exportation à faible charge opérationnelle pour la conservation des événements bruts, indépendamment du pipeline de traitement.
- Construire un pipeline de streaming Dataflow pour l’enrichissement et les récepteurs
- Action : Ingérer depuis
telemetry-stream-suben utilisant l’extraction en streaming (streaming pull) avec contrôle de flux. Valider par rapport au schéma, enrichir avec des données de référence et calculer des agrégats fenêtrés. Écrire dans BigQuery en utilisant la Storage Write API avec un flux nommé etinsertId = event_id; écrire des sauvegardes brutes dans Cloud Storage toutes les heures ; rediriger les enregistrements incorrects/échoués vers une table de lettres mortes BigQuery. - Justification : La Storage Write API offre des écritures à haut débit et faible latence avec idempotence via
insertId/décalages de flux. Une table de lettres mortes permet l’inspection sans bloquer le flux, et les archives Cloud Storage permettent la relecture.
- Gérer les doublons et la cohérence à terme dans les analyses
- Action : Pour les requêtes interactives qui doivent exclure les doublons, publiez
event_idetevent_timedans chaque enregistrement et utilisez une vue de déduplication : CREATE OR REPLACE VIEW analytics.latest_events AS SELECT * EXCEPT(rn) FROM ( SELECT e.*, ROW_NUMBER() OVER (PARTITION BY event_id ORDER BY event_time DESC) rn FROM analytics.events e ) WHERE rn = 1; Introduire un bref délai dans les requêtes basé sur la disponibilité observée du streaming BigQuery (par exemple, deux fois la latence médiane). - Justification : La livraison au moins une fois nécessite des écritures idempotentes et une déduplication au moment de la requête. L’attente réduit les oublis de lignes en cours de traitement, compte tenu de la latence de visibilité du streaming.
- Intégrer les dépôts de fichiers des partenaires avec Eventarc
- Action : Configurer Eventarc pour acheminer les événements
object.finalizedde Cloud Storage pour le bucketpartner-dropsvers un service Cloud Run qui lance une tâche Dataflow en mode batch pour charger les CSV dans BigQuery, en envoyant les erreurs d’analyse (parse errors) à une table de lettres mortes. - Justification : Eventarc fournit une orchestration événementielle avec un filtrage CloudEvents sur le bucket et le préfixe de l’objet. Une tâche Dataflow en mode batch sépare les lignes malformées pour analyse tout en chargeant rapidement les données valides.
- Sécuriser la plateforme
- Action : Utiliser des comptes de service distincts : les producteurs obtiennent
pubsub.publishersurtelemetry; le SA du worker Dataflow obtientpubsub.subscribersurtelemetry-stream-subet un accès en écriture aux datasets BigQuery et à Cloud Storage cibles ; le déclencheur Eventarc utilise un SA dédié avec le rôleinvokersur Cloud Run. Activer CMEK sur le sujettelemetryet les datasets BigQuery. Configurer les points de terminaison push, le cas échéant, avec OIDC et des vérifications d’audience. - Justification : Le principe du moindre privilège IAM et CMEK répondent aux exigences de sécurité et de conformité ; la livraison authentifiée empêche l’usurpation d’identité (spoofing).
- Opérer et mettre à l’échelle de manière fiable
- Action : Configurer l’autoscaling de Dataflow avec un
maxWorkersgénéreux pour absorber les pics. Surveillersubscription/oldest_unacked_message_ageetexpired_ack_deadline_count; alerter lorsque les seuils sont dépassés. Pour les changements de pipeline qui rompent la compatibilité, déployer avec l’option de drainage (drain) pour éviter la perte de messages. Si le retard augmente, augmentez le parallélisme des abonnés et prolongez les délais de confirmation (ack) proportionnellement au temps de traitement. - Justification : La surveillance proactive détecte tôt le retard et les nouvelles livraisons. L’autoscaling et des délais de confirmation (ack) bien réglés préviennent les tempêtes de doublons. Le drainage préserve les messages en cours de traitement pendant les mises à niveau.
Cette conception offre une ingestion en temps réel résiliente, sécurisée et observable avec une intégration batch événementielle, prend en charge
← Traitement en flux continu avec Dataflow et Apache Beam · Tous les domaines · Spark →
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 →