7.6

Voir en anglais

7.6 Données en temps réel et en flux

Vue d’ensemble et motivation

La plupart de ce que vous savez sur les pipelines de données suppose que les données se tiennent immobiles. Vous collectez une journée d’enregistrements, exécutez une tâche pendant la nuit, et lisez les résultats le matin. Les données en temps réel et en flux inversent cette hypothèse. Au lieu de traiter un tas fini de données, vous traitez un flux sans fin d’événements à mesure qu’ils arrivent, et vous produisez des réponses continuellement. C’est la différence entre le traitement par lots, qui opère sur un jeu de données borné et complet, et le traitement de flux, qui opère sur un flux non borné et jamais terminé.

Pour les grandes équipes, le streaming apparaît au moment où la latence commence à compter pour l’entreprise. Une décision de fraude qui arrive une heure en retard ne vaut rien. Un signal de personnalisation qui atterrit demain ne personnalise rien. Un tableau de bord opérationnel qui traîne derrière la réalité d’un quart de travail trompe les gens qui l’observent. Le chapitre 7.2 (ingénierie de données) argumente que vous devriez choisir le par lots par défaut et recourir au streaming seulement là où la latence paie véritablement, et ce chapitre vous emmène le reste du chemin : quand le temps réel gagne son coût, et comment le construire sans mettre le feu à votre budget d’opérations. Le streaming se trouve proche des motifs de messagerie pilotée par événement du chapitre 3.12 (architecture événementielle et messagerie), des choix de stockage du chapitre 3.4 (architecture et stockage de données), et des pratiques de télémétrie du chapitre 9.2 (observabilité et télémétrie).

Les contextes d’entreprise et gouvernementaux élèvent les enjeux. Une banque note chaque transaction de carte pour la fraude dans le temps qu’il faut à un lecteur de carte pour clignoter. Une agence de transit suit les véhicules et prédit les arrivées pour des millions de passagers. Une agence de prestations surveille les anomalies dans les réclamations tout en gardant un enregistrement auditable de chaque décision. Dans tous ces cas, la valeur vient d’agir sur des données pendant qu’elles sont encore fraîches, et le risque vient d’agir sur des données qui sont fausses, incomplètes, ou impossibles à reconstruire plus tard. Ce chapitre a une opinion sur les deux.

Principes clés

  • Recourez au streaming seulement quand la latence a une valeur d’affaires claire ; le par lots est moins cher et plus simple.
  • Distinguez les données bornées (finies) des données non bornées (jamais terminées), et concevez en conséquence.
  • Traitez le temps d’événement, pas le temps d’arrivée, comme la source de vérité, et planifiez pour les données tardives et désordonnées.
  • Les fenêtres et filigranes sont comment vous obtenez des réponses finies depuis des flux infinis.
  • Préférez des résultats effectivement-une-fois à travers des puits idempotents à de fragiles promesses exactement-une-fois.
  • Le traitement avec état a besoin de points de contrôle pour pouvoir récupérer sans perdre ou compter en double.
  • Concevez pour la contre-pression et le retraitement dès le premier jour, pas comme une pensée après coup.
  • Gardez la logique de streaming observable et auditable ; un flux silencieux est pire qu’un lot échoué.

Recommandations

Justifier le temps réel avant de le construire

La décision de streaming la plus importante est de savoir si streamer du tout. Le temps réel double approximativement votre complexité et coût opérationnels, parce que vous échangez une tâche qui s’exécute et s’arrête contre un système qui doit rester sain chaque seconde. Avant de vous engager, nommez la décision que les données fraîches permettent et le coût de cette décision arrivant en retard. Le scoring de fraude, l’alerte opérationnelle, et la personnalisation en direct franchissent habituellement la barre. Un tableau de bord qu’un humain regarde deux fois par jour ne le fait presque jamais, peu importe à quel point « temps réel » sonne satisfaisant dans une réunion de planification. Écrivez l’exigence de latence comme un nombre, en secondes ou minutes, et vérifiez-la contre la réalité. Une grande partie de ce que les gens appellent temps réel est bien servie par des micro-lots qui s’exécutent toutes les quelques minutes à une fraction du coût.

Concevoir autour du temps d’événement, pas du temps de traitement

