Node.js-taustakehityksen bootcamp · Oppitunti

Pub/Sub, Streams ja nopeusrajoitus Redisillä

Lähetä tapahtumia, rakenna pysyviä Redis Streameja ja toteuta token bucket -nopeusrajoittimia atomisesti.

Oppitunti 3/413 vaihetta

Pub/Sub, Streams ja nopeusrajoitus Redisillä on ilmainen Node.js-taustakehityksen bootcamp-oppitunti CoddyKitissä. Tämä on oppitunti 3/4. Voit lukea koko oppitunnin alta ilmaiseksi ja harjoitella sen jälkeen käytännössä selaimessa sisäänrakennetulla koodieditorilla ja ympäri vuorokauden käytettävissä olevan tekoälytuutorin avulla. Oppitunti kuuluu Node.js-taustakehityksen bootcamp-oppimispolkuun, ja edistymisesi synkronoituu verkon ja CoddyKit-sovelluksen välillä. Node.js-taustakehityksen bootcamp-kurssilla on yhteensä 4 oppituntia.

Kolme viestintäprimitiiviä yhdellä palvelimella

Redis on enemmän kuin avain-arvo-välimuisti. Node.js-taustajärjestelmässä se toimii myös kevyenä viestinvälittäjänä ja koordinaatiomoottorina. Tässä oppitunnissa opitte kolme tuotantokäyttöön sopivaa mallia, jotka kuuluvat jokaiseen Redis-asennukseen:

  • Pub/Sub — viestien välitön, ilman toimitustakuuta tapahtuva lähetys kaikille yhdistetyille kuuntelijoille.
  • Streams — vain lisäämiseen tarkoitettu, kestävä loki kuluttajaryhmineen ja uudelleenluvulla.
  • Rate limiting — atomiset laskurit, joilla väärin käyttäytyviä asiakkaita rajoitetaan.

Niitä yhdistävä keskeinen ajatus on, että yhdellä Redis-kutsukierroksella voidaan tehdä työ, joka muuten vaatisi tietokannan, jonopalvelimen ja lukkopalvelun. Käytämme koko ajan ioredis-asiakasohjelmaa, joka on Node-taustajärjestelmien de facto -standardi, koska se tukee putkitusta, Lua-komentosarjoja ja klusteritilaa.

Pub/Sub: lähetys kaikille tilaajille

Redis Pub/Sub mahdollistaa sen, että yksi prosessi julkaisee viestin kanavalle komennolla PUBLISH ja jokainen SUBSCRIBE-tilassa oleva asiakas vastaanottaa sen välittömästi. Se sopii erinomaisesti välimuistin mitätöinnin levittämiseen tai reaaliaikaisten päivitysten lähettämiseen WebSocket-palvelimille.

Kriittinen huomio: tilaustilassa oleva yhteys ei voi suorittaa tavallisia komentoja. Teidän on luotava erillinen tilaajayhteys, joka on omistettu tilaamiseen ja erotettu yhteydestä, jota käytätte komennoille GET/SET. Alla sub vain kuuntelee; julkaisemiseen käytettäisiin toista asiakasohjelmaa.

const Redis = require('ioredis');

const sub = new Redis();   // dedicated subscriber connection
const pub = new Redis();   // separate connection for publishing

sub.subscribe('cache:invalidate', (err, count) => {
  if (err) throw err;
  console.log('Subscribed to ' + count + ' channel(s)');
});

sub.on('message', (channel, message) => {
  console.log('[' + channel + '] ' + message);
});

// Another part of the app publishes an event
pub.publish('cache:invalidate', JSON.stringify({ key: 'user:42' }));

Kuvioihin perustuvat tilaukset ja ilman toimitustakuuta lähettämisen rajoitus

Tarkkojen kanavien lisäksi Redis tukee PSUBSCRIBE-komentoa glob-kuvioilla. Kuvioon order:* tilaamalla vastaanotetaan viestejä kanavista order:created, order:shipped ja niin edelleen. Callback on pmessage, ja se sisältää täsmäävän kuvion.

Keskeinen rajoitus, joka on tärkeää ymmärtää: Pub/Sub ei säilytä viestejä. Jos julkaisun hetkellä yksikään tilaaja ei ole yhteydessä, viesti katoaa lopullisesti — puskurointia, kuittausta tai uudelleenlukua ei ole. Jos tilaaja kaatuu ja muodostaa yhteyden uudelleen, se menettää kaiken yhteyden katkon aikana tapahtuneen.

  • Käyttäkää Pub/Subia hetkellisiin tapahtumiin, joissa viestin menettäminen on hyväksyttävää.
  • Tapahtumiin, joita ei saa menettää, tarvitaan Streams (seuraava osio).
