6 ms·
Making 768 servers look like 1
- mike_hearn 3mo agoI should start by disclosing that I work part time in the Oracle Database group, but - of course - my HN account is entirely my own, despite occasional mild shilling. The article itself is shilling for PlanetScale so that seems OK. The author - certainly not deliberately! - says some untrue things about relational databases. The most important one is this: > To understand why sharding is a necessary part of scaling relational databases... But sharding isn't a necessary part of scaling relational databases. It's really just a requirement of simple databases like Postgres and MySQL. If you want a relational database cluster that really does make 768 servers look like one, then you want what Oracle calls RAC ("real application clusters"). Your cloud will probably rent you access to one under a name like Autonomous Database. Self hosted it may be called ExaData, which is a unified hardware/software "rent a rack" style offering. You may be surprised to discover that it's not much more expensive than many managed Postgres offerings. RAC can scale a non-sharded relational database horizontally. That means all queries can access all data, all SQL features like sequences and joins work, any server can take part in transactions with any other and in general it looks exactly like a really big single machine would. In other words, it's a synchronous multi-write-master system. RAC scales very well. 768 servers is well within reach as long as your query patterns scale too i.e. don't all contend on writing to one row. Behind the scenes it uses a dedicated high speed RDMA network with lock coordination to transfer data blocks directly between nodes, never hitting disk for memory that's already in the buffer cache. Additionally RAC is fully HA and supports rolling upgrades of the cluster whilst live. You don't need any NLBs or routers either. The client drivers automatically discover and load balance between nodes without needing intermediaries, transactions can start on one node and fail over to another without applications noticing, and so on. There's plenty of opportunity for caching and replication. You can run asynchronously replicated failover clusters, run multiple clusters in a Raft-driven globally coordinated super-cluster and can deploy coherent read-through caches anywhere; the main clusters will inform them the moment data in them becomes stale. In other words, it can do a lot. If for some reason you do need sharding then that's also supported with features like automatic sharding key distribution to client drivers that transparently route queries correctly, but most apps don't need this. In case you're wondering why I say all this, firstly, obviously, I have a financial conflict of interest. But the database hasn't driven Oracle's stock price for a long time, so it's not a big one. These days it's all about cloud and AI. No, the main reason is that HN fills up every month with blog posts where engineers talk about the incredible pain involved with scaling and running Postgres. And almost always, it's clear that they don't realize there's any alternative to that pain. It's not that they considered the options and then explain why they picked this one, it's that they think - as claimed in this article - that it's almost some fundamental limitation of reality itself, imposed by the laws of computer science. There are lots of startups that lose time and money due to database problems they simply don't need to have. And that sucks. If they'd prefer to spend that time, pain and money to avoid using a DB from Oracle, fine, so be it. I won't argue with random devs about lawnmower memes. But if it's because they don't realize what's possible.... well, maybe someone will be helped by being aware of this. Database scaling problems are a choice, not an inevitability.
- marcosdumay 3mo agoLol about buying Oracle for performance! But anyway, do you have benchmarks from anybody not related to Oracle showing that performance? Because Oracle forbids talking about it...
- mike_hearn 3mo agoNot forbidden. You can email a specific address to share results/setup and get permission to publish if you want. Other commercial databases also do that, because there's so many ways to misconfigure a database to make it slow and competitors are strongly incentivized to do so. The question is what you'd want to benchmark. For example, imagine testing Postgres with a write load that goes well beyond what a single machine can do. It would collapse and query latencies would go through the roof. A horizontally scaled DB would keep going and process all those queries. Would you accept this as evidence that Postgres is slow or would you say it's not valid to benchmark at traffic loads Postgres physically cannot handle, given it never claimed to scale horizontally? Stuff like this is where benchmarking gets complicated. I used to work at a different company that sold a kind of database system. It was much faster than our nearest competitor, so we were surprised when that competitor claimed to an important customer they were just as fast as us. Their benchmark counted transactions that failed and rolled back due to overload (optimistic concurrency) as "successful".
- inigyou 3mo agoI suppose Oracle will only allow publication if your benchmark results look good. Can you tell us about a time you did a benchmark that made Oracle look bad, and they still let you publish it?
- marcosdumay 3mo ago> Other commercial databases also do that That's why you won't see many serious comments bragging about the performance of SQL Server either. Even though is has way fewer performance land-mines that you must design your entire architecture around than Oracle. Anyway, you are correct that RAC is a very impressive piece of software. If Postgres had something like it, it would be a beast. It's not enough to save Oracle, though.
- alightsoul 3mo agoLoad balancers, microservices and horizontal scaling?
- jdw64 3mo agoLooks like the GIF is fully built out in code. It's really nice to look at, well made, and easy to understand too. I wonder what program or code they used. I'd love to know. p.sI thought it was a GIF, but it's an iframe. That was a nice little surprise.
- gurjeet 3mo agoSpecifically, it's an JS-controlled/animated SVG embedded in an iframe.
- jdw64 3mo agoYeah, I'm looking at it in developer mode. It's using a GSAP timeline approach to update SVG properties. I'm curious how they handle security and caching for something like this. It looks like they're using Tailwind, at least. but this approach is really clean and nice. It really feels like the best way to learn is by studying other people's code.
- stavros 3mo agoWhat kind of security and caching concerns do they need to handle to animate an SVG?
- jdw64 3mo agoI didn't explain myself very well. With iframes, there can be security concerns in general. Though in practice, there's usually no real issue. It just surfaces the surface level risk that JavaScript could be tampered with, but it's almost never something to actually worry about. As for caching strategy, cache invalidation issues do pop up from time to time. Since the filenames usually include a content hash, I'm guessing they're using a Vite style caching strategy. Cache invalidation might almost never be a problem, but you still plan for it just in case.
- bddicken 3mo agoAuthor here, thank you. Technically, they are using js + gsap + svg embedded i the article with iframes. Process-wise, I drafted most of them as static images in excalidraw, passed the images along to cursor for a first draft, applied styling rules, and then did a bunch of fine-tuning.
- aarvin_roshin 3mo agoPreviously: https://news.ycombinator.com/item?id=48925420 https://news.ycombinator.com/item?id=48925420
- drdexebtjl 3mo agoWhat about sequences? The example shows an auto-incrementing user ID. How’s that possible without contention between all shards? Is the proxy responsible for sequences? What about foreign keys? Do they all have to live on the same shard? How do you do distributed transactions? On cross-shard reads: how do you do sorting? And cross-shard joins? I’d love to be proven wrong, but I suspect the 768 servers look like 1 only on the very surface, and you’ll get wildly different characteristics from cross-shard and single-shard queries. I personally would prefer if they _didn’t_ look like 1 if they can’t behave like 1.
- random3 3mo agoA 767 servers KV store should be enough for everyone
- vkazanov 3mo agoOf course 768 servers NEVER behave as 1. This is physically impossible. Global services using relational dbs typically severely restrict queries that run against the cluster. So no joins, no intervals, no grouping, etc. Transactional queries are usually limited to something like "get a single record, preferably from cache". For many typical web services this can go VERY FAR. Only a handful of global services needs more than a few dozen database servers and a caching cluster. In fact, i have seen major businesses running off a pair of very big postgres instances. Analytical stuff is extracted into dedicated storages optimized for throughput, like Snowflake or Redshift or BigQuery.
- ahk-dev 3mo agoThis seems like the important distinction: making the infrastructure look like one database to the application is different from making it behave like one unrestricted relational database. At what point does hiding the sharding become counterproductive? I imagine teams still need a fairly deep understanding of shard keys, query routing, and failure modes to avoid accidentally expensive cross-shard operations.
- vkazanov 3mo ago
- zinodaur 3mo agoSibling post has author answering questions in comments: https://news.ycombinator.com/item?id=48925420 https://news.ycombinator.com/item?id=48925420
- groundzeros2015 3mo agoI disagree with the opening premise: > A single database server cannot handle such demand, so we must spread the queries and data out across many servers with database sharding Did you max out the capacity of the best server you can buy? Such a database can serve millions of customers (the numbers given). You always want to scale up the other parts first, request handlers, caching, etc. The day you can no longer inspect the essential state of your system is the day your company better be included in NASDAQ and ready to pay a few hundred engineers 300k salaries.
- anonzzzies 3mo agoWell, they are selling this thing so they don't want you to buy a big server (with a read replica) as that's much cheaper.
- farslan 3mo agoThat's not true! We have demand for bigger machines and we also sell them. You can go checkout the PlanetScale website. See: https://x.com/samlambert/status/2077197049129587150 https://x.com/samlambert/status/2077197049129587150 There are truly customers that bigger machine no longer cuts. Disclaimer: I'm an Engineer at PS.
- anonzzzies 3mo agoI did not say there is no market, but I am willing to bet whatever that the majority of your client base would be much cheaper off with a simple setup where they can scale the entirety of their company’s existence with one machine. Not even a large one. This cloud stuff has its place but not for most. For them it’s either themselves overspending or wasting VC money.
- bddicken 3mo agoWe have tons of customers who do exactly this. It's great. Sharding is for customers who out grow this path.
- Hugsbox 3mo agoTook me an embarrassing amount of time to realize this is an ad.
- foxhill 3mo agoi do wonder how something like this can be generally implemented. i presume this must only support a subset of SQL/plpgsql, as some things would be.. utterly insane to manage manually. e.g., if i have a table with a btree-gist overlap constraint, or some inclusion-exclusion check-constraint (or literally any constraint that requires multiple rows to be fully determined - there are quite a lot of them), how on earth does this work? there's a reason why postgres writing is (mostly) serialised (asterisk) to a single writer (asterisk asterisk). something something ACID, but in short by having multiple writers improves availability, but weakens integrity.
- xbas 3mo ago[dead]
- Ellis_dev 3mo ago[dead]
- skeptic_ai 3mo agoHow does Sharding works when You do complex joins ? Seems tricky , runs on each server and gets data back and aggregate it again?
- KoleSeise1277 3mo ago[dead]
- BedVibe_Studios 3mo ago[flagged]
- kjellsbells 3mo agoI sometimes feel that when the industry moved from pets to cattle, what really happened is that the cattle turned out to need an exotic farm to live on, negating the savings. You can have a few honking servers or you can hand massage exotic k8s setups on your fleet. Pick your poison, but dont delude yourself that the TCO of the latter is lower than the former.
- metalliqaz 3mo agowhat did they use to make those diagrams/animations?
- bddicken 3mo agoHey, author here. Technically, they're powered by js + gsap + svg. Process wise (for most of them) I sketched them out in advance in excalidraw for figure out layout, then passed these along to cursor to have it build out an initial draft from the image, then used some styling rules to get all the styles inline, then did a bunch of fine-tuning.
- themgt 3mo agoEven with a large database servers (10s of CPU cores, 100s of gigabytes of RAM) bottlenecks arise pretty quickly. Err, do they? For what percent of real world use cases? The database can scale to handle more traffic by adding replicas. An extreme example of this is OpenAI's use of 50 replicas on a single Primary. So an extreme example is OpenAI needing 50 replicas, but we're doing five blades ... err, we're doing 768 servers because the need arose "pretty quickly"? When we needed to store a petabyte of data (one million gigabytes), we'd need many more shards For who? The United States government? How many end-users are running 1PB Postgres database on DBaaS?
- exiguus 3mo agoPersonally, I know several databases where single tables have +500GB and the database has +100TB. With this huge databases the restore and backup process over network become indeed a bottleneck. So I can agree with the author. Also, the author does not say that they can't start with a single database server and just read replicas and max hardware out. Real world use cases with +100TB I know about are Stock Market, Traffic Data, Warehouses, Analytics and Monitoring.
- bddicken 3mo ago> So an extreme example is OpenAI needing 50 replicas, but we're doing five blades ... err, we're doing 768 servers because the need arose "pretty quickly"? If you read the OpenAI article, you'll see that they actually used sharding to offload a bunch of work from their "1 primary 50 replicas setup" >>> "To mitigate these limitations and reduce write pressure, we’ve migrated, and continue to migrate, shardable (i.e. workloads that can be horizontally partitioned), write-heavy workloads to sharded systems such as Azure Cosmos DB..." The 768 servers and 1PB example is just one of many configurations. A business with 10TB may choose to go from a monolithic database to a 8-shard setup to improve backup times, eliminate single-point-of-failure, have more breathing room for scaling.
- hasyimibhar 3mo agoI'm surprised no one has mentioned Multigres yet, they are the competitor of Neki. I'm a big fan of both and has been following them since they were announced last year. I think this blog post is the first time they are talking about the internals of Neki. In contrast, Multigres is being built in public since day 1, you can see their high-level architecture here [1], though I'm still waiting for more details on their sharding model. [1] https://multigres.com/docs/architecture https://multigres.com/docs/architecture
- nazgulsenpai 3mo ago> Most applications you've ever used function in this way, or at least did early in their existence. I'm old enough that this is not true.
- perceptronas 3mo agoHow do you deal with master failures for specific shard? Does switch happen automatically once server becomes unresponsive?