L’idée la plus difficile du streaming est que les événements se produisent à un moment et sont traités à un autre. Le temps d’événement est quand la chose s’est réellement produite, par exemple quand un passager a tapé une carte. Le temps de traitement est quand votre système a eu le temps de le gérer. Ceux-ci dérivent constamment l’un de l’autre : un téléphone perd le signal dans un tunnel et téléverse trois minutes de tapotements à la fois, un hoquet réseau réordonne les messages, une partition traîne. Si vous calculez sur le temps de traitement, vos chiffres oscillent avec votre infrastructure plutôt que de refléter le monde. Ce problème de données tardives et désordonnées est le cœur de la discipline, et il se connecte directement à la modélisation d’événement dans l’architecture événementielle. Estampillez chaque événement avec son temps d’événement à la source, portez cet horodatage à travers tout le pipeline, et calculez vos résultats contre lui.

Utiliser des fenêtres et filigranes pour obtenir des réponses finies

Un flux non borné ne se termine jamais, donc « compter les événements » n’a pas de réponse jusqu’à ce que vous le borniez. Les fenêtres font ce bornage. Les fenêtres tumbling découpent le temps en compartiments fixes et non chevauchants, par exemple chaque minute. Les fenêtres glissantes se chevauchent, donc une fenêtre de cinq minutes qui avance chaque minute vous donne un chiffre mobile lisse. Les fenêtres de session groupent les rafales d’activité séparées par des écarts d’inactivité, ce qui convient bien aux sessions utilisateur. Une fois que vous avez des fenêtres, vous devez décider quand une fenêtre est terminée, parce que des données tardives pourraient encore arriver. Un filigrane est l’estimation du système qu’il a probablement vu tous les événements jusqu’à un temps d’événement donné. Quand le filigrane dépasse la fin d’une fenêtre, vous émettez le résultat. Réglez combien de temps vous attendez : gardez les fenêtres ouvertes plus longtemps et vous tolérez plus de retard au coût de la latence et de la mémoire, fermez-les plus vite et vous risquez de laisser tomber des retardataires. Décidez explicitement ce qui arrive aux données qui arrivent après qu’une fenêtre se ferme, si vous les laissez tomber, les journalisez, ou émettez une correction.

Rendre les puits idempotents et préférer effectivement-une-fois

Les garanties de livraison sonnent simples et ne le sont pas. La livraison au-moins-une-fois signifie que chaque événement est traité, mais certains peuvent être traités plus d’une fois après une retentative, donc les comptes peuvent s’inflater. Exactement-une-fois sonne idéal mais est coûteux et, pris littéralement à travers des systèmes externes arbitraires, souvent impossible. La cible pratique est effectivement-une-fois : le résultat observable est comme si chaque événement était traité une fois, même si la machinerie a retenté en dessous. Vous y arrivez en rendant vos puits idempotents sûrs à écrire à répétition, en utilisant des clés déterministes et des upserts pour qu’un événement rejoué écrase plutôt que de dupliquer. Combinez la livraison au-moins-une-fois avec des écritures idempotentes et vous obtenez des résultats corrects sans payer pour une coordination transactionnelle lourde partout. Réservez la vraie machinerie exactement-une-fois aux endroits étroits qui en ont véritablement besoin.

Faire des points de contrôle du traitement avec état pour qu’il puisse récupérer

De nombreux calculs de streaming utiles sont avec état : comptes courants, jointures à travers des flux, déduplication, modèles de fraude qui se souviennent du comportement récent. Cet état vit en mémoire et disparaîtrait quand un processus redémarre. Les points de contrôle instantanent périodiquement l’état et la position de flux ensemble, pour qu’après un crash le système reprenne depuis un point cohérent plutôt que de rejouer tout ou de perdre sa mémoire. Dimensionnez votre état délibérément, parce que l’état non borné est une façon courante de faire manquer de mémoire à une tâche de streaming en production. Utilisez l’expiration et la durée de vie sur l’état dont vous n’avez plus besoin, et surveillez la taille d’état comme une métrique de première classe. Le temps de récupération après un échec est une vraie préoccupation de niveau de service, donc testez-le avant vos utilisateurs.

Streamer depuis des bases de données opérationnelles avec la capture de changement de données

