Skip to content
Des files sans Kafka : SQS, EventBridge, et le seul endroit où nous avons vraiment eu besoin d'un stream
← ← Retour aux Réflexions Cloud

Des files sans Kafka : SQS, EventBridge, et le seul endroit où nous avons vraiment eu besoin d'un stream

Chaque plateforme atteint le moment où une requête en fait trop. Le handler de checkout envoie un email, met à jour l'index de recherche, notifie l'entrepôt, enregistre un événement d'analytique et, ah oui, débite aussi la carte. Ça prend quatre secondes et ça échoue si l'une des cinq choses est lente. Le correctif est une file, et la première proposition est généralement Kafka, parce que Kafka est ce que les grandes entreprises utilisent et que les conférences parlent de Kafka.

Nous ne faisons pas tourner Kafka. Pour une plateforme que nous exploitons pour un client américain, nous utilisons SQS pour le travail, EventBridge pour les événements, et un seul stream Kinesis pour l'unique workload qui avait besoin d'ordre et de relecture. Voici comment nous l'avons découpé, ce que coûte chaque partie, et le test pour savoir si vous avez besoin d'un stream.

Trois problèmes différents

« File » porte beaucoup de sens dans la plupart des conversations d'architecture. Trois formes se cachent dessous.

Du travail à faire, une fois, par quelqu'un. Envoyer cet email. Redimensionner cette image. Synchroniser cette commande avec l'entrepôt. Le producteur ne se soucie pas de qui le fait ni quand, seulement que ce soit fait, et fait une fois. Cela veut une file : un élément est livré à un consommateur, acquitté quand c'est fait, réessayé sinon, mis en dead-letter après trop d'échecs.

Quelque chose s'est produit, et tous les intéressés devraient le savoir. Une commande a été passée. Un client a changé son email. Le producteur ne sait pas qui écoute et ne devrait pas avoir à le savoir. Cela veut un bus d'événements : une publication, de nombreux abonnés, chacun avec sa propre file derrière, ajoutés et retirés sans toucher au producteur.

Un historique ordonné et rejouable. Chaque changement sur ce compte, dans l'ordre, qu'un consommateur peut lire depuis n'importe quel point et relire après un bug. Cela veut un stream, et c'est la seule des trois formes où Kafka est unique. C'est aussi le besoin le plus rare.

La plupart des plateformes ont beaucoup de la première, un peu de la deuxième, et une ou zéro de la troisième. Kafka fait les trois, au prix de faire tourner Kafka, ou de payer un Kafka géré qui commence à quelques centaines de dollars par mois avant d'avoir envoyé un message.

file · travailfais ça une fois, quelqu'unun consommateur le prendack · réessai · dead-letterSQS · 11 filesemails, redimensionnements, sync entrepôt~2 $ / mois bus d'événements · faitsceci s'est produitune publication, de nombreux abonnéschacun reçoit sa propre fileEventBridge · 1 bus · 9 règlesOrderPlaced, CustomerUpdated, RefundIssued~1 $ / mois stream · historiqueordonné par clé, rejouablelu depuis n'importe quel pointrelu après un bugKinesis · 1 stream · 1 shardla projection du grand livre, et seulement elle~15 $ / mois Kafka fait les trois. Kafka géré commence à quelques centaines par mois ; auto-hébergé commence à un opérateur. La plupart des plateformes ont besoin de beaucoup de la première, d'un peu de la deuxième, et d'une ou zéro de la troisième.

Les files : SQS

Onze files SQS, une par type de travail, chacune avec une file dead-letter et une Lambda ou un worker App Runner qui la consomme. Le handler de checkout qui prenait quatre secondes écrit maintenant la commande, débite la carte (la seule chose qui doit être synchrone), publie un événement et répond en 400 ms. Tout le reste est un consommateur de file.

Les réglages qui comptent, en CDK :

const dlq = new sqs.Queue(this, 'EmailDlq', { retentionPeriod: Duration.days(14) });
const emailQueue = new sqs.Queue(this, 'EmailQueue', {
  visibilityTimeout: Duration.seconds(90),      // > 6 × le délai du consommateur
  deadLetterQueue: { queue: dlq, maxReceiveCount: 5 },
  encryption: sqs.QueueEncryption.SQS_MANAGED,
});
new lambda.EventSourceMapping(this, 'EmailConsumer', {
  target: emailFn, eventSourceArn: emailQueue.queueArn,
  batchSize: 10, reportBatchItemFailures: true,     // succès partiel par lot
  maxConcurrency: 20,                                // protège le fournisseur d'email
});
new cloudwatch.Alarm(this, 'EmailDlqAlarm', {
  metric: dlq.metricApproximateNumberOfMessagesVisible(), threshold: 1, evaluationPeriods: 1,
});

Trois choses qui étaient fausses dans la première version et qui sont justes maintenant. Le délai de visibilité était égal au délai de la fonction, donc un message lent était relivré alors qu'il était encore en traitement, et les emails partaient deux fois ; c'est maintenant six fois le délai de la fonction, comme le dit la documentation que personne ne lit. reportBatchItemFailures était désactivé, donc un mauvais message dans un lot de dix faisait échouer les dix, et neuf bons emails étaient réessayés cinq fois chacun avant que le lot atteigne la dead-letter. Et la file dead-letter n'avait pas d'alarme, donc les messages y restaient une semaine avant que quelqu'un regarde ; elle alerte maintenant à un.

