Skip to content
Cozi fără Kafka: SQS, EventBridge și singurul loc în care chiar am avut nevoie de un stream
← ← Înapoi la Idei Cloud

Cozi fără Kafka: SQS, EventBridge și singurul loc în care chiar am avut nevoie de un stream

Fiecare platformă ajunge la momentul în care o cerere face prea multe. Handler-ul de checkout trimite un email, actualizează indexul de căutare, anunță depozitul, înregistrează un eveniment de analytics și, a, mai și debitează cardul. Durează patru secunde și eșuează dacă oricare dintre cele cinci lucruri e lent. Soluția e o coadă, iar prima propunere e de obicei Kafka, pentru că Kafka e ce folosesc companiile mari și despre Kafka sunt prezentările de la conferințe.

Noi nu rulăm Kafka. Pentru o platformă pe care o operăm pentru un client din SUA, rulăm SQS pentru muncă, EventBridge pentru evenimente și un singur stream Kinesis pentru unicul workload care avea nevoie de ordonare și replay. Iată cum am împărțit, cât costă fiecare și testul pentru a afla dacă ai nevoie de un stream.

Trei probleme diferite

„Coadă” duce mult în spate în majoritatea discuțiilor de arhitectură. Sub cuvânt se ascund trei forme.

Muncă de făcut, o dată, de cineva. Trimite emailul ăsta. Redimensionează imaginea asta. Sincronizează comanda asta cu depozitul. Producătorului nu-i pasă cine o face sau când, doar că se face, și se face o dată. Asta vrea o coadă: un element e livrat unui singur consumator, confirmat când e gata, reîncercat dacă nu, trimis în dead-letter după prea multe eșecuri.

S-a întâmplat ceva, și oricine e interesat ar trebui să afle. S-a plasat o comandă. Un client și-a schimbat emailul. Producătorul nu știe cine ascultă și n-ar trebui să fie nevoit să știe. Asta vrea un event bus: o publicare, mulți abonați, fiecare cu coada lui în spate, adăugați și scoși fără să atingi producătorul.

O istorie ordonată, care se poate reda. Fiecare schimbare a acestui cont, în ordine, pe care un consumator o poate citi din orice punct și reciti după un bug. Asta vrea un stream, și e singura dintre cele trei la care Kafka e unic de bun. E și nevoia cea mai rară.

Majoritatea platformelor au multă primă formă, ceva a doua și una sau zero din a treia. Kafka le face pe toate trei, cu prețul de a rula Kafka, sau de a plăti un Kafka gestionat care pornește de la câteva sute de dolari pe lună înainte să fi trimis un mesaj.

coadă · muncăfă asta o dată, cinevaun consumator o iaack · reîncercare · dead-letterSQS · 11 coziemailuri, redimensionări, sincronizare depozit~2 $ / lună event bus · fapteasta s-a întâmplato publicare, mulți abonațifiecare primește coada luiEventBridge · 1 bus · 9 reguliOrderPlaced, CustomerUpdated, RefundIssued~1 $ / lună stream · istorieordonat per cheie, redabilcitit din orice punctrecitit după un bugKinesis · 1 stream · 1 shardproiecția registrului, și doar ea~15 $ / lună Kafka le face pe toate trei. Kafka gestionat pornește de la câteva sute pe lună; self-hosted pornește de la un operator. Majoritatea platformelor au nevoie de multă primă formă, ceva a doua și una sau zero din a treia.

Cozi: SQS

Unsprezece cozi SQS, una pe fiecare tip de muncă, fiecare cu o coadă dead-letter și un Lambda sau un worker App Runner care o consumă. Handler-ul de checkout care dura patru secunde acum scrie comanda, debitează cardul (singurul lucru care trebuie să fie sincron), publică un eveniment și se întoarce în 400 ms. Tot restul e un consumator de coadă.

Setările care contează, în CDK:

const dlq = new sqs.Queue(this, 'EmailDlq', { retentionPeriod: Duration.days(14) });
const emailQueue = new sqs.Queue(this, 'EmailQueue', {
  visibilityTimeout: Duration.seconds(90),      // > 6 × timeout-ul consumatorului
  deadLetterQueue: { queue: dlq, maxReceiveCount: 5 },
  encryption: sqs.QueueEncryption.SQS_MANAGED,
});
new lambda.EventSourceMapping(this, 'EmailConsumer', {
  target: emailFn, eventSourceArn: emailQueue.queueArn,
  batchSize: 10, reportBatchItemFailures: true,     // succes parțial pe batch
  maxConcurrency: 20,                                // protejează furnizorul de email
});
new cloudwatch.Alarm(this, 'EmailDlqAlarm', {
  metric: dlq.metricApproximateNumberOfMessagesVisible(), threshold: 1, evaluationPeriods: 1,
});

Trei lucruri care erau greșite în prima versiune și sunt corecte acum. Visibility timeout-ul era egal cu timeout-ul funcției, deci un mesaj lent era relivrat în timp ce era încă procesat, și emailurile plecau de două ori; acum e de șase ori timeout-ul funcției, cum spune documentația pe care n-o citește nimeni. reportBatchItemFailures era oprit, deci un mesaj prost dintr-un batch de zece le pica pe toate zece, și nouă emailuri bune erau reîncercate de cinci ori fiecare înainte ca batch-ul să ajungă în dead-letter. Și coada dead-letter n-avea alarmă, deci mesajele stăteau în ea o săptămână înainte să se uite cineva; acum sună la unu.