Vous voulez souvent réagir aux changements dans une base de données qui n’a jamais été conçue pour émettre des événements. La capture de changement de données (CDC) résout cela en lisant le journal de transaction de la base de données et en transformant chaque insertion, mise à jour, et suppression en un flux d’événements de changement. C’est bien mieux que d’interroger la table sur une minuterie, ce qui est lent, manque les états intermédiaires, et martèle la source. La CDC vous laisse garder un index de recherche, un cache, un magasin d’analytique, ou un service en aval continuellement synchronisé avec un système d’enregistrement, et cela sans changements invasifs à l’application. Traitez le flux de changement comme un produit de données de première classe : versionnez son schéma, documentez son sens, et surveillez son retard, parce que tout en aval hérite de ce retard.

Préférer une architecture streaming d’abord au maintien de deux bases de code

L’architecture Lambda classique exécute une couche par lots pour un historique précis et complet aux côtés d’une couche de vitesse pour des résultats frais et approximatifs, puis les fusionne. Cela fonctionne, mais cela vous fait écrire et maintenir la même logique métier deux fois, dans deux systèmes, et réconcilier les différences pour toujours. L’architecture Kappa effondre cela : gardez un journal d’événements durable et rejouable et exécutez tout le traitement comme traitement de flux, retraitant l’historique en rejouant le journal quand la logique change. L’industrie a dérivé vers cette forme streaming d’abord parce qu’une seule base de code est dramatiquement moins chère à maintenir et raisonner. Si vous pouvez exprimer vos besoins par lots comme des replays sur un journal d’événements retenu, vous évitez entièrement la taxe de deux bases de code. Utilisez des courtiers basés sur journal qui retiennent l’historique pour que le retraitement soit une question de rembobinage, pas de reconstruction.

Exposer les flux comme SQL, vues matérialisées, et OLAP temps réel

Tout le monde qui a besoin de streaming ne devrait pas avoir à écrire du code de traitement de flux de bas niveau. Le SQL de streaming laisse les analystes et ingénieurs exprimer fenêtres, jointures, et agrégations dans un langage qu’ils connaissent déjà, et il garde les résultats continuellement à jour comme vues matérialisées. Pour les requêtes analytiques à faible latence sur des données fraîches, un magasin traitement analytique en ligne (OLAP) temps réel ingère le flux et répond aux requêtes de découpage en millisecondes, ce qui est ce qui alimente un tableau de bord opérationnel véritablement en direct. Associez ceux-ci aux pratiques d’analytique de produit du chapitre 7.4 (analytique de produit et expérimentation) quand l’objectif est un retour rapide sur les fonctionnalités et expériences. Choisissez ces outils de plus haut niveau où ils conviennent, et gardez les processeurs de flux écrits à la main pour la logique qu’ils ne peuvent pas exprimer.

Planifier pour la contre-pression et le retraitement dès le début

Un flux peut arriver plus vite que vous ne pouvez le traiter. La contre-pression est le mécanisme qui laisse un consommateur lent signaler en amont de ralentir plutôt que de s’effondrer ou de laisser tomber des données silencieusement. Assurez-vous que chaque étape dans votre pipeline l’honore, et surveillez le retard de consommateur comme une métrique phare, parce qu’un retard croissant est l’avertissement le plus précoce que vous perdez la course. Le retraitement est l’autre capacité que les gens souhaitent avoir construite. Quand vous trouvez un bogue ou changez une règle, vous voulez rejouer l’historique à travers la logique corrigée. Cela n’est possible que si votre journal d’événements retient assez d’historique et vos puits sont assez idempotents pour absorber le replay. Concevez les deux dès le premier jour ; les rétro-adapter sous pression d’incident est misérable.

Compromis : avantages et inconvénients

ChoixAvantagesInconvénientsMeilleur ajustement
Par lotsSimple, bon marché, facile à tester et rétro-remplirLatence élevée, périmé entre exécutionsRapport, la plupart de l’analytique
Micro-lots (minutes)Proche du temps réel, bien plus simple que le streamingPas véritablement instantanéTableaux de bord « temps réel »
Vrai streaming (sous-seconde)Réaction instantanée, résultats continusComplexe, coûteux, difficile à testerFraude, alerte, personnalisation en direct
Au-moins-une-fois + puits idempotentRésultats corrects, abordable, résilientExige une conception de clé disciplinéeLa plupart des pipelines de streaming
Machinerie exactement-une-foisGarantie forte de bout en boutCoûteux, limité à travers les systèmesChemins étroits à haut enjeu
Lambda (par lots + vitesse)Historique précis plus vue fraîcheDeux bases de code à maintenirMigrations héritées
Kappa (streaming d’abord)Une base de code, rejouableExige un journal durable et retenuNouvelles plateformes de streaming