const Redis = require('ioredis');
const sub = new Redis();

sub.psubscribe('order:*', (err, count) => {
  console.log('Pattern subscriptions: ' + count);
});

sub.on('pmessage', (pattern, channel, message) => {
  console.log('matched ' + pattern + ' on ' + channel + ': ' + message);
});

Streams: kestävä ja uudelleenluettava loki

Redis Stream on vain lisäämiseen tarkoitettu loki. Jokainen merkintä saa kasvavan tunnisteen, kuten 1718000000000-0 (millisekunnit-järjestysnumero). Toisin kuin Pub/Subissa, merkinnät tallennetaan, kunnes niitä karsitaan, joten myöhään liittyvät tai uudelleenkäynnistyvät kuluttajat voivat lukea historian uudelleen.

Merkintä lisätään komennolla XADD. Erityinen tunniste * käskee Redisiä luomaan seuraavan tunnisteen automaattisesti. Tunnisteen jälkeen välitetään kenttä-arvo-pareja aivan kuten hash-rakenteessa.

const Redis = require('ioredis');
const redis = new Redis();

async function publishOrder() {
  // XADD key * field value field value ...
  const id = await redis.xadd(
    'stream:orders', '*',
    'orderId', '42',
    'amount', '99.90',
    'status', 'created'
  );
  console.log('Appended entry with ID ' + id);

  // Read the latest 5 entries (newest last)
  const entries = await redis.xrange('stream:orders', '-', '+', 'COUNT', 5);
  console.log(JSON.stringify(entries, null, 2));
}

publishOrder();

Kuluttajaryhmät: skaalaus ilman kaksoiskäsittelyä

Streamsin varsinainen voima ovat kuluttajaryhmät. Ryhmä mahdollistaa yhden streamin kuorman jakamisen useiden työntekijäprosessien kesken: jokainen merkintä toimitetaan täsmälleen yhdelle ryhmän kuluttajalle, mikä mahdollistaa horisontaalisen skaalauksen.

Luokaa ryhmä kerran komennolla XGROUP CREATE. MKSTREAM-valitsin luo streamin, jos sitä ei vielä ole. Aloitustunniste $ tarkoittaa, että toimitetaan vain ryhmän luomisen jälkeen lisätyt merkinnät; arvolla 0 kulutetaan aivan alusta alkaen.

const Redis = require('ioredis');
const redis = new Redis();

async function setup() {
  try {
    await redis.xgroup(
      'CREATE', 'stream:orders', 'order-workers', '$', 'MKSTREAM'
    );
    console.log('Group created');
  } catch (e) {
    if (e.message.includes('BUSYGROUP')) {
      console.log('Group already exists, continuing');
    } else {
      throw e;
    }
  }
}

setup();

Stream-merkintöjen lukeminen ja kuittaaminen

Jokainen työntekijä lukee merkintöjä komennolla XREADGROUP ja välittää ryhmän nimen sekä yksilöllisen consumer-nimen. Erityinen tunniste > tarkoittaa: "anna merkinnät, joita ei ole toimitettu kenellekään tämän ryhmän kuluttajalle".

Kuittaamaton merkintä pysyy ryhmän odottavien merkintöjen luettelossa (PEL). Kun käsittely on valmis, kutsutaan XACK-komentoa, jotta Redis tietää merkinnän käsitellyksi turvallisesti. Jos työntekijä kaatuu ennen XACK-komentoa, merkintä jää odottamaan ja se voidaan ottaa uudelleen käsittelyyn — tämän ansiosta Streams tarjoaa vähintään kerran -toimituksen ja kestävyyden.

const Redis = require('ioredis');
const redis = new Redis();

async function consume() {
  const res = await redis.xreadgroup(
    'GROUP', 'order-workers', 'worker-1',
    'COUNT', 10, 'BLOCK', 5000,
    'STREAMS', 'stream:orders', '>'
  );
  if (!res) return; // BLOCK timed out with no new entries

  for (const [, entries] of res) {
    for (const [id, fields] of entries) {
      console.log('processing ' + id, fields);
      // ... do real work here ...
      await redis.xack('stream:orders', 'order-workers', id);
    }
  }
}

consume();

Jumiutuneiden viestien palauttaminen ja karsiminen