FIFO ou standard ? Standard, partout sauf pour la synchronisation avec l'entrepôt, qui a besoin qu'« annuler » arrive après « créer » pour la même commande. FIFO avec l'id de commande comme groupe de messages donne ça, à peu près au même prix, avec un plafond de débit par groupe dont nous sommes loin. Ne mettez pas FIFO par défaut ; il ajoute un contrat de déduplication et d'ordre que la plupart du travail ne veut pas, et qui fait qu'un message bloqué bloque tout ce qui le suit dans son groupe.

Les événements : EventBridge

Un bus personnalisé. Les producteurs déposent des événements avec un detail-type et une source ; les règles filtrent et routent vers des cibles, qui sont presque toujours une file SQS appartenant au service consommateur, avec une Lambda derrière. Le producteur de OrderPlaced n'a aucune idée que cinq services y sont abonnés, et quand le sixième apparaît le trimestre prochain, c'est une nouvelle règle et une nouvelle file, sans changement au checkout.

const bus = new events.EventBus(this, 'PlatformBus');
new events.Rule(this, 'OrderPlacedToSearch', {
  eventBus: bus,
  eventPattern: { source: ['platform.orders'], detailType: ['OrderPlaced'] },
  targets: [new targets.SqsQueue(searchIndexQueue)],
});
new events.Rule(this, 'AllEventsToArchive', {
  eventBus: bus, eventPattern: { source: [{ prefix: 'platform.' }] },
  targets: [new targets.CloudWatchLogGroup(eventArchive)],   // 30 jours, pour « que s'est-il passé ? »
});

La règle d'archive est l'astuce bon marché qui vous donne l'essentiel de ce que les gens attendent d'un stream : un enregistrement consultable, ordonné dans le temps, de chaque événement, dans le même magasin de logs que tout le reste, sans stream. Elle ne peut pas rejouer dans un consommateur, mais elle peut répondre à « avons-nous émis OrderPlaced pour la commande 4412 ? » en une requête, ce qui est la question réellement posée.

Le contrat d'EventBridge est au-moins-une-fois, non ordonné, avec une limite de charge utile de 256 Ko. Chaque consommateur est idempotent, indexé sur l'id d'événement, et tout événement de plus de quelques Ko porte un pointeur vers S3, pas la charge utile. Ces deux règles couvrent chaque problème que nous avons eu avec lui.

Le stream : Kinesis, une seule fois

Le grand livre. Chaque mouvement financier sur la plateforme, dans l'ordre par compte, projeté en soldes et rapports par un consommateur qui doit pouvoir être reconstruit de zéro si un bug de projection est trouvé. C'est le problème en forme de stream : ordre par clé, et relecture depuis n'importe quel point.

Un stream Kinesis, un shard, 7 jours de rétention, un consommateur Lambda avec un checkpoint. Quand un bug de projection a été trouvé au quatrième mois, le correctif a été de corriger le consommateur, de remettre le checkpoint au début de la rétention et de le laisser relire trois jours d'événements dans une table neuve. SQS ne peut pas faire ça ; un message consommé est parti. EventBridge ne peut pas faire ça ; l'archive est un log, pas un curseur.

checkout · 400 msécrire la commandedébiter la carte (sync)publier 1 événement OrderPlaced EventBridge9 règles SQS · email → LambdaSQS · index de recherche → LambdaSQS FIFO · entrepôt, par id de commandeSQS · analytique → LambdaCloudWatch Logs · archive, 30 j écriture grand livre Kinesis · 1 shard projection des soldes · rejouable chaque consommateuridempotent sur l'id d'événementsa DLQ, son alarme le seul endroitqui a besoin d'ordre + relecture

Le test pour « avez-vous besoin d'un stream ? »

Demandez : si un consommateur avait un bug mardi dernier, devez-vous lui refournir les messages de mardi dans l'ordre ? Si la réponse honnête est « nous relancerions un job batch sur la base de données », vous avez besoin d'une file et d'une base de données, pas d'un stream. Si la réponse est « oui, et la base de données n'a pas l'historique », vous avez besoin d'un stream, pour ce consommateur. Un stream, pour un consommateur, n'est pas une raison de faire migrer toute la plateforme vers Kafka.

Ce que ça coûte

Mensuel Exploitation
SQS, 11 files + 11 DLQ, ~4 M de messages ~2 $ Zéro. Des alarmes sur les DLQ
EventBridge, 1 bus, 9 règles, ~1 M d'événements ~1 $ Zéro. Règle d'archive pour « que s'est-il passé »
Kinesis, 1 shard, 7 jours de rétention ~15 $ Surveillance du checkpoint ; nombre de shards si ça compte un jour
Total ~18 $
Kafka géré, plus petit cluster utile 300–600 $ Topics, partitions, groupes de consommateurs, rétention, une version de broker à suivre

Dix-huit dollars et aucun broker. La plateforme traite quelques millions de messages par mois ; SQS en traiterait quelques milliards avec la même forme, la facture augmentant linéairement et rien à ré-architecturer.

La version courte

Découpez « file » en trois problèmes. Le travail va dans SQS avec une file dead-letter et une alarme. Les faits vont dans EventBridge avec une file par abonné et une règle d'archive. L'historique, si vous avez vraiment un consommateur qui doit rejouer dans l'ordre, va dans un stream Kinesis, pour ce consommateur. Kafka est la bonne réponse quand vous avez beaucoup du troisième problème, et il est très bien de ne pas avoir ce problème.

Si votre checkout prend quatre secondes parce qu'il fait cinq choses, nous pouvons en sortir quatre du chemin de la requête en une semaine, pour environ deux dollars par mois.