La tension centrale est la latence contre la complexité. Chaque pas vers le temps réel vous coûte en fardeau opérationnel, difficulté de test, et argent, et les retours ne sont pas linéaires : passer de quotidien à toutes-les-quelques-minutes est bon marché et souvent suffisant, tandis que passer de minutes à sous-seconde est là où la dépense se concentre. Résolvez la tension en évaluant la décision, pas la technologie. Demandez quelle action la fraîcheur permet et ce que le retard coûte, puis achetez seulement autant de réduction de latence que cette action justifie. Quand vous avez besoin de streaming, appuyez-vous sur la livraison au-moins-une-fois avec des puits idempotents et un journal streaming d’abord, parce que cette combinaison vous donne la correction et la rejouabilité sans les garanties les plus lourdes.

Questions à discuter avec votre équipe

  1. Quelle décision les données temps réel permettent-elles réellement pour nous, et que coûte-t-il quand ces données arrivent une minute en retard au lieu d’instantanément ? C’est la question qui devrait conditionner chaque projet de streaming, parce que le streaming double approximativement votre coût et complexité opérationnels comparé au par lots. Une grande équipe peut brûler des trimestres à construire une plateforme temps réel qui sert des tableaux de bord qu’un humain vérifie deux fois par jour, ce qui est de l’argent brûlé. Apportez l’action concrète que les données pilotent, que ce soit bloquer une transaction frauduleuse, appeler un opérateur, ou changer ce qu’un utilisateur voit, et mettez un chiffre sur le coût de latence pour chacune. Si la réponse honnête est qu’un micro-lot de cinq minutes servirait le besoin, c’est une découverte qui vaut la peine d’être célébrée, pas cachée. La réponse devrait directement changer si vous construisez du vrai streaming, vous contentez de micro-lots, ou restez en par lots.

  2. Comment gérons-nous les événements tardifs et désordonnés, et que se passe-t-il avec les données qui arrivent après qu’une fenêtre se ferme ? Les données tardives et désordonnées sont la partie difficile du streaming, et les équipes qui sautent cette question la découvrent en production quand leurs chiffres refusent de se réconcilier. Les pressions concurrentes sont la latence et la correction : gardez les fenêtres ouvertes plus longtemps pour attraper les retardataires et vous retardez chaque résultat et consommez plus de mémoire, fermez-les plus vite et vous laissez tomber silencieusement de vraies données. Apportez une preuve de combien de retard vos données ont réellement, mesuré comme l’écart entre le temps d’événement et le temps de traitement à travers vos sources, puisqu’une source mobile dans des tunnels se comporte très différemment d’un événement côté serveur. Décidez explicitement si les données tardives sont laissées tomber, journalisées, ou déclenchent une correction, et assurez-vous que tout le monde en aval sait laquelle. Dans un contexte gouvernemental où les chiffres doivent être défendables, laisser tomber silencieusement des événements tardifs peut être un problème de conformité, donc la politique doit être délibérée et documentée.

  3. Nos puits sont-ils assez idempotents pour que nous puissions rejouer l’historique en sécurité, et notre journal d’événements retient-il assez pour rendre le replay possible ? Le retraitement est la capacité que les équipes souhaitent le plus souvent avoir construite et n’ont le plus souvent pas, et cela dépend de deux choses fonctionnant ensemble : des puits idempotents qui absorbent les événements rejoués sans dupliquer, et un journal durable qui retient assez d’historique pour rejouer depuis. Sans les deux, corriger un bogue de logique signifie que vous ne pouvez pas recalculer proprement la période affectée, et vous êtes coincé à corriger les chiffres à la main sous pression. Apportez votre fenêtre de rétention actuelle et un test concret : choisissez un vrai bogue du dernier trimestre et demandez si vous auriez pu rejouer la logique corrigée sur les données affectées. L’attraction contre cela est le coût, puisque retenir l’historique et concevoir des écritures idempotentes prend du stockage et de la discipline en amont. Mais l’alternative fait surface au pire moment possible, pendant un incident, donc la réponse façonne combien vous investissez dans la rejouabilité avant d’en avoir besoin.

  4. Quand une tâche de streaming plante, à quelle vitesse doit-elle récupérer, combien d’état est-elle autorisée à tenir, et avons-nous réellement chronométré une récupération sous charge de production ? Une tâche par lots qui meurt peut être réexécutée demain, mais un flux toujours actif qui meurt est une panne en cours, et les tâches avec état qui tiennent des comptes courants, jointures, ou modèles de fraude peuvent perdre des minutes de mémoire ou prendre longtemps à recharger l’état après un redémarrage. Pour une grande équipe, c’est là qu’un détail peu glamour fixe discrètement votre vraie disponibilité : l’état non borné grandit jusqu’à ce qu’une tâche manque de mémoire, et une restauration de point de contrôle lente transforme un blip de dix secondes en un de dix minutes. Les pressions concurrentes sont la fraîcheur contre la sécurité, parce que des points de contrôle plus fréquents raccourcissent la récupération mais ajoutent de la surcharge, et une rétention d’état généreuse améliore la précision mais risque l’épuisement de mémoire. Apportez un objectif de temps de récupération concret, votre taille d’état actuelle et sa courbe de croissance, votre intervalle de point de contrôle, et les résultats d’un vrai exercice de basculement plutôt qu’une estimation optimiste. Dans les contextes d’entreprise et gouvernementaux où le flux soutient le scoring de fraude ou un flux de sécurité publique, un chemin de récupération non testé est un risque opérationnel que vous avez accepté sans le mesurer, donc traitez l’exercice comme une exigence, pas une gentillesse.

  5. Exécutons-nous une base de code streaming d’abord unique ou une couche par lots séparée et couche de vitesse, et que nous coûte-t-il réellement de garder les deux réconciliées ? Le motif Lambda d’une couche par lots pour l’historique précis plus une couche de vitesse pour les résultats frais vous force à écrire la même logique métier deux fois, dans deux systèmes, et réconcilier leurs réponses pour toujours, tandis qu’une forme streaming d’abord (Kappa) garde un journal durable et rejouable et exécute tout le traitement comme traitement de flux. Pour une grande organisation la logique dupliquée est là où la dérive et les chiffres disputés se reproduisent, parce qu’une règle change dans une couche et pas l’autre, et les ingénieurs passent du vrai temps à expliquer pourquoi les deux sont en désaccord. L’attraction vers garder les deux est l’inertie et le confort d’une couche par lots éprouvée, donc pesez cela contre la taxe de maintenance honnêtement. Apportez la liste des calculs que vous exécutez actuellement aux deux endroits, les incidents causés par le désaccord des deux couches, et une évaluation de si votre journal d’événements retient assez d’historique pour exprimer les besoins par lots comme replays. Dans les contextes gouvernementaux et d’entreprise audités, deux couches qui peuvent rapporter des chiffres différents pour la même période sont elles-mêmes un passif de conformité, puisque vous devez pouvoir dire quel chiffre fait autorité et pourquoi.

  6. Qui exploite ce système toujours actif quand il se casse à trois heures du matin, et avons-nous budgétisé la charge d’astreinte et les compétences spécialisées qu’il exige, ou supposons-nous une dotation en forme de par lots ? Le streaming déplace le coût de la construction vers l’exploitation : le système doit rester sain chaque seconde, ce qui signifie une vraie couverture d’astreinte, des ingénieurs fluides en temps d’événement, filigranes, état, et sémantique de livraison, et un test plus difficile que pour une tâche qui s’exécute et s’arrête. Les équipes approuvent routinièrement une plateforme de streaming sur la force de ses capacités et ne financent jamais les gens qui la gardent en vie, donc la plateforme se dégrade et la confiance s’érode. Le compromis est la portée contre la durabilité : chaque pipeline temps réel supplémentaire est une chose de plus qui peut appeler quelqu’un, donc la question est si la latence qu’il achète justifie un engagement opérationnel permanent. Apportez un inventaire honnête de qui possède chaque flux en production, votre rotation d’astreinte actuelle et sa marge, et où se trouve réellement l’expertise en temps d’événement, que ce soit une embauche, un partenaire, ou un service géré. Pour un organisme public ou une grande entreprise, ajoutez les délais d’approvisionnement et d’embauche et toute option de service géré, parce qu’une plateforme temps réel qui dépend d’un talent rare que vous ne pouvez pas recruter ou retenir est un plan d’exploiter un système sujet aux pannes sous-doté.