Mitä tapahtuu, jos worker-1 kaatuu kesken käsittelyn? Sen merkinnät pysyvät PEL-luettelossa ikuisesti, ellei niitä oteta uudelleen käsittelyyn. Käyttäkää komentoa XAUTOCLAIM (Redis 6.2+), jotta kynnysajan käyttämättöminä olleet merkinnät voidaan siirtää toimivalle kuluttajalle:

  • XAUTOCLAIM stream group consumer min-idle-time start — hakee merkinnät, jotka ovat olleet käyttämättöminä yli min-idle-time millisekunnin ajan.
  • Tarkistakaa käsittelemättömät työt komennolla XPENDING ennen niiden palauttamista.

Streams kasvaa rajatta, joten rajoittakaa muistin käyttöä rajoitetulla XADD-komennolla: XADD key MAXLEN ~ 10000 * .... Merkki ~ tarkoittaa "likimääräisesti", minkä ansiosta Redis voi karsia merkintöjä tehokkaasti kokonaisina makrosolmuina tarkan lukumäärän sijaan.

const Redis = require('ioredis');
const redis = new Redis();

async function recover() {
  // Reclaim entries idle > 30s, hand them to worker-2
  const [cursor, claimed] = await redis.xautoclaim(
    'stream:orders', 'order-workers', 'worker-2',
    30000, '0', 'COUNT', 25
  );
  console.log('Reclaimed ' + claimed.length + ' entries; next cursor ' + cursor);

  // Append with an approximate cap to bound memory
  await redis.xadd('stream:orders', 'MAXLEN', '~', 10000, '*', 'orderId', '99');
}

recover();

Miksi naiivi nopeusrajoitus ei toimi

Siirrytään seuraavaksi rajoittamiseen. Ensimmäinen ajatus on hakea laskuri komennolla GET, tarkistaa se Nodessa ja tallentaa se sitten takaisin komennolla SET. Tämä on klassinen kilpailutilanne: GET- ja SET-komentojen välissä toinen samanaikainen pyyntö lukee saman vanhentuneen arvon, ja molemmat luulevat olevansa rajan alapuolella. Kuormituksen alaisena päästätte läpi huomattavasti enemmän pyyntöjä kuin sallittu.

Ratkaisu on atomisuus. Redis suorittaa jokaisen komennon (ja jokaisen Lua-komentosarjan) yksisäikeisesti ja jakamattomana. Yksinkertainen kiinteän aikaikkunan rajoitin käyttää komentoja INCR ja EXPIRE, jolloin lue-muokkaa-kirjoita-toiminto tapahtuu palvelimella ilman katkosta. Alla oleva pelkkä JavaScript-esimerkki havainnollistaa kilpailutilannetta ennen kuin siirrämme sen Redisiin.

// Demonstrates WHY check-then-set races. Two 'requests' interleave.
let counter = 0;
const LIMIT = 3;

function tryRequest(name) {
  const current = counter;        // read
  if (current < LIMIT) {
    // imagine an await here: another request runs before we write
    counter = current + 1;        // write (stale!)
    return name + ': allowed (' + counter + ')';
  }
  return name + ': blocked';
}

// Both read 0 before either writes -> over-admission
const a = tryRequest('reqA');
const b = tryRequest('reqB');
console.log(a);
console.log(b);
console.log('Final counter: ' + counter);

Kiinteän aikaikkunan rajoitin komennoilla INCR + EXPIRE

Yksinkertaisin oikein toimiva rajoitin: laskurin avaimessa yhdistetään asiakas ja aikaikkuna, laskuria kasvatetaan komennolla INCR ja TTL asetetaan ensimmäisellä kerralla, kun aikaikkuna avautuu. Koska INCR palauttaa uuden arvon atomisesti, kilpailutilannetta ei synny.

Kiinteiden aikaikkunoiden heikkous on rajan kohdalla tapahtuva purske: asiakas voi lähettää koko kiintiön hetkellä 0:59 ja uudelleen hetkellä 1:00, jolloin nopeus hetkellisesti kaksinkertaistuu. Tämä on hyväksyttävää monissa rajapinnoissa, mutta token bucket -menetelmä (seuraavaksi) tasoittaa tämän.

const Redis = require('ioredis');
const redis = new Redis();

async function allow(userId, limit = 100, windowSec = 60) {
  const key = 'rl:' + userId + ':' + Math.floor(Date.now() / 1000 / windowSec);
  const count = await redis.incr(key);
  if (count === 1) {
    await redis.expire(key, windowSec); // set TTL only on first hit
  }
  return count <= limit;
}

allow('user:42').then((ok) => {
  console.log(ok ? 'request allowed' : '429 Too Many Requests');
});

Atominen token bucket Lua-komentosarjalla

