4 ms·
118M Queries per Second on Neki
- jeffbee 22d agoThe fact that you can just pay to scale out point reads is not news to anyone.
- AdamProut 22d agoyeah, this is a definitely a "best case" workload for a sharded database. Single row reads on the key used to shard with no hotspots (no shard to shard network traffic at all).
- samlambert 22d agoIt cost $250,000 to do this run but it feels worth it.
- handfuloflight 22d agoI did not know men could build such things.
- whalesalad 22d agoI estimated the cluster to achieve this was ~$3-4k per-hour. I am thinking there is a typo on the r8g.16xlarge and they are actually r8gd.16xlarge (notice the d) which comes with directly attached nvme disks.
- svuiv 22d agoWe used r8g.16xlarge instances with EBS disks, no nvmes
- whalesalad 22d agowould love to hear more about the ebs volumes, iops/size/raid configuration
- svuiv 22d agoeach shard had a 4TiB volume with 65k IOPS and 1,500 mbps of throughput
- rcrowley 22d agoIt felt rude to take so many r8gd instances away from our customers who really love those (and i8g and i8ge).
- jeffbee 22d agoThat's roughly 250x more than it would cost to perform this stunt using on-demand Cloud Bigtable, if my math checks out (~1150 nodes @ 85¢/hour for 1h).
- dataviz1000 22d ago> 512 shards, each with one Postgres primary each on an r8g.16xlarge > 480 Neki routers, each on its own 8xlarge instance > We sustained 118,538,803 QPS for 16 minutes across 512 shards and 1.22 PiB of data. Our largest recording was 118,747,267. Component Detail Monthly Hourly 16-min burst --------------------------------------------------------------------------------------------- Shard compute 512x r8g.16xlarge $1.41M $1,930 $515 Router compute 480x r8g.8xlarge* $661K $905 $241 Storage (gp3 floor) 1.22 PiB @ $0.08/GB-mo $102K $140 $37 Storage (io2 floor) 1.22 PiB @ $0.125/GB-mo $160K $219 $58 IOPS (io2, light) 5K IOPS/shard, tiered rate $166K $228 $61 IOPS (io2, medium) 20K IOPS/shard, tiered rate $666K $912 $243 IOPS (io2, worst-case) 231,517 IOPS/shard (0% cache) $4.56M $6,251 $1,667 --------------------------------------------------------------------------------------------- Total (gp3 floor) $2.17M $2,975 $793 Total (io2 floor) $2.23M $3,054 $814 Total (io2 + light IOPS) $2.40M $3,282 $875 Total (io2 + medium IOPS) $2.90M $3,966 $1,058 Total (io2 + worst-case IOPS) $6.79M $9,305 $2,481
- deleted 22d ago[deleted]
- znpy 22d agoIf this is closed source then i have zero interest in it.
- noir_lord 22d agoSomeone posted a twitch conversation yesterday about this, I poked around on the page realised there was no open source version and noped out immediately. I'm sure it's a great product (it seems like planetscale do good engineering and the folks I know who use them seem fine with it) but I don't do vendor lock-in as a service personally, I'll use whatever employer uses because that's the deal but for personal stuff, well this isn't designed for that really, wrong order of magnitude on scaling.
- jjice 22d agoI believe multigress is the similarly aged open equivalent from Supabase. Haven't used it myself and don't know what the differences are in usability, but I'm a bit more interested in that since it's open.
- spongeboi 22d agounfortunately they haven't been able to move forward w/ the project, it can't even shard yet
- jjice 22d agoIt appears to be under active development, but you're right that sharding doesn't appear implemented. I'm excited to see what they can build out over the next few years.
- samlambert 22d agoNeki right now: Multiple live shards: yes Query routing across shards: yes Online shard splitting: yes Zero-downtime resharding: yes Multiple independent shard groups: yes Data topology management: yes HA / automated failover: yes Multi-AZ: yes Connection pooling: yes Online schema changes: yes Workflow-driven migrations/cutovers: yes Zero-downtime imports: yes CDC / logical replication: yes Online Postgres version upgrade workflows: yes Cross-shard transactions: coming Multigres today: Multiple live shards: no Query routing across shards: no Online shard splitting: no Resharding: no Multiple shard groups: no HA / failover: yes Multi-AZ: yes Connection pooling: yes Logical replication/import work: in progress Distributed migration/resharding workflows: no How it is an it's an alternative? Do you just say things without validating?
- AdamProut 22d agoI'm curious why the test needed so many router hosts: 512 shards, each with one Postgres primary each on an r8g.16xlarge 480 Neki routers, each on its own 8xlarge instance That's ~250K queries/sec per router which seems lowish for this type of workload? The routers won't be doing very much (parse query, route it to proper shard?).
- svuiv 22d agothat's over 13k queries/sec per router core, about 50% of it is spent doing syscalls, the other 50%: parsing, doing grpc, tls, go gc, resolving the shards, waiting for the responses neki is still in platform preview, this experimentation allowed us to collect profiles at such scale and ship some nice optimizations, more are coming
- danbruc 22d ago87.3 % served from cache. Does that mean it returned a result existing in the cache because the very same query was executed before? Probably still a relevant result, if you have to process millions of queries every second, it seems not unlikely that you will see a lot of repeated queries. But at that point you are measuring cache performance more than query performance. But unless you run some standardized query benchmark, a single queries per second number is not that informative anyway because query complexity and therefore execution time can span many others of magnitude. Looking up a name by ID and aggregating across a billion rows from seventeen tables joined together are both a single query.
- farazbabar 22d agoIn 2015, I was able to get to 1 million read/write queries per second on only a couple nodes and tested this with multiple databases, it required (at the time) decent network tuning and node placement inside AWS but it cost me about 10 to 15 dollars per run if I recall correctly, obviously there is the matter of scaling such performance and so I want to recognize the engineering effort gone into this but this is too much money. This reminds of when one of my teams used Hadoop to process only a a few terabytes of offline data and were able to process the WHOLE THING in only a few hours. I did not have the heart or courage to tell them during the demo that this was overkill, but I did write a very simple (and small) piece of code that could extract all the signals from the offline files in mere seconds with careful network planning and storage optimization and invited them for a demo/lunch and learn next week.
- tomnipotent 22d ago> all the signals from the offline files in mere seconds That sounds unlikely unless those machines had access to crazy disk I/O. A local RAID 5/10 with 8 drives would still take 30-50 minutes just to read that much data. Even with a mid-range SAN you still would have spent 15-25 minutes just reading data. This assumes 7.2K SAS/SATA since SSD/NVMe were not ubiquitous in 2015, but even with 2015-era SSDs you're still looking at half that much time spent reading.
- farazbabar 22d agoLots of EBS volumes mounted via 25GBPs or higher network, positioned carefully onto a single rack where possible and carefully tuned network/ip/os for both clients and servers. It is not feasible in most deployments as this would not have scaled to real production loads (RAID configurations requiring redundancy alone would slow you down, not to mention costs of using that many EBS volumes on extra high network/IO/provisioned IOPS nodes). This was only meant to prove what was possible in AWS at the time.
- tomnipotent 21d agoNot in 2015. 25GBps didn't even arrive until 2016, and dedicated EBS bandwidth didn't arrive until 2017.
- _zoltan_ 22d ago> The benchmark was very simple. A single-shard point select, one row fetched per-query by primary key. No writes, joins, or cross-shard queries. The workload that each shard receives is isolated, in that there are no single queries that span multiple shards. I mean... What's the point of this "benchmark"?
- voodoo_child 19d ago> Worth being clear about this run: the shards were primary-only with no replicas, the workload is read-only across queries ranging in complexity, and we did not fail over during the measured window. No HA either, or backups I assume.(single AZ?) But think the point of it was not to benchmark PostgreSQL but to demonstrate the linear scalability of the shared-nothing/stateless routing architecture. I think this was already known from Vitess though, which shares the same architecture. Also, as part of a new product release benchmarking shit like this internally is cool/interesting, why not blog about it. The more cynical side of me says it was a money burning exercise to give the marketing team a number to tweet about. :)
- cbg0 22d agoI did some quick ChatGPT math for the same performance/storage as the benchmark: Neki 1 primary + 2 replicas: ~$5.0M/month (just the AWS bill) Google Spanner w/ 3 replicas built in: ~$3.85M/month
- deleted 22d ago[deleted]
- stephenlf 22d agoI watched an interview that Casey Muratori did with Tyler Cloutier (SpacetimeDB founder and spokesperson) [1]. One of the points that Tyler is that, given modern CPU architecture with cache lines, a distributed database needs to fan out to at least 50-100 nodes to beat the throughput of a cache-optimized, single node database. It’s cool to see the flip side of that argument. Planet scale is answering the question, “what does it look like when you DO fan out your workload to >100 nodes?” There’s a place for both technologies. Very cool stuff. [1] https://youtu.be/ONxwjqFjP3A?is=awlEJwGLxQmRE25i https://youtu.be/ONxwjqFjP3A?is=awlEJwGLxQmRE25i
- titanomachy 22d agoIME people usually go multi-node for availability and durability reasons, not throughput.
- malisper 21d ago> I watched an interview that Casey Muratori did with Tyler Cloutier (SpacetimeDB founder and spokesperson) [1]. One of the points that Tyler is that, given modern CPU architecture with cache lines, a distributed database needs to fan out to at least 50-100 nodes to beat the throughput of a cache-optimized, single node database. Do you have the timestamp where they are talking about this? The claim doesn't pass the smell test for me. If you're talking about latency, then perhaps. On throughput, I don't understand how a single-node system could deliver higher throughput than a three-node system
- jamesblonde 22d agoIn 2015, MySQL Cluster (NDB Cluster engine) benchmarked 200m transactions/second on commodity hardware [ref]. It was read-committed transactions, not snapshot isolation, but still impressive. NDB has now become RonDB, but is still based on a non-blocking 2-phase commit protocol and is GPL-v2. RonDB now has support for infiniband, so ought to blow through the 1B ops/sec. For reference, that is 1 GHz of transactions/sec. [ref] https://www.slideshare.net/frazerClement/200-million-qps-on-commodity-hardware-getting-started-with-mysql-cluster-74 https://www.slideshare.net/frazerClement/200-million-qps-on-...
- weekendcode 21d agoIts great to see scale and progress, but being closed source is HUGE DEALBREAKER. Clickhouse is also on the right track of building some amazing opensource integrations with postgres, they have superior*[1] managed postgres looks like from their recent blog. I hope they do some OSS sharded postgres solution. [1] - https://clickhouse.com/blog/benchmarking-nvme-managed-postgres-planetscale-vs-clickhouse https://clickhouse.com/blog/benchmarking-nvme-managed-postgr...
- ahachete 21d agoIf you want an OSS sharded Postgres, just use Citus. Feature-full, proven, mature, boring.
- weekendcode 21d agoits not feature-full. like schema changes locking, coordinator node.. Maybe I am wrong, but I am yet to read stories on operating tens of TB scale workloads on citus.
- saisrirampur 21d agoMost Citus workloads were 10s of TB with largest at around a few PB or so. Heap was a couple PB, back then, if I remember correctly. It is a brilliant piece of technology that supported mission critical workloads across mid/late stage startups to huge enterprises. The planner/executor are very advanced supporting a multitude of features and decade of intricate effort. The biggest problem of Citus was migration effort, transition from single node to multi-node was not trivial. Here I’m not talking about single table use-cases, more classic relational, multi-tenant apps with 100s to 1000s of tables. This is partly expected with most sharding technologies, though. Sharing some insights based on my multiple years of experience working with Citus! Here are few customer use-cases I could found: https://docs.citusdata.com/en/v10.0/get_started/what_is_citus.html?l#how-far-can-citus-scale https://docs.citusdata.com/en/v10.0/get_started/what_is_citu... https://info.citusdata.com/rs/235-CNE-301/images/Citus_Data_Case_Study_-_MixRank.pdf https://info.citusdata.com/rs/235-CNE-301/images/Citus_Data_...? https://www.youtube.com/watch?v=F6df3HV6kP0 https://www.youtube.com/watch?v=F6df3HV6kP0