Regard sectoriel

Jeune pousse. Le streaming est rarement votre premier mouvement, et monter une plateforme lourde peut couler une minuscule équipe. Choisissez le seul signal qui touche votre valeur centrale, placez les événements sur un seul courtier basé sur journal retenu, et exécutez un processeur léger avec des puits clés et idempotents pour qu’une retentative au-moins-une-fois ne compte jamais en double. Gardez quelques jours d’historique pour pouvoir rejouer à travers la logique fixe, et préférez un service de streaming géré à exploiter votre propre cluster, parce que votre ressource la plus rare est l’attention d’ingénierie.

Petite entreprise. Vous n’avez probablement pas de spécialiste de streaming et aucun appétit pour exploiter une infrastructure toujours active, donc traitez le temps réel comme quelque chose que vous achetez dans des outils que vous utilisez déjà plutôt qu’un système que vous dotez. Cadrez le besoin comme une question de latence avec un chiffre attaché, et dans la plupart des cas un micro-lot qui se rafraîchit toutes les quelques minutes le satisfera à une fraction du coût et risque. Choisissez des fournisseurs dont les fonctionnalités temps réel sont transparentes sur le retard et faciles à replier, et réservez le streaming personnalisé pour le rare cas où les données fraîches pilotent directement le revenu ou la sécurité.