Token bucket antaa kullekin asiakkaalle säiliön, jonka kapasiteetti on C ja joka täyttyy nopeudella R tokenia sekunnissa. Jokainen pyyntö kuluttaa yhden tokenin; jos säiliö on tyhjä, pyyntö hylätään. Menetelmä sallii hallitut purskeet ja samalla ylläpitää tasaista keskimääräistä nopeutta.

Täyttö ja kulutus on tehtävä yhtenä atomisena vaiheena, joten välitämme ne Lua-komentosarjana komennolla EVAL. Redis suorittaa koko komentosarjan ilman muiden komentojen lomittumista. Tallennamme hash-rakenteeseen kaksi kenttää — nykyiset tokens-tokenit ja viimeisimmän täytön ajan ts — ja laskemme laiskasti, kuinka monta tokenia on syntynyt edellisen kutsun jälkeen.

const Redis = require('ioredis');
const redis = new Redis();

const LUA = [
  "local cap = tonumber(ARGV[1])",
  "local refill = tonumber(ARGV[2])",
  "local now = tonumber(ARGV[3])",
  "local cost = tonumber(ARGV[4])",
  "local b = redis.call('HMGET', KEYS[1], 'tokens', 'ts')",
  "local tokens = tonumber(b[1])",
  "local ts = tonumber(b[2])",
  "if tokens == nil then tokens = cap; ts = now end",
  "local delta = math.max(0, now - ts)",
  "tokens = math.min(cap, tokens + delta * refill)",
  "local allowed = 0",
  "if tokens >= cost then allowed = 1; tokens = tokens - cost end",
  "redis.call('HMSET', KEYS[1], 'tokens', tokens, 'ts', now)",
  "redis.call('EXPIRE', KEYS[1], 3600)",
  "return allowed"
].join('\n');

async function take(userId) {
  const now = Date.now() / 1000;
  // capacity 10, refill 1 token/sec, cost 1
  const allowed = await redis.eval(LUA, 1, 'tb:' + userId, 10, 1, now, 1);
  return allowed === 1;
}

take('user:42').then((ok) => console.log(ok ? 'allowed' : 'throttled'));

Rajoittimen liittäminen Express-väliohjelmistoon

Oikeassa Node-taustajärjestelmässä rajoitin sijaitsee middleware-välikerroksessa, joten jokainen reitti on suojattu. Hylkäyksen yhteydessä palautetaan HTTP 429 sekä kohteliaisuutena Retry-After-otsake, joka kertoo asiakkaalle, milloin sen kannattaa yrittää uudelleen.

Muistettavat parhaat käytännöt:

  • Käyttäkää avaimena pysyvää tunnistetta (API-avainta tai todennetun käyttäjän tunnistetta), ei pelkkää IP-osoitetta — NAT-verkkojen takana IP-osoitteet ovat jaettuja.
  • Päättäkää tarkoituksella, avataanko vai suljetaanko toiminta virhetilanteessa: jos Redis ei ole tavoitettavissa, päättäkää, sallitaanko liikenne (saatavuus) vai estetäänkö se (suojaus). Tehkää valinta tietoisesti, älkää jättäkö sitä sattuman varaan.
  • Palauttakaa nopeusrajoituksen otsakkeet, jotta asianmukaisesti toimivat asiakkaat voivat rajoittaa liikennettä itse.
function rateLimit(redis, take) {
  return async (req, res, next) => {
    const id = req.user?.id || req.ip;
    try {
      if (await take(id)) return next();
      res.set('Retry-After', '1');
      return res.status(429).json({ error: 'Too Many Requests' });
    } catch (err) {
      // Redis down: fail OPEN here (prioritize availability)
      console.error('rate limiter degraded', err.message);
      return next();
    }
  };
}

module.exports = { rateLimit };

Pikatarkistus: oikean primitiivin valinta

Rakennatte tilausten käsittelyputkea. Työntekijäprosessit voivat kaatua ja käynnistyä uudelleen, eikä yksikään tilaustapahtuma saa koskaan kadota; jokainen tapahtuma on käsiteltävä täsmälleen yhden työntekijän toimesta, ja kaatuneen työn käsittely on yritettävä automaattisesti uudelleen. Mikä Redis-ominaisuus sopii tähän?

Kertaus: valitkaa takuun mukainen työkalu

