Votre mission
MoveNow gère une flotte de véhicules. Les applications envoient leur position, mais les tableaux de suivi prennent du retard dès que le traitement s’interrompt. Les analystes veulent garder un historique exploitable, et l’équipe mobile annonce déjà plusieurs versions de son format de message. Votre groupe construit un pipeline de laboratoire et doit pouvoir expliquer où vont les événements quand tout ne se passe pas comme prévu.
« Si le traitement s’arrête, les positions doivent pouvoir être traitées plus tard. Et un message incorrect ne doit pas bloquer tout le flux. »
Ce dont le produit a besoin
- Recevoir des positions fictives et les retrouver dans un historique analytique.
- Découpler la publication du traitement, pour encaisser une interruption temporaire.
- Repérer les doublons, le retard et les messages invalides.
- Prévoir une procédure de correction et de rejeu, sans promettre une absence de perte que vous n’avez pas démontrée.
Le périmètre du laboratoire
Vous réalisez Pub/Sub, le transfert vers BigQuery, une table partitionnée, le traitement des échecs et des identités dédiées. L’application mobile, la carte, le dispatch et la géolocalisation réelle sont hors périmètre.
Le formateur vous donnera les volumes métier quand vous les lui demanderez. Ils servent à raisonner, pas à être reproduits : les essais de charge du laboratoire ont leurs propres limites, fixées à l’étape 5.
L’application sert seulement à faire passer des requêtes ou des messages. Votre travail porte sur Terraform, le réseau, IAM, le pipeline, le monitoring et la reprise. Personne ne vous demande une application métier complète, et elle ne sera pas évaluée. Le kit de démarrage est là pour vous éviter de l’écrire.
Ce qui compte en premier
Le parcours suit six étapes, et le temps passe vite. Le socle, c’est ce que vous devez pouvoir montrer quoi qu’il arrive. La version complète vient ensuite, si le temps le permet. Un socle solide et bien expliqué vaut mieux qu’une version complète vérifiée à moitié. À chaque étape, notez en quelques lignes ce que vous avez décidé et ce que vous avez mesuré.
| Étape | Socle | Version complète |
|---|---|---|
| 01 Cadrer | Une fiche d’hypothèses et trois critères de réussite validés par le formateur. | Les volumes obtenus, une limite de coût et ce que vous écartez, écrits noir sur blanc. |
| 02 Architecture | Le schéma cible expliqué flux par flux, avec votre découpage Terraform. | Trois choix d’infrastructure argumentés, avec les alternatives écartées. |
| 03 Construire | Le chemin nominal déployé par Terraform, avec un état distant et des identités dédiées : un event_id suivi de la publication jusqu’à BigQuery. | Des modules documentés, un plan sans dérive et un test négatif, c’est-à-dire un refus ou une erreur attendus. |
| 04 Automatiser | Une CI qui s’authentifie par fédération d’identité et lance fmt, validate et plan à chaque changement. | L’apply du plan relu après approbation, un test du parcours nominal et une destruction approuvée. |
| 05 Éprouver | Un dashboard, une alerte, un essai de charge borné et une analyse argumentée de la nouvelle contrainte. | Deux alertes dont la notification a bien été reçue, une expérience menée jusqu’au retour au service, et un essai borné sur la nouvelle contrainte. |
| 06 Présenter | 20 minutes sur ce qui fonctionne vraiment, un ordre de grandeur des coûts et le nettoyage des ressources. | Un calcul de coûts détaillé et les écarts avec une cible de production. |
Avant de commencer
Avant de créer la moindre ressource, faites le point avec le formateur :
- Le projet GCP du laboratoire. La facturation doit être active, les API et les quotas disponibles, et vous devez pouvoir créer les ressources et déléguer les permissions nécessaires.
- La région, le préfixe de votre groupe, les labels à poser et le plafond de dépense. Une alerte de budget prévient, elle ne coupe rien.
- Vos outils : Terraform, Git et gcloud. Si votre sujet utilise des conteneurs, il vous faut aussi de quoi construire une image et un registre pour la stocker.
- Un dépôt Git pour le groupe et une plateforme de CI. En local, vous travaillez avec les identifiants prévus par le cours. En CI, vous préparerez une identité fédérée aux droits limités.
- Un endroit où ranger vos preuves : schémas, captures, résultats horodatés et décisions, sans aucun secret. L’un exécute, l’autre relit, le troisième mesure, et vous changez de rôle régulièrement.
Votre jeu de démonstration
Un producteur cadencé et dix événements JSON synthétiques avec event_id, vehicle_id, event_time, latitude et longitude. Documentez les types et les formats. Les positions sont inventées, et l’event_id doit permettre de suivre chaque événement. Le kit ci-dessous fournit ce producteur.
Vous savez dans quel projet vous travaillez, ce que vous avez le droit de créer et comment tout arrêter puis supprimer. S’il vous manque un accès, demandez-le plutôt que de le contourner avec des permissions trop larges.
Kit de démarrage
Ce kit vous évite d’écrire l’application. Il contient le code à déployer, avec une interface qui montre ce que fait votre infrastructure, et un docker compose qui imite les services GCP pour tout faire tourner sur votre machine avant le premier déploiement. Il ne contient pas de Terraform : l’écrire reste le cœur du TP.
Vous pouvez le modifier, mais il ne sera pas évalué. Lisez-le avant de l’utiliser, car vous devrez pouvoir l’expliquer comme n’importe quel code que vous n’avez pas écrit vous-même.
Un producteur de positions fictives et un tableau de bord : la carte des véhicules, le débit écrit dans la table, la dead-letter et la réconciliation des lots par identifiant. En local, le flux passe par le vrai émulateur Pub/Sub de Google, et le tableau imite la BigQuery subscription et la dead-letter policy. Un bouton coupe le transfert pour observer l’accumulation, puis le rattrapage. Sans Docker, un mode démo fait tout tourner en mémoire.
Télécharger le kit, archive zipVoir la démo en ligne
| Rôle | En local | Sur GCP |
|---|---|---|
| Topic et subscriptions | émulateur Pub/Sub | Pub/Sub |
| Transfert vers la table | imité par le tableau de bord | BigQuery subscription |
| Table analytique | mémoire du tableau de bord | table BigQuery partitionnée |
| Dead-letter | imitée, cinq essais puis topic d’échec | dead-letter policy de la subscription |
Pour l’essayer : décompressez l’archive, puis lancez docker compose up --build dans le dossier movenow et ouvrez http://localhost:8080. Le README de l’archive explique comment passer sur GCP.
docker-compose.yml, le lancement local
# MoveNow en local : docker compose up --build, puis http://localhost:8080
name: movenow-local
services:
# L'émulateur Pub/Sub officiel de Google. Il remplace le vrai service, sans IAM ni BigQuery.
pubsub:
image: gcr.io/google.com/cloudsdktool/google-cloud-cli:emulators
command: ["gcloud", "beta", "emulators", "pubsub", "start", "--project=movenow-local", "--host-port=0.0.0.0:8085"]
healthcheck:
test: ["CMD", "bash", "-c", "echo > /dev/tcp/127.0.0.1/8085"]
interval: 3s
retries: 40
# Tableau de bord. En local, il imite aussi la BigQuery subscription et la dead-letter policy.
tableau:
build: ./tableau
ports:
- "8080:8080"
environment:
PUBSUB_EMULATOR_HOST: pubsub:8085
PUBSUB_PROJECT_ID: movenow-local
LOTS_DIR: /lots
volumes:
- lots:/lots
depends_on:
pubsub:
condition: service_healthy
# Le producteur publie en continu des lots d'une minute, avec deux messages invalides par lot.
producteur:
build: ./producteur
command: ["--project", "movenow-local", "--topic", "positions", "--rate", "5", "--duration", "60",
"--invalid", "2", "--batch", "lot-local", "--out", "/lots", "--loop"]
environment:
PUBSUB_EMULATOR_HOST: pubsub:8085
volumes:
- lots:/lots
depends_on:
tableau:
condition: service_healthy
volumes:
lots:
producteur/producteur.mjs, le producteur
// Producteur de positions fictives pour MoveNow.
// Publie des lots cadencés sur un topic Pub/Sub et écrit, pour chaque lot, la liste des événements
// attendus dans <dossier>/<lot>.jsonl : c'est la base de la réconciliation avec BigQuery.
//
// node producteur.mjs --project mon-projet --topic positions --rate 10 --duration 60 --batch lot-01
// node producteur.mjs --topic positions --rate 5 --duration 10 --batch lot-02 --invalid 3
// node producteur.mjs --rate 2 --duration 3 --dry-run essai sans GCP
//
// Sur GCP, l'authentification passe par ADC. Si PUBSUB_EMULATOR_HOST est défini, la bibliothèque
// Pub/Sub publie vers l'émulateur local à la place.
import fs from 'node:fs';
import path from 'node:path';
import { parseArgs } from 'node:util';
const { values: options } = parseArgs({
options: {
topic: { type: 'string' },
project: { type: 'string' },
rate: { type: 'string', default: '10' },
duration: { type: 'string', default: '60' },
batch: { type: 'string', default: `lot-${new Date().toISOString().slice(0, 19).replace(/[-:T]/g, '')}` },
invalid: { type: 'string', default: '0' },
vehicles: { type: 'string', default: '20' },
out: { type: 'string', default: '.' },
loop: { type: 'boolean', default: false },
'dry-run': { type: 'boolean', default: false },
},
});
const rate = Number(options.rate);
const duration = Number(options.duration);
const invalid = Number(options.invalid);
const vehicles = Number(options.vehicles);
for (const [name, value, max] of [['rate', rate, 100], ['duration', duration, 600], ['invalid', invalid, 1000], ['vehicles', vehicles, 1000]]) {
// Plafonds volontaires : ce script sert à des essais bornés, pas à reproduire le trafic métier.
if (!Number.isInteger(value) || value < 0 || value > max) throw new Error(`--${name} doit être un entier entre 0 et ${max}`);
}
if (!/^[a-zA-Z0-9_-]+$/.test(options.batch)) throw new Error('--batch ne doit contenir que des lettres, chiffres, - et _');
if (!options['dry-run'] && !options.topic) throw new Error('--topic est requis, sauf avec --dry-run');
fs.mkdirSync(options.out, { recursive: true });
let topic = null;
if (!options['dry-run']) {
const { PubSub } = await import('@google-cloud/pubsub');
const projectId = options.project || process.env.GOOGLE_CLOUD_PROJECT || process.env.PUBSUB_PROJECT_ID;
topic = new PubSub(projectId ? { projectId } : {}).topic(options.topic);
}
// Compteur global, écrit chaque seconde dans progression.json : en local, le tableau de bord s'en sert
// pour estimer le nombre de messages en attente.
let publiesTotal = 0;
const progression = (batch) => fs.writeFileSync(path.join(options.out, 'progression.json'), JSON.stringify({ lot: batch, publies: publiesTotal, maj: new Date().toISOString() }));
// Chaque véhicule part d'un point au hasard autour de Saint-Quentin, puis se déplace un peu à chaque position.
const flotte = Array.from({ length: vehicles }, () => ({ lat: 49.82 + Math.random() * 0.06, lon: 3.25 + Math.random() * 0.08 }));
function deplacer(n) {
const v = flotte[n];
v.lat = Math.min(49.88, Math.max(49.82, v.lat + (Math.random() - 0.5) * 0.002));
v.lon = Math.min(3.33, Math.max(3.25, v.lon + (Math.random() - 0.5) * 0.003));
return v;
}
async function publierLot(batch) {
const total = rate * duration;
// Les événements invalides sont répartis régulièrement dans le lot.
const invalidCount = Math.min(invalid, total);
const invalidAt = new Set(Array.from({ length: invalidCount }, (_, i) => 1 + Math.floor(((i + 0.5) * total) / invalidCount)));
const records = [];
const pending = [];
const start = Date.now();
let seq = 0;
for (let second = 0; second < duration; second++) {
for (let i = 0; i < rate; i++) {
seq += 1;
const n = seq % vehicles;
const position = deplacer(n);
const event = {
event_id: `${batch}-${String(seq).padStart(6, '0')}`,
vehicle_id: `vh-${String(n).padStart(4, '0')}`,
event_time: new Date().toISOString(),
latitude: Number(position.lat.toFixed(5)),
longitude: Number(position.lon.toFixed(5)),
};
// Un message « invalide » est accepté par le topic mais ne respecte pas le schéma de la table :
// sa latitude devient une chaîne de caractères.
if (invalidAt.has(seq)) event.latitude = 'nord';
const record = { event_id: event.event_id, invalide: invalidAt.has(seq), publie_le: event.event_time };
records.push(record);
publiesTotal += 1;
if (topic) {
pending.push(topic.publishMessage({ data: Buffer.from(JSON.stringify(event)) })
.then((messageId) => { record.message_id = messageId; })
.catch((error) => { record.erreur = error.message; }));
} else if (seq <= 3) {
console.log(JSON.stringify(event));
}
}
if (topic) progression(batch);
const wait = start + (second + 1) * 1000 - Date.now();
if (wait > 0) await new Promise((resolve) => setTimeout(resolve, wait));
}
await Promise.all(pending);
const file = path.join(options.out, `${batch}.jsonl`);
fs.writeFileSync(file, records.map((record) => JSON.stringify(record)).join('\n') + '\n');
const failed = records.filter((record) => record.erreur).length;
console.log(`${batch} : ${records.length} événements, ${invalidAt.size} invalides, ${failed} échecs de publication${topic ? '' : ' (essai à blanc, rien publié)'}. Liste attendue : ${file}`);
return failed;
}
if (options.loop) {
// Mode continu pour la démonstration locale : un nouveau lot numéroté à la suite du précédent.
for (let n = 1; ; n++) await publierLot(`${options.batch}-${String(n).padStart(3, '0')}`);
} else {
process.exitCode = (await publierLot(options.batch)) ? 1 : 0;
}
tableau/sources.mjs, les sources du tableau de bord
// Sources de données du tableau de bord MoveNow. Trois modes :
// demo aucune variable : un flux simulé en mémoire, sans GCP ni Docker
// emulateur PUBSUB_EMULATOR_HOST défini : l'émulateur Pub/Sub local, avec une imitation
// de la BigQuery subscription et de la dead-letter policy
// bigquery SOURCE=bigquery : lecture de la vraie table BigQuery sur GCP
import fs from 'node:fs';
import path from 'node:path';
import { bloc, blocCloudRun, blocIdentite, enCache } from './gcp.mjs';
const MAX_TENTATIVES = Number(process.env.MAX_DELIVERY_ATTEMPTS || 5);
const log = (severity, message, extra = {}) => console.log(JSON.stringify({ severity, message, ...extra }));
// Le contrat de la table : les mêmes contrôles que BigQuery quand la subscription écrit avec use_table_schema.
export function valider(data) {
let event;
try {
event = JSON.parse(data);
} catch {
return { erreur: 'JSON illisible' };
}
if (!event || typeof event !== 'object') return { erreur: 'message vide' };
if (typeof event.event_id !== 'string' || !event.event_id) return { erreur: 'event_id manquant' };
if (typeof event.vehicle_id !== 'string' || !event.vehicle_id) return { erreur: 'vehicle_id manquant', event };
if (typeof event.event_time !== 'string' || Number.isNaN(Date.parse(event.event_time))) return { erreur: 'event_time n’est pas un TIMESTAMP', event };
for (const champ of ['latitude', 'longitude']) {
if (typeof event[champ] !== 'number' || !Number.isFinite(event[champ])) return { erreur: `${champ} n’est pas un FLOAT`, event };
}
return { event };
}
const lireJson = (texte) => { try { return JSON.parse(texte); } catch { return null; } };
// Une « table » en mémoire, partagée par les modes demo et emulateur.
class TableLocale {
constructor() {
this.lignes = [];
this.parId = new Map();
this.derniers = new Map();
this.deadLetter = [];
this.idsDeadLetter = new Set();
this.recus = 0;
this.debit = [];
this.depuisDernierPoint = 0;
setInterval(() => {
this.debit.push({ t: Date.now(), n: this.depuisDernierPoint });
this.depuisDernierPoint = 0;
if (this.debit.length > 60) this.debit.shift();
}, 2000).unref();
}
ecrire(event) {
const ligne = { ...event, ingere_le: new Date().toISOString() };
this.lignes.push(ligne);
if (this.lignes.length > 20000) this.lignes.splice(0, this.lignes.length - 20000);
this.parId.set(event.event_id, (this.parId.get(event.event_id) || 0) + 1);
this.derniers.set(event.vehicle_id, ligne);
this.recus += 1;
this.depuisDernierPoint += 1;
}
mettreDeCote(entree) {
this.deadLetter.unshift(entree);
if (this.deadLetter.length > 200) this.deadLetter.pop();
if (entree.event_id) this.idsDeadLetter.add(entree.event_id);
}
positions() {
return [...this.derniers.values()].map(({ vehicle_id, latitude, longitude, event_time }) => ({ vehicle_id, latitude, longitude, event_time }));
}
evenements() {
return this.lignes.slice(-25).reverse();
}
lot(nom) {
const prefixe = `${nom}-`;
const recus = {};
for (const [id, n] of this.parId) if (id.startsWith(prefixe)) recus[id] = n;
return { recus, dead_letter: [...this.idsDeadLetter].filter((id) => id.startsWith(prefixe)) };
}
nomsDeLots() {
const noms = new Set();
for (const id of this.parId.keys()) noms.add(id.replace(/-\d{6}$/, ''));
return [...noms].reverse().slice(0, 20);
}
}
// Mode demo : un producteur et une imitation du transfert, tout en mémoire.
function sourceDemo() {
const table = new TableLocale();
const file = [];
const lots = new Map();
let actif = true;
let publies = 0;
let numeroLot = 0;
let lot = null;
const flotte = Array.from({ length: 20 }, () => ({ lat: 49.82 + Math.random() * 0.06, lon: 3.25 + Math.random() * 0.08 }));
// Cinq événements par seconde, par lots d'une minute, avec un message invalide par lot.
setInterval(() => {
if (!lot || lot.seq >= 300) {
numeroLot += 1;
lot = { nom: `lot-demo-${String(numeroLot).padStart(3, '0')}`, seq: 0, attendus: [], invalide: 1 + Math.floor(Math.random() * 300) };
lots.set(lot.nom, lot.attendus);
if (lots.size > 10) lots.delete(lots.keys().next().value);
}
lot.seq += 1;
const n = lot.seq % flotte.length;
const v = flotte[n];
v.lat = Math.min(49.88, Math.max(49.82, v.lat + (Math.random() - 0.5) * 0.002));
v.lon = Math.min(3.33, Math.max(3.25, v.lon + (Math.random() - 0.5) * 0.003));
const event = {
event_id: `${lot.nom}-${String(lot.seq).padStart(6, '0')}`,
vehicle_id: `vh-${String(n).padStart(4, '0')}`,
event_time: new Date().toISOString(),
latitude: Number(v.lat.toFixed(5)),
longitude: Number(v.lon.toFixed(5)),
};
const invalide = lot.seq === lot.invalide;
if (invalide) event.latitude = 'nord';
lot.attendus.push({ event_id: event.event_id, invalide });
file.push({ data: JSON.stringify(event), tentatives: 0 });
publies += 1;
}, 200).unref();
// Le transfert traite au plus vingt messages toutes les 200 ms : le rattrapage reste visible.
setInterval(() => {
if (!actif) return;
for (let i = 0; i < 20 && file.length; i++) {
const message = file.shift();
const resultat = valider(message.data);
if (!resultat.erreur) { table.ecrire(resultat.event); continue; }
message.tentatives += 1;
if (message.tentatives < MAX_TENTATIVES) { file.push(message); continue; }
table.mettreDeCote({ event_id: resultat.event?.event_id || null, erreur: resultat.erreur, tentatives: message.tentatives, recu_le: new Date().toISOString(), extrait: message.data.slice(0, 160) });
}
}, 200).unref();
return {
mode: 'demo',
libelle: 'Démo en mémoire',
pret: Promise.resolve(),
async etat() {
return { compteurs: { recus: table.recus, dead_letter: table.deadLetter.length, en_attente: file.length, publies }, debit: table.debit, pas_secondes: 2, transfert: { pilotable: true, actif } };
},
async positions() { return table.positions(); },
async evenements() { return table.evenements(); },
async deadLetter() { return table.deadLetter.slice(0, 30); },
async lots() { return [...lots.keys()].reverse().map((nom) => ({ nom, attendus: lots.get(nom).length })); },
async lot(nom) { return { nom, attendus: lots.get(nom) || null, ...table.lot(nom) }; },
async piloter(etat) { actif = etat; return { actif }; },
async deploiement() {
return { plateforme: 'local', blocs: [
bloc('producteur', 'Producteur', 'simule', 'producteur intégré au tableau, cinq positions par seconde'),
bloc('topic', 'Topic Pub/Sub', 'simule', 'file en mémoire'),
bloc('subscription', 'Transfert vers la table', 'simule', actif ? 'transfert imité, actif' : 'transfert imité, coupé'),
bloc('deadletter', 'Dead-letter', 'simule', `${table.deadLetter.length} message(s) mis de côté après ${MAX_TENTATIVES} essais`),
bloc('inspection', 'Subscription d’inspection', 'simule', 'liste affichée dans le tableau'),
bloc('table', 'Table BigQuery', 'simule', `${table.recus} positions en mémoire`),
blocCloudRun('Tableau de bord sur Cloud Run'),
bloc('identite', 'Identité du tableau', 'simule', 'pas de compte de service en local'),
bloc('obs', 'Monitoring et alertes', 'manuel', 'backlog, âge du plus ancien message et alertes : à montrer dans Monitoring'),
] };
},
};
}
// Mode emulateur : vrai Pub/Sub local. L'émulateur ne sait pas écrire dans BigQuery ni appliquer
// une dead-letter policy : ce module les imite, pour que le flux se comporte comme sur GCP.
async function sourceEmulateur(env) {
const { PubSub } = await import('@google-cloud/pubsub');
const pubsub = new PubSub({ projectId: env.PUBSUB_PROJECT_ID || 'movenow-local' });
const noms = {
topic: env.TOPIC || 'positions',
subscription: env.SUBSCRIPTION || 'positions-vers-bigquery',
deadLetterTopic: env.DEAD_LETTER_TOPIC || 'positions-dead-letter',
inspection: env.DEAD_LETTER_SUBSCRIPTION || 'positions-dead-letter-inspection',
};
const dossierLots = env.LOTS_DIR || '/lots';
const table = new TableLocale();
const tentatives = new Map();
let abonnement = null;
// Crée les ressources si elles n'existent pas encore (code 6 : ALREADY_EXISTS).
const assurer = async (creer) => { try { await creer(); } catch (error) { if (error.code !== 6) throw error; } };
async function initialiser() {
await assurer(() => pubsub.createTopic(noms.topic));
await assurer(() => pubsub.createTopic(noms.deadLetterTopic));
await assurer(() => pubsub.topic(noms.topic).createSubscription(noms.subscription, { ackDeadlineSeconds: 10 }));
await assurer(() => pubsub.topic(noms.deadLetterTopic).createSubscription(noms.inspection));
ouvrir();
pubsub.subscription(noms.inspection).on('message', (message) => {
const event = lireJson(message.data.toString());
table.mettreDeCote({
event_id: event?.event_id || null,
erreur: message.attributes.erreur || 'inconnue',
tentatives: Number(message.attributes.tentatives) || null,
recu_le: new Date().toISOString(),
extrait: message.data.toString().slice(0, 160),
});
message.ack();
});
log('INFO', 'topics et subscriptions prêts dans l’émulateur', noms);
}
// L'émulateur peut mettre quelques secondes à répondre : on réessaie pendant une minute.
const pret = (async () => {
for (let essai = 1; ; essai++) {
try {
return await initialiser();
} catch (error) {
if (essai >= 30) throw error;
log('WARNING', `émulateur pas encore prêt : ${error.message}`);
await new Promise((resolve) => setTimeout(resolve, 2000));
}
}
})();
function ouvrir() {
const deadLetterTopic = pubsub.topic(noms.deadLetterTopic);
abonnement = pubsub.subscription(noms.subscription, { flowControl: { maxMessages: 500 } });
abonnement.on('message', async (message) => {
const resultat = valider(message.data.toString());
if (!resultat.erreur) {
table.ecrire(resultat.event);
tentatives.delete(message.id);
message.ack();
return;
}
const n = (tentatives.get(message.id) || 0) + 1;
tentatives.set(message.id, n);
if (n < MAX_TENTATIVES) { message.nack(); return; }
tentatives.delete(message.id);
await deadLetterTopic.publishMessage({ data: message.data, attributes: { erreur: resultat.erreur, tentatives: String(n), source: noms.subscription } });
message.ack();
});
abonnement.on('error', (error) => log('ERROR', error.message, { code: error.code }));
}
function lireAttendus(nom) {
const fichier = path.join(dossierLots, `${nom}.jsonl`);
if (!/^[\w-]+$/.test(nom) || !fs.existsSync(fichier)) return null;
return fs.readFileSync(fichier, 'utf8').trim().split('\n').map(lireJson).filter(Boolean).map(({ event_id, invalide }) => ({ event_id, invalide }));
}
return {
mode: 'emulateur',
libelle: 'Émulateur Pub/Sub local',
pret,
async etat() {
const progression = fs.existsSync(path.join(dossierLots, 'progression.json')) ? lireJson(fs.readFileSync(path.join(dossierLots, 'progression.json'), 'utf8')) : null;
const publies = progression?.publies ?? null;
const enAttente = publies === null ? null : Math.max(0, publies - table.recus - table.deadLetter.length);
return { compteurs: { recus: table.recus, dead_letter: table.deadLetter.length, en_attente: enAttente, publies }, debit: table.debit, pas_secondes: 2, transfert: { pilotable: true, actif: Boolean(abonnement) } };
},
async positions() { return table.positions(); },
async evenements() { return table.evenements(); },
async deadLetter() { return table.deadLetter.slice(0, 30); },
async lots() {
const fichiers = fs.existsSync(dossierLots) ? fs.readdirSync(dossierLots).filter((f) => f.endsWith('.jsonl')).map((f) => f.slice(0, -6)) : [];
const noms = [...new Set([...table.nomsDeLots(), ...fichiers])].sort().reverse().slice(0, 20);
return noms.map((nom) => ({ nom, attendus: lireAttendus(nom)?.length ?? null }));
},
async lot(nom) { return { nom, attendus: lireAttendus(nom), ...table.lot(nom) }; },
async deploiement() {
const progression = fs.existsSync(path.join(dossierLots, 'progression.json')) ? lireJson(fs.readFileSync(path.join(dossierLots, 'progression.json'), 'utf8')) : null;
const recent = progression && Date.now() - Date.parse(progression.maj) < 10000;
return { plateforme: 'local', blocs: [
bloc('producteur', 'Producteur', 'simule', recent ? `conteneur producteur actif, lot ${progression.lot}` : 'aucune publication récente du conteneur producteur'),
bloc('topic', 'Topic Pub/Sub', 'simule', `émulateur Pub/Sub, topic ${noms.topic}`),
bloc('subscription', 'Transfert vers la table', 'simule', `subscription ${noms.subscription}, transfert imité ${abonnement ? 'actif' : 'coupé'}`),
bloc('deadletter', 'Dead-letter', 'simule', `topic ${noms.deadLetterTopic}, ${MAX_TENTATIVES} essais imités`),
bloc('inspection', 'Subscription d’inspection', 'simule', `${noms.inspection}, ${table.deadLetter.length} message(s) reçu(s)`),
bloc('table', 'Table BigQuery', 'simule', `${table.recus} positions en mémoire`),
blocCloudRun('Tableau de bord sur Cloud Run'),
bloc('identite', 'Identité du tableau', 'simule', 'pas de compte de service en local'),
bloc('obs', 'Monitoring et alertes', 'manuel', 'backlog, âge du plus ancien message et alertes : à montrer dans Monitoring'),
] };
},
// Couper le transfert imite le retrait du droit d'écriture : les messages s'accumulent dans Pub/Sub.
async piloter(actif) {
if (actif && !abonnement) ouvrir();
if (!actif && abonnement) { await abonnement.close(); abonnement = null; }
return { actif: Boolean(abonnement) };
},
};
}
// Mode bigquery : lecture seule de la table réelle, rafraîchie toutes les REFRESH_SECONDS secondes.
// Chaque rafraîchissement lance quelques petites requêtes filtrées sur la partition.
async function sourceBigQuery(env) {
const { BigQuery } = await import('@google-cloud/bigquery');
if (!/^[\w-]+\.[\w-]+\.[\w-]+$/.test(env.BQ_TABLE || '')) throw new Error('BQ_TABLE doit être de la forme projet.dataset.table');
const bigquery = new BigQuery(env.GOOGLE_CLOUD_PROJECT ? { projectId: env.GOOGLE_CLOUD_PROJECT } : {});
const table = `\`${env.BQ_TABLE}\``;
const requete = async (query, params = {}) => (await bigquery.query({ query, params, location: env.BQ_LOCATION || undefined }))[0];
const valeur = (v) => (v && typeof v === 'object' && 'value' in v ? v.value : v);
const cache = { compteurs: { recus: 0, dead_letter: null, en_attente: null, publies: null }, debit: [], positions: [], evenements: [], deadLetter: [], erreur: null };
let lireDeadLetter = async () => [];
if (env.DEAD_LETTER_SUBSCRIPTION) {
// Lecture sans acquittement : les messages sont aussitôt rendus à la subscription d'inspection.
const { v1 } = await import('@google-cloud/pubsub');
const client = new v1.SubscriberClient();
const subscription = env.DEAD_LETTER_SUBSCRIPTION.includes('/') ? env.DEAD_LETTER_SUBSCRIPTION : `projects/${env.GOOGLE_CLOUD_PROJECT}/subscriptions/${env.DEAD_LETTER_SUBSCRIPTION}`;
lireDeadLetter = async () => {
const [reponse] = await client.pull({ subscription, maxMessages: 20 });
const messages = reponse.receivedMessages || [];
if (messages.length) await client.modifyAckDeadline({ subscription, ackIds: messages.map((m) => m.ackId), ackDeadlineSeconds: 0 });
return messages.map(({ message }) => {
const data = Buffer.from(message.data || '').toString();
return {
event_id: lireJson(data)?.event_id || null,
erreur: 'voir les journaux de la subscription',
tentatives: Number(message.attributes?.CloudPubSubDeadLetterSourceDeliveryCount) || null,
recu_le: message.publishTime ? new Date(Number(message.publishTime.seconds) * 1000).toISOString() : null,
extrait: data.slice(0, 160),
};
});
};
}
async function rafraichir() {
try {
const [compte] = await requete(`SELECT COUNT(*) AS recus FROM ${table} WHERE event_time > TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL 1 HOUR)`);
const debit = await requete(`SELECT UNIX_SECONDS(TIMESTAMP_TRUNC(event_time, MINUTE)) AS t, COUNT(*) AS n FROM ${table}
WHERE event_time > TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL 30 MINUTE) GROUP BY t ORDER BY t`);
const positions = await requete(`SELECT vehicle_id, latitude, longitude, event_time FROM ${table}
WHERE event_time > TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL 10 MINUTE)
QUALIFY ROW_NUMBER() OVER (PARTITION BY vehicle_id ORDER BY event_time DESC) = 1`);
const evenements = await requete(`SELECT event_id, vehicle_id, event_time, latitude, longitude FROM ${table}
WHERE event_time > TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL 1 HOUR) ORDER BY event_time DESC LIMIT 25`);
try {
cache.deadLetter = await lireDeadLetter();
cache.erreurDeadLetter = null;
} catch (error) {
cache.erreurDeadLetter = error.message;
}
cache.compteurs = { recus: compte.recus, dead_letter: env.DEAD_LETTER_SUBSCRIPTION ? cache.deadLetter.length : null, en_attente: null, publies: null };
cache.debit = debit.map(({ t, n }) => ({ t: t * 1000, n }));
cache.positions = positions.map((p) => ({ ...p, event_time: valeur(p.event_time) }));
cache.evenements = evenements.map((e) => ({ ...e, event_time: valeur(e.event_time) }));
cache.erreur = null;
} catch (error) {
cache.erreur = error.message;
log('ERROR', error.message);
}
}
const pret = rafraichir();
setInterval(rafraichir, Number(env.REFRESH_SECONDS || 20) * 1000).unref();
return {
mode: 'bigquery',
libelle: `BigQuery, ${env.BQ_TABLE}`,
pret,
async etat() { return { compteurs: cache.compteurs, debit: cache.debit, transfert: { pilotable: false, actif: null }, erreur: cache.erreur, pas_secondes: 60 }; },
async positions() { return cache.positions; },
async evenements() { return cache.evenements; },
async deadLetter() { return cache.deadLetter; },
async lots() {
const lignes = await requete(`SELECT REGEXP_EXTRACT(event_id, r'^(.*)-[0-9]{6}$') AS nom, COUNT(*) AS n FROM ${table}
WHERE event_time > TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL 1 DAY) GROUP BY nom ORDER BY MAX(event_time) DESC LIMIT 20`);
return lignes.filter((l) => l.nom).map((l) => ({ nom: l.nom, attendus: null }));
},
async lot(nom) {
const lignes = await requete(`SELECT event_id, COUNT(*) AS n FROM ${table}
WHERE STARTS_WITH(event_id, @prefixe) AND event_time > TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL 3 DAY) GROUP BY event_id`, { prefixe: `${nom}-` });
return { nom, attendus: null, recus: Object.fromEntries(lignes.map((l) => [l.event_id, l.n])), dead_letter: cache.deadLetter.map((d) => d.event_id).filter((id) => id?.startsWith(`${nom}-`)) };
},
async piloter() { throw Object.assign(new Error('sur GCP, coupez le transfert en retirant le droit d’écriture du service agent'), { http: 409 }); },
deploiement: enCache(() => carteBigQuery(env, bigquery, cache), 20),
};
}
// Carte du déploiement sur GCP. La table se lit avec bigquery.dataViewer, déjà nécessaire au tableau.
// La configuration de la subscription demande roles/pubsub.viewer sur elle, ce qui reste facultatif.
const contrat = { event_id: 'STRING', vehicle_id: 'STRING', event_time: 'TIMESTAMP', latitude: 'FLOAT', longitude: 'FLOAT' };
async function carteBigQuery(env, bigquery, cache) {
const blocs = [blocCloudRun('Tableau de bord sur Cloud Run'), await blocIdentite()];
const recus = cache.compteurs.recus;
blocs.push(recus > 0
? bloc('producteur', 'Producteur', 'ok', `${recus} position(s) écrite(s) dans la dernière heure`, 'requête sur la table')
: bloc('producteur', 'Producteur', 'inconnu', 'aucune position récente dans la table : producteur lancé, transfert configuré ?'));
const [projet, dataset, nomTable] = env.BQ_TABLE.split('.');
try {
const [meta] = await bigquery.dataset(dataset, { projectId: projet }).table(nomTable).getMetadata();
const champs = Object.fromEntries((meta.schema?.fields || []).map((f) => [f.name, f.type === 'FLOAT64' ? 'FLOAT' : f.type]));
const ecarts = Object.entries(contrat).filter(([nom, type]) => champs[nom] !== type).map(([nom, type]) => `${nom} : ${champs[nom] || 'absent'} au lieu de ${type}`);
const partition = meta.timePartitioning?.field ? `partitionnée sur ${meta.timePartitioning.field}` : meta.timePartitioning ? 'partitionnée par date d’ingestion' : 'non partitionnée';
const filtre = meta.requirePartitionFilter ? ', filtre de partition exigé' : '';
blocs.push(ecarts.length
? bloc('table', 'Table BigQuery', 'alerte', `schéma différent du contrat du producteur : ${ecarts.join(', ')}`, env.BQ_TABLE)
: bloc('table', 'Table BigQuery', 'ok', `${partition}${filtre}, schéma conforme au contrat du producteur`, env.BQ_TABLE));
} catch (error) {
blocs.push(bloc('table', 'Table BigQuery', 'echec', 'métadonnées de la table illisibles', error.message));
}
if (!env.SUBSCRIPTION) {
blocs.push(bloc('topic', 'Topic Pub/Sub', 'inconnu', 'renseignez SUBSCRIPTION pour afficher ce bloc'), bloc('subscription', 'BigQuery subscription', 'inconnu', 'renseignez SUBSCRIPTION pour afficher ce bloc'), bloc('deadletter', 'Dead-letter', 'inconnu', 'renseignez SUBSCRIPTION pour afficher ce bloc'));
} else {
try {
const { PubSub } = await import('@google-cloud/pubsub');
const [sub] = await new PubSub(env.GOOGLE_CLOUD_PROJECT ? { projectId: env.GOOGLE_CLOUD_PROJECT } : {}).subscription(env.SUBSCRIPTION).getMetadata();
const court = (nom = '') => nom.split('/').pop();
blocs.push(bloc('topic', 'Topic Pub/Sub', 'ok', `topic ${court(sub.topic)}`, sub.topic));
blocs.push(sub.bigqueryConfig?.table
? bloc('subscription', 'BigQuery subscription', 'ok', `écrit dans ${sub.bigqueryConfig.table}, use_table_schema ${sub.bigqueryConfig.useTableSchema ? 'oui' : 'non'}`, court(sub.name))
: bloc('subscription', 'Transfert vers la table', 'ok', 'subscription sans export BigQuery : un consommateur fait le transfert', court(sub.name)));
blocs.push(sub.deadLetterPolicy?.deadLetterTopic
? bloc('deadletter', 'Dead-letter', 'ok', `vers ${court(sub.deadLetterPolicy.deadLetterTopic)}, ${sub.deadLetterPolicy.maxDeliveryAttempts} tentatives au plus`, 'deadLetterPolicy')
: bloc('deadletter', 'Dead-letter', 'inconnu', 'pas de dead-letter policy sur la subscription'));
} catch (error) {
const droit = error.code === 7 ? 'pour afficher ce bloc, donnez au tableau roles/pubsub.viewer sur la subscription (facultatif)' : 'configuration de la subscription illisible';
blocs.push(bloc('topic', 'Topic Pub/Sub', 'inconnu', droit), bloc('subscription', 'BigQuery subscription', 'inconnu', droit, error.message), bloc('deadletter', 'Dead-letter', 'inconnu', droit));
}
}
blocs.push(!env.DEAD_LETTER_SUBSCRIPTION
? bloc('inspection', 'Subscription d’inspection', 'inconnu', 'renseignez DEAD_LETTER_SUBSCRIPTION pour afficher ce bloc')
: cache.erreurDeadLetter
? bloc('inspection', 'Subscription d’inspection', 'inconnu', 'lecture impossible : roles/pubsub.subscriber sur la subscription d’inspection ?', cache.erreurDeadLetter)
: bloc('inspection', 'Subscription d’inspection', 'ok', `lisible, ${cache.deadLetter.length} message(s) visible(s)`, env.DEAD_LETTER_SUBSCRIPTION));
blocs.push(bloc('obs', 'Monitoring et alertes', 'manuel', 'backlog, âge du plus ancien message et alertes : à montrer dans Monitoring'));
return { plateforme: process.env.K_SERVICE ? 'Cloud Run' : 'local', blocs };
}
export async function creerSource(env = process.env) {
if (env.SOURCE === 'bigquery') return sourceBigQuery(env);
if (env.PUBSUB_EMULATOR_HOST) return sourceEmulateur(env);
return sourceDemo();
}
README.md, le passage sur GCP
# MoveNow, kit de démarrage
Un producteur de positions fictives et un tableau de bord qui montre ce que le pipeline a vraiment écrit : la carte des véhicules, le débit, la dead-letter et la réconciliation des lots. Le kit est un démonstrateur : il sert à vérifier votre infrastructure, il n’est pas évalué.
## Ce qu’il contient
```text
producteur/ Publie des lots cadencés sur Pub/Sub et écrit la liste des événements attendus
tableau/ Tableau de bord, déployable sur Cloud Run en lecture seule
server.mjs Routes de l’API et de l’interface
sources.mjs Démo en mémoire, émulateur Pub/Sub local ou BigQuery réel
public/ Interface : carte, débit, flux, dead-letter, réconciliation
docker-compose.yml Émulateur Pub/Sub, tableau et producteur en continu
```
## Lancer en local
Le plus rapide, sans Docker, sans GCP et sans rien installer : un flux simulé entièrement en mémoire.
```sh
cd tableau
node server.mjs
```
Avec Docker, le flux passe par le vrai émulateur Pub/Sub de Google. Le premier lancement télécharge son image, qui est volumineuse.
```sh
docker compose up --build
```
Ouvrez ensuite http://localhost:8080. Les véhicules bougent sur la carte, quelques messages invalides partent en dead-letter après cinq tentatives. Le bouton « Couper le transfert » imite le retrait du droit d’écriture : la file d’attente grossit, puis se résorbe quand vous rétablissez.
| Rôle | En local | Sur GCP |
| --- | --- | --- |
| Topic et subscriptions | émulateur Pub/Sub | Pub/Sub |
| Transfert vers la table | imité par le tableau (`sources.mjs`) | BigQuery subscription |
| Table analytique | mémoire du tableau | table BigQuery partitionnée sur `event_time` |
| Dead-letter | imitée : cinq essais, puis publication sur le topic d’échec | dead-letter policy de la subscription |
| Listes attendues | volume partagé `lots` | fichiers `.jsonl` à charger dans l’interface |
L’émulateur ne connaît ni IAM, ni les BigQuery subscriptions, ni les métriques de Monitoring. Ce que vous observez en local montre le principe, pas le comportement exact de GCP.
## Passer sur GCP
### Le producteur
Il tourne depuis votre poste ou Cloud Shell, avec ADC. Son identité n’a besoin que de publier sur le topic.
```sh
cd producteur
npm install
node producteur.mjs --project votre-projet --topic positions --rate 1 --duration 10 --batch lot-01
node producteur.mjs --project votre-projet --topic positions --rate 5 --duration 10 --batch lot-02 --invalid 3
```
Chaque lot laisse un fichier `lot-XX.jsonl` : la liste des identifiants attendus. Le débit est plafonné à 100 événements par seconde et la durée à 10 minutes par lot.
### Le contrat de la table
Le producteur envoie du JSON avec ces champs, qui doivent correspondre au schéma de la table quand la subscription écrit avec `use_table_schema` :
| Champ | Type |
| --- | --- |
| `event_id` | STRING |
| `vehicle_id` | STRING |
| `event_time` | TIMESTAMP |
| `latitude` | FLOAT |
| `longitude` | FLOAT |
L’option `--invalid` remplace la latitude par du texte : le topic accepte le message, la table le refuse.
### Le tableau de bord
Il se déploie sur Cloud Run et lit la table en lecture seule. Son identité a besoin de lire la table et de lancer des requêtes, et de s’abonner à la subscription d’inspection si vous voulez voir la dead-letter.
```sh
REGION=europe-west9
PROJECT_ID=votre-projet
IMAGE="$REGION-docker.pkg.dev/$PROJECT_ID/movenow/tableau:v1"
gcloud builds submit --tag "$IMAGE" tableau
```
| Variable | Rôle |
| --- | --- |
| `SOURCE` | `bigquery` pour lire la vraie table. |
| `BQ_TABLE` | `projet.dataset.table`. |
| `BQ_LOCATION` | La localisation du dataset, par exemple `europe-west9`. |
| `GOOGLE_CLOUD_PROJECT` | Le projet, utile pour les noms courts. |
| `DEAD_LETTER_SUBSCRIPTION` | Facultatif. La subscription d’inspection, lue sans acquitter les messages. |
| `REFRESH_SECONDS` | Intervalle de rafraîchissement, 20 secondes par défaut. |
Chaque rafraîchissement lance quelques petites requêtes filtrées sur la partition. C’est peu, mais pas gratuit : fermez le tableau de bord après la démonstration. Le nombre de messages en attente ne se lit pas dans la table ; sur GCP, regardez `num_undelivered_messages` dans Monitoring.
## La carte du déploiement
En haut de l’interface, l’architecture cible est dessinée bloc par bloc. Chaque bloc s’allume selon ce que l’application constate elle-même, avec la preuve affichée dessous :
| Couleur | Sens |
| --- | --- |
| vert, « prouvé » | l’application l’a vérifié elle-même |
| orange, « à revoir » | ça fonctionne, mais c’est un anti-pattern connu |
| rouge, « en échec » | l’application a essayé et ça ne marche pas |
| pointillés, « pas encore détecté » | rien de visible pour l’instant |
| gris, « à prouver vous-même » | invisible depuis l’application : montrez-le dans la console |
| violet, « simulé en local » | l’équivalent local, en attendant le déploiement |
La carte constate, elle ne note pas : un bloc vert ne dit pas que votre choix est le bon, seulement qu’il est en place.
Pour MoveNow, en mode `bigquery`, le tableau de bord lit :
| Bloc | Comment | Droit nécessaire |
| --- | --- | --- |
| Producteur | positions écrites dans la dernière heure | `roles/bigquery.dataViewer` et `roles/bigquery.jobUser`, déjà nécessaires |
| Table BigQuery | partitionnement et comparaison du schéma au contrat du producteur | idem |
| Topic, BigQuery subscription, dead-letter | configuration de la subscription donnée par `SUBSCRIPTION` | `roles/pubsub.viewer` sur la subscription, facultatif |
| Subscription d’inspection | lecture sans acquittement | `roles/pubsub.subscriber` sur elle |
| Variable | Rôle |
| --- | --- |
| `SUBSCRIPTION` | Facultatif. La subscription de transfert, pour afficher sa configuration sur la carte. |
| `DEMO_PUBLIQUE` | `true` pour les démos hébergées : un transfert coupé se rétablit seul au bout d’une minute. |
01 Cadrer
Comprendre la demande
Laissez Terraform de côté pour l’instant. Partez du message du responsable et traduisez-le en un parcours utilisateur et en critères de réussite que l’on peut observer.
À faire
- Choisissez dans la liste ci-dessous les trois questions qui vous semblent prioritaires et posez-les au formateur. Les autres pourront venir ensuite.
- Séparez ce qui est confirmé, ce que vous supposez pour l’instant et ce qui reste à décider.
- Fixez un minimum démontrable, une limite de coût et ce que vous laissez de côté dans la première version.
Les questions à poser
- Combien de conducteurs publient des positions ?
- À quelle fréquence ?
- Quelle est la taille moyenne d’un événement ?
- Combien de temps faut-il conserver l’historique ?
- Combien de temps une panne du traitement doit-elle être tolérée ?
- Les événements invalides sont-ils fréquents, et que doit-on en faire ?
Modèle de fiche de cadrage
| Besoin ou question | Réponse ou hypothèse | Preuve prévue |
|---|---|---|
| À compléter | Qui l’a confirmée ? | Comment la vérifier ? |
Montrez au formateur une fiche courte avec vos hypothèses et trois critères de réussite. Il valide le périmètre et le plafond de charge avant la suite.
02 Architecture
Comprendre l’architecture cible
Cette architecture est votre point de départ. Votre travail consiste à la déployer avec Terraform, à la sécuriser, à l’automatiser et à vérifier qu’elle se comporte comme prévu.
Vous avez le choix entre une BigQuery subscription et un consommateur explicite. Dans la variante native, c’est la subscription elle-même qui fait le transfert : le bloc « Transfert » du schéma n’implique pas un service de plus. Vérifiez quelles erreurs sont réellement prises en charge et quels droits demande le dead-letter.
flowchart TB P["Producteur de positions fictives"] --> T["Topic Pub/Sub"] T --> S["Subscription retry et rétention"] S --> W["Transfert BigQuery natif ou consommateur justifié"] W --> BQ["BigQuery table partitionnée"] S -->|"échecs persistants"| DL["Dead-letter topic"] DL --> DS["Dead-letter subscription inspection et reprise"] SA["Service agent ou identité droits explicites"] -.-> S SA -.-> BQ S -.-> M["Backlog, retard et alertes"]
Sur un petit écran, faites défiler le schéma horizontalement.
Voir et copier le code Mermaid
flowchart TB
P["Producteur de positions fictives"] --> T["Topic Pub/Sub"]
T --> S["Subscription
retry et rétention"]
S --> W["Transfert BigQuery
natif ou consommateur justifié"]
W --> BQ["BigQuery
table partitionnée"]
S -->|"échecs persistants"| DL["Dead-letter topic"]
DL --> DS["Dead-letter subscription
inspection et reprise"]
SA["Service agent ou identité
droits explicites"] -.-> S
SA -.-> BQ
S -.-> M["Backlog, retard et alertes"]
Comment ça circule
Le producteur publie une position
Un événement de test représente un conducteur, un horodatage et une position. Vous choisissez son format et vous le documentez.
Pub/Sub découple les composants
Le topic reçoit les publications. La subscription porte la livraison, et les nouvelles tentatives quand le traitement échoue.
Les données alimentent l’historique
Le mécanisme choisi écrit les événements valides dans BigQuery. Le partitionnement doit suivre vos besoins de rétention et de requêtes.
Les échecs sont mis de côté
Une dead-letter policy, avec les droits du service agent, déplace les messages en échec. Le nombre de tentatives est approximatif : montrez comment se comporte votre configuration.
Calculez le débit avant de choisir. Le nombre d’événements par seconde, c’est le nombre de conducteurs divisé par l’intervalle entre deux positions, en secondes. Par exemple, 2 000 véhicules qui émettent toutes les 4 secondes produisent 500 événements par seconde. Cet exemple n’est pas une donnée du sujet.
Approfondir dans la documentation Google Cloud
À faire à partir du schéma
- Suivez une requête ou un événement de bout en bout et expliquez le rôle de chaque service GCP.
- Listez les ressources Terraform nécessaires, leurs dépendances et ce que chaque module expose aux autres.
- Repérez ce qui reste à décider : région et zones, exposition réseau, identités et permissions, dimensionnement, seuils et politique de reprise.
- Annotez trois choix d’infrastructure avec leurs compromis. Gardez un schéma fidèle à ce que vous déployez réellement.
Présentez au formateur votre découpage Terraform, les flux autorisés et les paramètres retenus. On parle ici d’infrastructure : concevoir une application métier n’est pas l’objet du TP.
03 Construire
Construire le premier parcours
Avancez par petits pas : vous déployez, vous vérifiez, puis vous ajoutez la dépendance suivante. Les ressources finales sont décrites dans Terraform. La console reste utile pour observer et diagnostiquer.
À faire
- Créez le topic, la subscription, puis le dataset et la table de destination avec une stratégie de partitionnement.
- Configurez le transfert choisi et ses identités. Publiez un événement et retrouvez son identifiant dans BigQuery.
- Ajoutez le circuit d’échec, les permissions nécessaires et un moyen de consulter les messages qui y arrivent.
- Publiez un lot traçable et comparez les identifiants attendus, reçus, absents et éventuellement en double.
À montrer
Un event_id suivi de la publication jusqu’à BigQuery, la requête de contrôle, les permissions du transfert et un événement en échec expliqué.
À expliquer
Que prouve un nombre de messages publiés égal au nombre de lignes reçues ? Comment détecter une perte masquée par un doublon ?
Besoin d’un indice ?
Indice 1, une piste
Suivez les identifiants, pas seulement le nombre de lignes. Décidez où le format est validé avant de préparer un message invalide.
Indice 2, où chercher
Transfert de Pub/Sub vers BigQuery
Les ressources du provider sont ensuite listées dans la rubrique Documentation Terraform.
Indice 3, un coup de pouce
Une subscription BigQuery évite d’écrire un consommateur pour un transfert simple. Sa livraison est « au moins une fois » : prévoyez d’analyser les doublons. Et si le topic refuse un message dès la publication, ce message n’arrivera jamais dans une dead-letter subscription.
Organiser Terraform
Découpez vos modules par responsabilité. Écrivez d’abord leurs entrées, leurs sorties et leurs dépendances : le module racine se contente de les assembler. Chaque module a une courte documentation. Inutile en revanche de créer un module par ressource.
Un découpage à discuter, une fois votre propre proposition faite
Chaque module contient main.tf, variables.tf, outputs.tf et un court README. Le module racine relie les sorties des uns aux entrées des autres. Ce découpage est une proposition : à vous de le défendre ou de l’améliorer.
| Module | Responsabilité | Entrées principales | Sorties utiles |
|---|---|---|---|
| messaging | Le topic principal, le topic et la subscription de dead-letter. | retention, labels, noms des topics | topic_id, dead_letter_topic_id, dead_letter_subscription |
| analytics | Le dataset, le schéma, la table partitionnée et son expiration. | location, schema, partition_field, retention | dataset_id, table_id |
| delivery | La subscription de transfert, le retry et le dead-letter, et les droits du service agent. | topic_id, table_id, dead_letter_topic_id, retry, service_agent | writer_identity, export_subscription_id |
| observability | Le backlog, le retard, les erreurs et le volume en dead-letter. | subscription_ids, seuils, canaux | dashboard_id, alert_policy_ids |
Une organisation de dépôt possible :
infra/
bootstrap/ # backend GCS et accès CI, cycle séparé
envs/
lab/ # module racine déployé, backend et variables
prod-design/ # variante documentée, sans déploiement imposé
modules/
messaging/
analytics/
delivery/
observability/
app/ # le kit de démarrage, si vous l’utilisez
tests/
smoke/ # vérification du chemin nominal
load/ # scripts et jeux synthétiques
docs/
architecture.mmd
decisions/ # décisions et alternatives (ADR)
runbooks/ # diagnostic et reprise
evidence/ # résultats, sans secret
README.mdÉtat, environnements et bootstrap
L’état Terraform vit dans un bucket GCS créé par un bootstrap séparé, lancé avant le laboratoire et qui n’en dépend pas. Activez le versioning de ce bucket et limitez qui peut y accéder. Le backend GCS gère le verrouillage, et un préfixe par environnement évite de mélanger les états.
Ne commitez ni l’état, ni les plans, ni les clés, ni les valeurs de secrets. Attention, sensitive = true masque une valeur à l’affichage mais ne la retire pas du state. Gardez plutôt un fichier d’exemple avec les variables non sensibles. Voir la documentation des modules Terraform.
Le chemin nominal fonctionne. terraform fmt -check et terraform validate passent, et un nouveau plan ne propose rien d’inattendu. Un autre membre du groupe sait expliquer les dépendances.
04 Automatiser
Rendre le déploiement reproductible
Un collègue doit pouvoir relire puis déployer un changement sans refaire vos manipulations dans la console.
À faire
- À chaque changement, la CI vérifie le format, initialise sans backend et lance
validate. Les versions sont figées et le fichier de verrouillage est commité. - La CI s’authentifie par fédération d’identité, limitée à votre dépôt et aux branches ou environnements autorisés. Aucune clé JSON durable dans Git.
- Le plan est produit sur une branche de confiance, relu, puis appliqué tel quel après une approbation explicite. Un seul déploiement à la fois par environnement.
- Un test du parcours nominal suit l’apply. Si quelque chose échoue, le pipeline s’arrête là.
- La destruction est un job manuel et approuvé, suivi d’un inventaire des ressources restantes. Le bootstrap garde son propre cycle.
À montrer
Un petit changement suivi de son commit jusqu’au test final, et une erreur de validation volontaire arrêtée avant tout apply. Les plans sont des artefacts privés, conservés peu de temps.
Besoin d’un indice ?
Indice 1, une piste
Dessinez les jobs et demandez-vous quelle identité exécute chacun d’eux. Qu’est-ce qui prouve que le plan appliqué est bien celui qui a été relu ?
Indice 2, où chercher
La page du backend GCS et celle sur la fédération d’identité pour les pipelines.
Indice 3, un coup de pouce
Séparez quatre temps : une validation sans accès au cloud, un plan authentifié, une approbation, puis l’apply. Rattachez le plan au commit et à l’environnement. Si le code, les variables ou l’état changent, le plan doit être recalculé et relu à nouveau.
Montrez une exécution de la CI avec son test final. Une procédure manuelle documentée aide au diagnostic, mais elle ne remplace pas cette preuve d’automatisation.
05 Éprouver
Observer, tester, expliquer
Annoncez votre hypothèse avant chaque test. Un graphique doit répondre à une question précise : avoir un dashboard ne prouve pas, à lui seul, que le système est fiable.
Préparer les mesures
L’âge du plus ancien message non acquitté, le backlog, le délai jusqu’à BigQuery et les événements qui partent dans le circuit d’échec.
Déployez avec Terraform un dashboard et deux alertes utiles. Pour chacune, précisez la métrique, le filtre, la fenêtre, le seuil, le destinataire et la première action de diagnostic. Vérifiez que la notification arrive bien, puis qu’elle se referme quand tout revient à la normale.
Des signaux à adapter
| Signal | Où le trouver | Décision ou exemple de seuil |
|---|---|---|
| Âge du plus ancien message non acquitté | Pub/Sub, oldest_unacked_message_age | Exemple pour le labo : plus de 60 s pendant 2 minutes. Vérifiez alors les droits, le schéma et le consommateur. |
| Messages en attente | Pub/Sub, num_undelivered_messages | Regardez la pente et le temps de résorption, pas seulement la valeur à un instant donné. |
| Erreurs de transfert et dead-letter | Logs et métriques du transfert retenu | Alertez sur un volume anormal et gardez l’event_id pour réconcilier et rejouer. |
L’essai de charge du laboratoire
Commencez à 10 événements par seconde pendant 60 secondes. Une fois le budget validé, comparez 10, 25 puis 50 événements par seconde, 60 secondes chacun. Arrêtez ensuite le producteur et mesurez le temps de rattrapage.
Avant de lancer quoi que ce soit, fixez la cible autorisée, le plafond de ressources, la durée, le volume et la règle d’arrêt. Arrêtez si le coût dérape ou si la dégradation persiste. Fixez aussi un seuil de retard ou de backlog et un temps de drainage maximal, à faire valider avant le test.
Une expérience pour vous entraîner
Injectez un petit lot numéroté qui ne respecte pas le contrat choisi. Repérez à quel stade il est rejeté et montrez comment le corriger. Pour observer le circuit d’échec, il faut un cas accepté à l’ingestion puis refusé au traitement.
Écrivez la procédure de retour avant de commencer et ne changez qu’un paramètre à la fois. Si vous modifiez quelque chose à la main pour diagnostiquer, remettez ensuite Terraform en accord avec l’état voulu.
À montrer
Un tableau d’identifiants réconciliés, les courbes de retard et une procédure de rejeu qui tient compte des doublons.
Trame du compte rendu
- L’hypothèse et le résultat attendu.
- Le périmètre, la charge, l’heure de début et l’heure de fin.
- Ce que vous avez observé avant, pendant et après, et l’impact sur le parcours utilisateur.
- Les pistes de diagnostic, les vérifications faites et la cause retenue.
- La correction, la preuve du retour au service et les limites de votre conclusion.
En cours de route, le formateur vous annoncera un événement supplémentaire. Vous devrez en expliquer l’impact et proposer une réponse. Selon le temps et le budget, cette réponse pourra mêler une expérience bornée et une évolution d’architecture argumentée.
Votre compte rendu sépare ce qui a été mesuré, ce qui reste une hypothèse et ce qui demanderait un test plus large. Chaque membre sait lire les courbes.
06 Présenter
Présenter ce que vous avez réalisé
Chaque groupe dispose de 20 minutes de présentation technique, démonstration comprise. Les trois membres prennent la parole. Montrez la solution telle qu’elle est, avec ses résultats et ses limites, et gardez des traces de secours au cas où la démonstration en direct échouerait.
| Séquence | Durée | À montrer |
|---|---|---|
| Besoin et périmètre | 2 min | Les hypothèses, les objectifs et le minimum réellement livré. |
| Architecture réalisée | 4 min | Le schéma, les flux, la sécurité, vos choix et les alternatives écartées. |
| Terraform et pipeline | 4 min | Les modules et leurs interfaces, l’état distant et la trace d’un déploiement par la CI. |
| Démonstration et résultats | 6 min | Le chemin nominal, l’essai de charge, l’incident, le monitoring et la reprise. Prévoyez des captures au cas où la démonstration échouerait. |
| Coûts et limites | 3 min | Les coûts estimés, les écarts avec la production, ce qui n’est pas fait et la suite. |
| Conclusion | 1 min | Le bilan technique et ce que le groupe retient. |
Ce que votre support doit référencer
- Le dépôt, son README de déploiement et de destruction, le contrat de chaque module et une trace du pipeline.
- Le schéma réellement déployé et trois décisions argumentées.
- Les résultats de charge, le compte rendu d’incident, les alertes et le runbook.
- Une estimation des principaux coûts, ce qui n’a pas été réalisé et l’écart avec une cible de production.
Après la présentation
À la date convenue avec le formateur, détruisez les ressources du laboratoire, vérifiez ce qui reste et signalez ce que vous gardez volontairement. Tenir le budget et nettoyer font partie du travail.
Répétez avec un chronomètre. Chaque membre doit pouvoir expliquer un flux, une permission, une panne et une limite sans se contenter de lire les diapositives.
Où en êtes-vous ?
Socle
Version complète
Ces cases sont enregistrées dans votre navigateur uniquement. Le formateur ne les voit pas.
Si vous êtes en avance
Choisissez une seule extension, une fois le parcours principal terminé et son coût validé. Ajouter des services ne remplace pas des preuves de qualité.
- Ajouter une validation de schéma à l’entrée et comparer les deux types de rejet.
- Comparer des requêtes BigQuery avec et sans filtre de partition.
- Décrire une stratégie d’idempotence et la démontrer sur un lot rejoué.
Documentation Terraform
Ouvrez ces aides quand vous en avez besoin. Le rôle de chaque service est indiqué, mais leur assemblage, leurs paramètres et leurs permissions restent à concevoir.
Ressources Terraform utiles
Cette liste est un point de départ et votre architecture demandera d’autres ressources. Chaque lien ouvre la page du provider Google, avec ses arguments et ses exemples.
| Ressource | Rôle | Conseil |
|---|---|---|
google_pubsub_topic | Recevoir les événements et créer le topic de dead-letter. | Donnez à chaque événement un event_id et un horodatage pour suivre son parcours. |
google_pubsub_subscription | Configurer la livraison, les nouvelles tentatives, la rétention, le dead-letter et éventuellement l’export BigQuery. | Choisissez le mode de livraison. En export natif, bigquery_config remplace un consommateur à développer. |
google_bigquery_dataset | Créer le conteneur analytique. | Fixez la localisation avant de créer les tables. |
google_bigquery_table | Définir le schéma et le partitionnement. | Vérifiez les types, la colonne de partition et leur correspondance avec les messages produits. |
google_pubsub_topic_iam_member | Autoriser la publication, notamment vers le topic de dead-letter. | Le service agent Pub/Sub a besoin de droits sur le parcours de dead-letter, et aussi sur la subscription source. |
google_monitoring_alert_policy | Détecter une accumulation ou un retard. | Commencez par l’âge du plus ancien message et observez comment il redescend après la reprise. |
Ressources communes et configuration du provider
google_project_service active les API nécessaires, et google_service_account crée une identité dédiée. Donnez à chaque identité les permissions dont elle a besoin, au bon niveau, plutôt qu’un rôle Owner ou Editor sur tout le projet.
Pour l’observabilité, il vous faudra un dashboard, une politique d’alerte et un canal de notification. Vérifiez que la notification arrive réellement.
terraform {
required_providers {
google = {
source = "hashicorp/google"
version = "8.6.0"
}
}
}
provider "google" {
project = var.project_id
region = var.region
}
Cette version est un exemple, vérifié lors de la préparation du support. Contrôlez la dernière version publiée avant le TP. Déclarez vos variables et commitez .terraform.lock.hcl. Si votre projet utilise déjà une autre version, n’en changez pas sans relire le plan et les notes de migration. En local, utilisez ADC ; en CI, une identité fédérée. Jamais de clé JSON dans le code.
Critères d’évaluation
L’évaluation porte sur ce que le groupe a réalisé, sur ses résultats et sur la façon dont il les explique. Il n’y a pas de barème chiffré.
| Critère | Ce qui doit être observable |
|---|---|
| Cadrage et architecture | Des hypothèses explicites, des flux lisibles, des alternatives et des compromis documentés. |
| Modules et reproductibilité | Des responsabilités cohérentes, des interfaces claires, un état distant, des versions figées et un redéploiement documenté. |
| Pipeline Terraform | Une validation automatique, un plan relu puis appliqué tel quel, et des identités CI limitées. |
| Réseau, IAM et secrets | Des expositions justifiées, le moindre privilège, aucun secret dans Git et un state protégé. |
| Charge et résilience | Un protocole reproductible, des mesures avant et après, un incident analysé et un retour au service démontré. |
| Observabilité | Un dashboard utile, des seuils justifiés, une notification vérifiée et une procédure de diagnostic. |
| Coût et périmètre | Des ordres de grandeur, les limites du laboratoire, la destruction et le contrôle des ressources restantes. |
| Présentation et défense collective | Le groupe tient ses 20 minutes et montre ce qui fonctionne vraiment, avec ses limites. Chaque membre explique les choix, y compris le code proposé par un LLM. |