Grande entreprise. Le problème est la cohérence et le coût à travers de nombreuses équipes : une plateforme partagée basée sur journal, une politique standard de temps d’événement et données tardives, et des puits idempotents pour que les groupes arrêtent de réinventer des pipelines fragiles. Budgétisez explicitement les opérations toujours actives et le fardeau d’astreinte, standardisez sur un journal streaming d’abord pour éviter une base de code par lots dupliquée, et gérez les flux comme des produits de données gouvernés avec propriétaires, versionnage de schéma, et retard surveillé plutôt qu’une dispersion de tâches sur mesure. Suivez la latence, le temps de récupération, et le coût par flux comme métriques de portefeuille.

Gouvernement. L’auditabilité et la responsabilité publique façonnent chaque choix. Retenez chaque événement traité dans un journal durable pour que les chiffres rapportés aux organismes de surveillance, fréquentation, anomalies de prestations, décisions de fraude, puissent être reconstruits exactement, et rendez la politique de données tardives explicite et documentée plutôt que de laisser tomber silencieusement des événements. L’approvisionnement devrait exiger la portabilité de données et la divulgation des garanties de livraison et rétention d’un service géré, et toute reformulation après un changement de règle devrait être un replay défendable à travers la logique corrigée, pas un correctif manuel que personne ne peut tracer.

Exemples

Jeune pousse. Une application grand public veut montrer aux utilisateurs un flux d’activité en direct et signaler les connexions suspectes à mesure qu’elles se produisent. L’équipe résiste à monter une plateforme de streaming lourde. Ils placent les événements sur un seul courtier basé sur journal retenu, exécutent un processeur de flux léger pour la logique de risque de connexion, et alimentent un magasin OLAP temps réel qui alimente le flux d’activité. Chaque puits est clé et idempotent, donc une retentative au-moins-une-fois ne compte jamais en double. Quand ils trouvent plus tard un bogue dans la règle de risque, ils rejouent simplement le journal à travers la logique fixe pendant la nuit, parce qu’ils ont gardé une semaine d’historique et n’ont jamais eu besoin d’une seconde base de code par lots.

Grande entreprise. Une banque de détail note chaque transaction de carte pour la fraude dans la fenêtre d’autorisation, joignant le flux de transaction en direct contre un modèle avec état de comportement de compte récent. Le point de contrôle laisse le service de scoring récupérer d’un échec de nœud en secondes sans perdre sa mémoire des dernières minutes. Séparément, la capture de changement de données diffuse les mises à jour depuis la base de données bancaire centrale vers un index de recherche et un service de personnalisation, gardant les deux frais sans interrogation. Les tableaux de bord opérationnels lisent depuis un magasin OLAP temps réel pour que les équipes de risque et d’opérations observent l’entreprise à mesure qu’elle bouge, et le pipeline entier émet la télémétrie de retard et débit décrite au chapitre 9.2.