Teillä on nyt kolme Redis-viestintä- ja hallintamallia ja ennen kaikkea kyky valita niiden välillä:

  • Pub/Sub — välitön fan-out, mutta viestit ovat tilapäisiä. Käyttäkää erillistä tilaajayhteyttä. Sopii erinomaisesti välimuistin mitätöintiin ja reaaliaikaisiin ilmoituksiin, joissa yksittäisen viestin menettäminen ei haittaa.
  • Streams — pysyvä ja uudelleen toistettavissa oleva loki. Consumer group -ryhmät takaavat toimituksen täsmälleen yhdelle N:stä työntekijästä; XACK sekä Pending Entries List ja XAUTOCLAIM mahdollistavat vähintään kerran tapahtuvan käsittelyn ja palautumisen kaatumisista. Rajoittakaa kasvua käyttämällä komentoa MAXLEN ~.
  • Nopeusrajoitus — älkää koskaan tehkö tarkistusta ja asetusta erillisinä vaiheina sovelluskoodissa, koska niiden välille syntyy kilpailutilanne. Käyttäkää kiinteille aikajaksoille atomista yhdistelmää INCR+EXPIRE tai tasaisten purskeiden hallintaan EVAL-komennolla toteutettua Lua-tokenämpäriä. Käärikää ratkaisu middleware-välikerrokseen, palauttakaa 429 ja Retry-After ja päättäkää tietoisesti, avataanko vai suljetaanko toiminta virhetilanteessa.

Yhdistävä periaate on tämä: antakaa Redisin tehdä atominen työ yhdellä edestakaisella yhteydellä ja valitkaa primitiivi sen kestävyystakuun perusteella, jota käyttötapauksenne todella edellyttää.

Aloita maksutta

Opi JavaScript tekoälytuutorin avulla — ilmaiseksi

Kirjoita ja suorita oikeaa koodia selaimessa, saa välitöntä apua tekoälytuutorilta ympäri vuorokauden ja jatka siitä, mihin jäit, verkossa tai sovelluksessa.

Kurssit
22
Oppitunnit
92

Usein kysytyt kysymykset

Onko oppitunti ”Pub/Sub, Streams ja nopeusrajoitus Redisillä” ilmainen?

Kyllä – oppitunnin ”Pub/Sub, Streams ja nopeusrajoitus Redisillä” koko tekstin voi lukea täällä verkossa ilmaiseksi. Jos haluat harjoitella interaktiivisesti sisäänrakennetulla koodieditorilla ja ympäri vuorokauden käytettävissä olevan tekoälytuutorin avulla sekä avata koko Node.js-taustakehityksen bootcamp-kurssin, päivitä CoddyKit PROhon. Node.js-taustakehityksen bootcamp-kurssilla on yhteensä 4 oppituntia.

Mitä opin oppitunnilla ”Pub/Sub, Streams ja nopeusrajoitus Redisillä”?

Lähetä tapahtumia, rakenna pysyviä Redis Streameja ja toteuta token bucket -nopeusrajoittimia atomisesti. Harjoittelet Node.js-taustakehityksen bootcamp-aihetta koodilla, jonka suoritat suoraan selaimessa. Ympäri vuorokauden käytettävissä oleva tekoälytuutori vastaa kysymyksiisi oppitunnin aikana.

Tarvitsenko kokemusta aloittaakseni Node.js-taustakehityksen bootcamp-opiskelun?

Aiempi kokemus ei ole tarpeen. CoddyKitin Node.js-taustakehityksen bootcamp-oppimispolku sopii vasta-alkajista edistyneisiin, joten voit aloittaa tästä tai alusta ja edetä omaan tahtiisi. Tämä on oppitunti 3/4.

Kuinka kauan ”Pub/Sub, Streams ja nopeusrajoitus Redisillä”-oppitunnin suorittaminen kestää?

Useimmat CoddyKitin oppitunnit kestävät noin 5–10 minuuttia. Jokainen oppitunti on lyhyt ja interaktiivinen, joten edistyt tasaisesti ja voit jatkaa siitä, mihin jäit – sekä verkossa että sovelluksessa.

Voinko kirjoittaa ja suorittaa koodia tällä Node.js-taustakehityksen bootcamp-oppitunnilla?

Kyllä. Jokainen Node.js-taustakehityksen bootcamp-oppitunti sisältää sisäänrakennetun koodieditorin, joten voit kirjoittaa ja suorittaa oikeaa koodia suoraan selaimessa ja saada välitöntä palautetta tekoälyltä – paikallista asennusta ei tarvita.

Kaikki tämän kurssin oppitunnit

  1. Cache-aside-, write-through- ja TTL-strategiat
  2. Hajautetut lukot ja Redlock-algoritmi
  3. Pub/Sub, Streams ja nopeusrajoitus Redisillä
  4. Cache stampede- ja thundering herd -ilmiöiden estäminen
← Takaisin: Node.js-taustakehityksen bootcamp