FIFO sau standard? Standard, peste tot în afară de sincronizarea cu depozitul, care are nevoie ca „anulează” să ajungă după „creează” pentru aceeași comandă. FIFO cu id-ul comenzii ca message group dă asta, cam la același preț, cu un plafon de throughput per grup de care suntem departe. Nu pune FIFO implicit; adaugă un contract de deduplicare și ordonare pe care majoritatea muncii nu-l vrea și care face ca un mesaj blocat să blocheze tot ce e în spatele lui în grup.

Evenimente: EventBridge

Un bus custom. Producătorii pun evenimente cu un detail-type și o sursă; regulile potrivesc și rutează către ținte, care sunt aproape întotdeauna o coadă SQS deținută de serviciul consumator, cu un Lambda în spate. Producătorul lui OrderPlaced habar n-are că cinci servicii se abonează la el, și când apare al șaselea trimestrul viitor, e o regulă nouă și o coadă nouă, fără nicio schimbare la 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 de zile, pentru „ce s-a întâmplat?”
});

Regula de arhivă e trucul ieftin care îți dă mare parte din ce vor oamenii de la un stream: o înregistrare căutabilă, ordonată în timp, a fiecărui eveniment, în același store de loguri ca tot restul, fără un stream. Nu poate reda într-un consumator, dar poate răspunde la „am emis OrderPlaced pentru comanda 4412?” într-o singură interogare, care e întrebarea care chiar se pune.

Contractul EventBridge e cel-puțin-o-dată, neordonat, cu o limită de payload de 256 KB. Fiecare consumator e idempotent, cu cheie pe id-ul evenimentului, și orice eveniment mai mare de câțiva KB poartă un pointer către S3, nu payload-ul. Cele două reguli acoperă fiecare problemă pe care am avut-o cu el.

Stream-ul: Kinesis, o singură dată

Registrul. Fiecare mișcare financiară de pe platformă, în ordine per cont, proiectată în solduri și rapoarte de un consumator care trebuie să poată fi reconstruit de la zero dacă se găsește un bug de proiecție. Asta e problema în formă de stream: ordonare per cheie și replay din orice punct.

Un stream Kinesis, un shard, retenție de 7 zile, un consumator Lambda cu checkpoint. Când s-a găsit un bug de proiecție în luna a patra, soluția a fost corectarea consumatorului, resetarea checkpoint-ului la începutul retenției și lăsarea lui să recitească trei zile de evenimente într-o tabelă proaspătă. SQS nu poate face asta; un mesaj consumat e dus. EventBridge nu poate face asta; arhiva e un log, nu un cursor.

checkout · 400 msscrie comandadebitează cardul (sincron)publică 1 eveniment OrderPlaced EventBridge9 reguli SQS · email → LambdaSQS · index de căutare → LambdaSQS FIFO · depozit, după id comandăSQS · analytics → LambdaCloudWatch Logs · arhivă, 30 z scriere registru Kinesis · 1 shard proiecție de solduri · redabilă fiecare consumatoridempotent pe id-ul evenimentuluiDLQ propriu, alarmă proprie singurul loccare are nevoie de ordine + replay

Testul pentru „ai nevoie de un stream?”

Întreabă: dacă un consumator a avut un bug marțea trecută, trebuie să-i redai mesajele de marți în ordine? Dacă răspunsul sincer e „am rerula un job batch pe baza de date”, ai nevoie de o coadă și o bază de date, nu de un stream. Dacă răspunsul e „da, și baza de date nu are istoricul”, ai nevoie de un stream, pentru consumatorul ăla. Un stream, pentru un consumator, nu e un motiv să muți toată platforma pe Kafka.

Cât costă

Lunar Operațiuni
SQS, 11 cozi + 11 DLQ-uri, ~4 M mesaje ~2 $ Zero. Alarme pe DLQ-uri
EventBridge, 1 bus, 9 reguli, ~1 M evenimente ~1 $ Zero. Regulă de arhivă pentru „ce s-a întâmplat”
Kinesis, 1 shard, retenție 7 zile ~15 $ Monitorizarea checkpoint-ului; numărul de shard-uri dacă va conta vreodată
Total ~18 $
Kafka gestionat, cel mai mic cluster util 300–600 $ Topic-uri, partiții, consumer group-uri, retenție, o versiune de broker de urmărit

Optsprezece dolari și niciun broker. Platforma gestionează câteva milioane de mesaje pe lună; SQS ar gestiona câteva miliarde cu aceeași formă, cu factura scalând liniar și nimic de re-arhitecturat.

Versiunea scurtă

Împarte „coadă” în trei probleme. Munca merge în SQS cu o coadă dead-letter și o alarmă. Faptele merg în EventBridge cu o coadă per abonat și o regulă de arhivă. Istoria, dacă ai cu adevărat un consumator care trebuie să redea în ordine, merge într-un singur stream Kinesis, pentru consumatorul ăla. Kafka e răspunsul corect când ai multe din a treia problemă, și e în regulă să nu ai problema asta.

Dacă checkout-ul tău durează patru secunde pentru că face cinci lucruri, putem scoate patru dintre ele de pe calea cererii într-o săptămână, pentru vreo doi dolari pe lună.