Gouvernement. Une autorité de transit métropolitain ingère les positions de véhicule et tapotements de tarif pour prédire les arrivées et surveiller l’encombrement en temps réel, alimentant à la fois les applications publiques et un centre d’opérations. Parce que les passagers dans les tunnels téléversent les tapotements en rafales retardées, l’équipe calcule la fréquentation sur le temps d’événement avec des filigranes réglés au retard observé, et journalise tout événement qui arrive après que sa fenêtre se ferme plutôt que de le laisser tomber silencieusement. Chaque événement traité est retenu dans un journal auditable pour que les chiffres de fréquentation rapportés aux organismes de surveillance puissent être reconstruits exactement. Quand une règle de tarif change, ils rejouent la période affectée à travers la logique corrigée et produisent une reformulation défendable.

Argumentaire économique : motivations, retour sur investissement et coût total de possession

Le retour sur les données temps réel vient d’agir pendant que l’action compte encore. La fraude attrapée pendant l’autorisation prévient une perte qu’un lot nocturne ne ferait que rapporter. La personnalisation qui répond dans une session élève la conversion d’une façon que la recommandation de demain ne peut pas. La surveillance opérationnelle qui reflète le présent vous laisse intervenir avant qu’un petit problème ne devienne une panne ou un incident public. Dans chaque cas, la valeur est le delta entre agir maintenant et agir plus tard, et ce delta est ce que vous devriez quantifier quand vous faites valoir votre cas.

Le coût total de possession est plus élevé que le par lots, et l’honnêteté à ce sujet protège votre crédibilité. Vous payez pour l’infrastructure toujours active, pour des ingénieurs qui comprennent le temps d’événement, les filigranes, l’état, et la sémantique de livraison, et pour le test plus difficile et fardeau d’astreinte d’un système qui doit rester sain continuellement plutôt que s’exécuter et s’arrêter. Une architecture streaming d’abord sur un journal retenu abaisse le coût continu en vous épargnant une base de code par lots dupliquée, et choisir au-moins-une-fois avec des puits idempotents évite la dépense d’une machinerie exactement-une-fois de bout en bout. L’erreur la plus coûteuse est de construire du temps réel là où le micro-lot ou par lots suffirait, donc l’argument de coût le plus fort est souvent une décision de ne pas streamer. Cadrez l’argumentaire à la direction autour de décisions spécifiques sensibles à la latence et leur gain mesurable, et soyez également clair sur où rester en par lots économise de l’argent sans perte de valeur.

Anti-patterns et pièges

  • Construire du streaming pour le prestige quand un micro-lot toutes les quelques minutes satisferait le besoin.
  • Calculer sur le temps de traitement, pour que vos chiffres oscillent avec votre infrastructure plutôt que le monde.
  • Ignorer les données tardives et désordonnées jusqu’à ce que la réconciliation échoue en production.
  • Poursuivre littéralement exactement-une-fois partout au lieu d’au-moins-une-fois avec des puits idempotents.
  • État non borné sans expiration, grandissant discrètement jusqu’à ce qu’une tâche manque de mémoire.
  • Aucun point de contrôle, donc un redémarrage perd l’état ou force un replay complet.
  • Interroger les bases de données opérationnelles sur une minuterie au lieu d’utiliser la capture de changement de données.
  • Maintenir une couche par lots Lambda et couche de vitesse avec une logique dupliquée et dérivante.
  • Une fenêtre de rétention trop courte pour rejouer l’historique quand vous trouvez un bogue.
  • Des flux sans métriques de retard, débit, ou fraîcheur, échouant silencieusement.

