118 miljoen queries per seconde op Neki

Gisteren hebben we Neki in platform preview uitgebracht. Om deze release te vieren, wilden we testen of we 1 miljoen queries per seconde (QPS) op Neki konden draaien. We bereikten dit doel vrij snel met 5 shards en besloten te kijken hoe ver we konden gaan.

De volgende run eindigde met 512 shards die 118 miljoen queries per seconde verwerkten met 1,22 PiB aan data.

Lineaire schaalbaarheid

De benchmark was zeer eenvoudig: een single-shard point select, waarbij per query één rij werd opgehaald via de primaire sleutel. Er waren geen schrijfacties, joins of cross-shard queries. De werklast die elke shard ontvangt is geïsoleerd, wat betekent dat er geen queries zijn die over meerdere shards verspreid zijn.

Ons doel was om 200k QPS per shard vast te houden en vervolgens het cluster uit te breiden om de doorvoer te verhogen. We startten met vijf shards, gingen daarna naar vijftig en uiteindelijk naar 512.

ShardsRoutersGeleverde QPSQPS per shard
51999.624199.925
50489.923.900198.478
512480118.538.803231.521

Tien keer zoveel shards betekende tien keer zoveel doorvoer, en daarna opnieuw tien keer zoveel. Van 5 naar 50 shards bleef de snelheid per shard binnen een marge van 0,8%. Bij 512 shards was er nog steeds ruimte over, dus we lieten de loadgenerator deze benutten, waardoor elke shard uiteindelijk op 231k QPS uitkwam in plaats van 200k.

118,5 miljoen QPS

We hielden 118.538.803 QPS gedurende 16 minuten vast over 512 shards en 1,22 PiB aan data. Onze hoogste registratie was 118.747.267 QPS.

Specificaties van de run:

  • 512 shards, elk met één Postgres primary op een r8g.16xlarge instance.
  • 480 Neki-routers, elk op een eigen 8xlarge instance.
  • p99 latency van 6,06 ms bij de router en 13,95 ms bij de client.
  • 67 fouten per seconde, wat neerkomt op ongeveer één query op 1,8 miljoen.
  • 15,8M read IOPS over de gehele vloot.
  • Meer dan 2 Tb per seconde op het netwerk.

Om duidelijk te zijn over deze run: de shards bestonden uitsluitend uit primaries zonder replica's, de werklast was read-only over queries met variërende complexiteit, en er hebben geen fail-overs plaatsgevonden tijdens het gemeten tijdsbestek.

In een toekomstig artikel zullen we dieper ingaan op de technische inspanningen en de interessante uitdagingen waar we voor stonden om de 100 miljoen QPS te bereiken.