Modèle de maturité

  • Niveau 1 (Initier) : Tout est par lots, ou quelques tâches de streaming construites à la main s’exécutent réactivement sans surveillance. Les chiffres sont calculés sur le temps de traitement, les données tardives sont ignorées, et un redémarrage perd l’état. Personne ne peut rejouer l’historique pour corriger un bogue, et les problèmes sont découverts quand les chiffres en aval refusent de se réconcilier.
  • Niveau 2 (Développer) : Certaines équipes exécutent des pipelines de streaming centraux sur un courtier basé sur journal avec points de contrôle, et distinguent le temps d’événement du temps de traitement et utilisent des fenêtres de base. La pratique est incohérente d’équipe à équipe : la livraison est au-moins-une-fois mais tous les puits ne sont pas idempotents, la gestion des données tardives est improvisée, et le retard est observé informellement plutôt qu’alerté.
  • Niveau 3 (Standardiser) : Le temps d’événement, les filigranes, et une politique explicite de données tardives sont documentés et appliqués à travers l’organisation. Les puits sont idempotents pour des résultats effectivement-une-fois, l’état a une expiration, et la capture de changement de données alimente les systèmes en aval par convention. Un journal retenu soutient le replay, et le retard, débit, et fraîcheur sont surveillés avec des alertes comme norme à l’échelle de l’organisation plutôt qu’une habitude par équipe.
  • Niveau 4 (Gérer) : Le parc de streaming est mesuré et contrôlé contre des références. Chaque pipeline porte des objectifs de niveau de service pour la latence de bout en bout, le retard de consommateur, le temps de récupération, l’asymétrie de temps d’événement, le taux d’événement tardif, la taille d’état, et le coût par million d’événements, tous suivis contre des cibles convenues et alertant sur régression. La récupération est exercée et chronométrée plutôt que supposée, la marge de contre-pression et la croissance d’état sont observées comme signaux de capacité, et un nouveau flux doit franchir ces métriques avant d’aller en production.
  • Niveau 5 (Orchestrer) : Une architecture streaming d’abord sert à la fois les besoins frais et historiques depuis un journal rejouable unique, et le SQL de streaming, les vues matérialisées, et l’OLAP temps réel rendent les données fraîches largement accessibles. Le retraitement est routinier et testé, la plateforme s’auto-échelonne et se rééquilibre contre la charge et le coût mesurés, et les flux sont retirés, recadrés, ou remplacés sur preuve. Le streaming est intégré à la planification d’affaires et de risque, et chaque flux est observable et auditable de bout en bout à mesure que le paysage de charge et coût change.

Pistes de réflexion

  1. Où dans votre pile le « temps réel » gagne-t-il réellement son coût, et où est-ce un souhait non examiné ?
  2. Quelle est la taille de l’écart entre le temps d’événement et le temps de traitement à travers vos sources, et le mesurez-vous ?
  3. Pourriez-vous effondrer une configuration Lambda par lots-et-vitesse en une seule base de code streaming d’abord, et qu’est-ce qui bloquerait cela ?
  4. Lesquels de vos puits sont véritablement idempotents, et pourriez-vous rejouer en sécurité les données du dernier trimestre à travers la logique corrigée aujourd’hui ?
  5. Quelle est votre politique pour les données qui arrivent après qu’une fenêtre se ferme, et tout le monde en aval la connaît-il ?
  6. Comment la capture de changement de données changerait-elle la façon dont vous gardez la recherche, les caches, et l’analytique synchronisés ?

Points clés à retenir

  • Recourez au streaming seulement quand une décision sensible à la latence le paie ; le par lots et micro-lots sont des valeurs par défaut moins chères.
  • Calculez sur le temps d’événement, et traitez les données tardives et désordonnées comme le problème central, géré avec des fenêtres et filigranes.
  • Préférez la livraison au-moins-une-fois avec des puits idempotents pour des résultats effectivement-une-fois plutôt qu’exactement-une-fois littéral partout.
  • Faites des points de contrôle du traitement avec état, bornez votre état, et surveillez le retard de consommateur comme métrique phare.
  • Utilisez la capture de changement de données pour streamer depuis des bases de données opérationnelles au lieu d’interroger.
  • Favorisez une architecture streaming d’abord sur un journal retenu et rejouable plutôt que de maintenir deux bases de code.
  • Exposez les flux à travers le SQL de streaming, les vues matérialisées, et l’OLAP temps réel, et gardez chaque flux observable et auditable.

Références et lectures complémentaires

  • Tyler Akidau, Slava Chernyak, et Reuven Lax, « Streaming Systems ».
  • Martin Kleppmann, « Designing Data-Intensive Applications ».
  • Nathan Marz et James Warren, « Big Data » (architecture Lambda).
  • Jay Kreps, « Questioning the Lambda Architecture » (O’Reilly Radar).
  • Fabian Hueske et Vasiliki Kalavri, « Stream Processing with Apache Flink ».
  • Ben Stopford, « Designing Event-Driven Systems ».
  • Tyler Akidau et collègues, « The Dataflow Model » (article VLDB sur le fenêtrage et les filigranes).