The database ruins all good ideas
squarism.com
squarism.com
A few other points. First, horizontally scaling database reads is actually reasonably straight-forward these days: one leader, multiple replicas, load balance reads to the replicas and use a mechanism such that users that have just performed a write have their read traffic sent to the leader for the next 10 seconds or so.
Not trivial, but also not impossibly difficult - plenty of places implement this without too much trouble.
Scaling writes is a lot harder - but a well specc'd relational database server will handle tens of thousands of writes per second, so the vast majority of projects will never have to solve this.
When you do need to solve this, patterns for horizontally sharding your data exist. They're not at all easy to implement, but it's not an impossible problem either.
The article talks briefly about mocking your database: definitely never do this. How your database behaves should be considered part of your application code under test. Running a temporary database for your tests is a solved problem for most development frameworks these days (Django supports this out of the box).
Overall, my experience is that the database /enables/ all good ideas. Building stateful applications without a relational database in the mix is usually the wrong choice.
There's one feature that I'd like to see in that area: partitioned writable replicas. In the same vein that you can partition the table storage across an index, I'd like it to be possible to assign different writers to different parts of a table/database. Of course, you'd still need a single primary replica to handle the transactions that transcend the configured partitions, but we already have routing engines that can transparently redirect an incoming query to any available replica, so it's partly there.
There's probably corner cases lurking that I can't even think of, but in my mind it's the only thing missing from building a truly web-scale (multiple zones, multiple datacenters) ACID-preserving relational database.
This is getting more common these days. Google’s Spanner database does this. The database is divided into “tablets” which are each given their own Paxos state. Spanner moves tablets between server instances to do load balancing. There’s no need for you to do the partitioning manually, although if you run into performance problems, you may need to understand how tablets work.
Spanner’s definitely not unique in doing this, but Spanner does have a nice writeup of how it works:
https://static.googleusercontent.com/media/research.google.c... (PDF)
Thus the write throughput can be very high but the write latency is also very high.
When you do it manually you are very aware of queries that need to go to multiple shards and really do whatever it takes to avoid them.
We sharded by account id and all an individual users data would be on that shard.
Basically the only queries that need to transcend the shard are when you need to enumerate accounts by some non-ID value like an email address. You have to ask all the shards in that case.
Things like creating an account are tricky if you want to make sure that each has a unique email address. You need to use two-phase commit across all the databases to do it reliably.
If you're happy letting the database do that for you, that's how current CockroachDB works.
The issues described in this blog are about the challenges of persisting a consistent set of data, when you need to support a lot of read and write access. That is just a fundamentally difficult problem. Databases are the most sophisticated tools we have for doing that, but really they just tend to be a set of (hopefully) well implemented trade offs that attempt to nicely balance all of those conflicting requirements.
If a person can’t find a way to adequately persist their data, then their idea is probably bad, or they just don’t have the competence to implement it. An equally suitable title for this post could have been “why doesn’t magic exist”.
If you can't actually retain your business data, it really doesn't matter how many "good ideas" or fancy deployment strategies for the presentation tier you have.
I honestly thought you meant to never ‘ridicule’ your database, and I was thinking “that might be a healthy level of respect” :)
Spring Boot for Java does this too (generally using H2 in memory database). I have found that for CRUD services, doing this vs. pure unit tests that mock the db, yields far better benefit (provided you don't involve heavy volumes in the tests to avoid test data maintenance issues). With this approach, I have found database side edge cases/constraint violations that wouldn't be caught by mocking the database. Also, it is a much more comprehensive test if you take the approach of a) loading an initial state, b) doing business logic mutations and c) compare the final state with the expected state; e.g. you might find side effects that you hadn't considered testing for and so on.
It is easy enough to setup and tear down using docker.
https://vladmihalcea.com/how-to-run-integration-tests-at-war...
Otherwise I have a single docker-compose.yml file that sets up postgres using "Docker for Desktop" on Mac and Windows and works natively on Linux. CI (Linux) uses the same docker-compose.yml.
I haven't used docker in production, but it has proved very useful in test.
https://github.com/zonkyio/embedded-database-spring-test
also helps with having it verify all your database migrations as well
Agreed re: db mocks, database bugs and quirks should be captured and addressed in tests, not discovered only in production. I don't care how mature the DB you're using is, there's always going to be some edge case behavior that you can't predict until you actually see it happen, even just simple stuff like exactly how an error filters back to the application when you send data that violates a constraint. I've seen way too many tests that make assumptions about these things that turn out to be incorrect, leading to horrible behavior like writes getting through on prod that should have been blocked because of pre-checks that didn't turn out to be enforced the way the test writer assumed they would be.
I definitely recommend doing this for same tests: it's the only way to check how your application behaves in case of failures.
The best thing would be to have a DB where you can inject failures, but I'm not aware of any. So test - the happy path on the real db - the sad path on the mock/fake (which might be a light wrapper around a real db, but with the ability of injecting failures)
* Read-only queries can easily be scaled out to read replicas.
* Write transactions.
As an example of numbers publicly available, GitLab.com runs more than 250K read-only txs and more than 60K write txs on a single Postgres cluster, with room for further vertical scalability [1].
Do you need millions of transactions per second? Then go ahead. Don't? Then ou can scale very very far with a properly tuned and administered relational database like Postgres.
[1]: https://about.gitlab.com/blog/2020/09/11/gitlab-pg-upgrade/
I'll also add, that the point is that a devtools unicorn is able to run with a single postgres cluster, so maybe it isn't the rdbms that's going to be the limiting factor on your startup with a hundred users.
Sometimes you'll run into situations where an instance of any software has degraded performance for reasons beyond either your control or understanding. That's one of the situations where horizontal scaling let's you deal with such circumstances better, instead of outright impacting everything and everyone connected to it.
A front end web server container crumbles under the load? No worries, just redirect the traffic to others. A back end API container has problems with the JVM reserved memory and/or the thread pool for processing connections has been filled and all new incoming requests just get queued up? Just have more instances to cope with the load and don't let one instance under load affect all others.
But what are you supposed to do if long running SQL exhausts the resources of your DBMS? What if a background process needs to run some really complicated writes, but your application needs write capability without slowing down?
That's not to say that bad configurations or bad code shouldn't be addressed, but any means of fighting it and minimizing the impact is worthwhile in my eyes.
You can do basically the same thing with read replicas, promote one of them to primary and shut down the machine causing issues.
You're far better off starting with a system that does seamless active-active by default, and then figuring out what kind of transactional guarantees you need. An RDBMS might make sense for a few specialised use cases, but most of the time it doesn't.
That being said, if your code doesn't catch and handle errors, then there's a couple of deeper problems already.
Well if you can't get input that violates that integrity then what are you gaining by enforcing that integrity?
> That being said, if your code doesn't catch and handle errors, then there's a couple of deeper problems already.
Sure, but how can you do error handling without data storage, given that your application itself is stateless? Maybe you retry the failure a few times, but ultimately your only option is to drop the data on the floor, which is almost always not what you want.
Fixing up data that was written wrongly takes work, there's no getting away from that. But if you have a record of the original write attempt then at least you can do that work and recover the data (if you decide the cost/benefit is worth it). Whereas if you have the datastore reject it at write time then you're SOL before you've even started.
Have you ever actually encountered any RDBMS that just silently "drops" your write without informing you, or is this just how you think they work?
Otherwise, isn't this just like any other data write erroring out -- MongoDB driver saying "Service down!", simple flat-file write saying "Disk full!", etc -- in that if your app doesn't react to this, it's not the data store but your code that sucks?
Perhaps I'm wording this badly?
Your application should _not_ be generating data which can't go into the database correctly. If it is, that's a different set of problems. ;)
Validation is vital but the datastore is not the place to do it, because handling invalid data by dropping it on the floor is almost never the right behaviour.
There is no substitute for actually understanding your data model, but once you do 99% of the time you'll find using an RDBMS comes with minimal benefits and significant costs.
Maybe this for consumer web applications, or for CRUD apps.
But my experience with B2B and SaaS applications is that sooner or later, I always need transactions (or locks) for something. Maybe it's for the billing code. Maybe it's for batch jobs or workflow logic. And that's when I'm really happy to have a transactional database.
The alternative to having transactions is setting up a Raft or Paxos server to handle distributed locks, and those require a lot more ops effort than a copy of PostgreSQL.
PostgreSQL works surprisingly well for something like this. It has a rich set of coordination and locking primitives, and it can usually be recovered and restarted in under 10 minutes if the server fails. Yes, connection overhead from 300 workers is a real problem, so it doesn't hurt to put a thin layer in front of it.
One thing that an RDBMS buys me is flexible feature set and a solid ecosystem. Want write-ahead logs and automatic point-in-time recovery? It's out there. Need a specific locking primitive? PostgreSQL probably supports it. Etc.
If I were looking at 2,000 workers, yeah, it would be time to set up a Raft server. But a good RDBMS will do the job for a surprisingly long time.
I do see where you're coming from, but I find they're more trouble than they're worth in the long term. YMMV I guess.
It's possibly the sanest default for a data persistence layer. I'd argue otherwise: use a RDBMS unless you know what you are doing.
> the number of web applications that make effective use of database-level transactions is approximately zero
Almost every driver and layer (ORMs, etc) that connect to databases use transactions. Transparently or not, but it uses them. So actually they are using transactions all the time.
> You're far better off starting with a system that does seamless active-active
Think about what means active-active. Either it means that you serialize everything in every active node (pointless); or you need to deal with concurrency and data synchronization. Which in turn means either supporting distributed transactions (e.g. via Paxos or RAFT protocols) or doing conflict resolution (which essentially is a nice way of saying "data loss"). Both problems are much, much harder than the (usually harmless) consequences of using transactions on a "classical" RDBMS.
The only drawbacks of transactions are that are limited to scale in write on a single node. But if this single node can scale to the numbers I mentioned ---and they do, and more-- most use cases will likely fit within this envelope.
Therefore being RDBMSs like Postgres the best and safest, the default, approach that should be taken for data persistence layers.
They're using them but they're not getting any useful effect out of them, since the database-level transactions don't actually align with business-level transactions. Fundamentally you can't actually use database-level transactions in a useful way in a web application because you can't share a transaction between multiple HTTP requests (without introducing much bigger problems).
> Either it means that you serialize everything in every active node (pointless)
Not pointless - you gain redundancy which is the whole point. Performance will be bad but for a low-traffic system it's often actually fine.
> Which in turn means either supporting distributed transactions (e.g. via Paxos or RAFT protocols) or doing conflict resolution (which essentially is a nice way of saying "data loss"). Both problems are much, much harder than the (usually harmless) consequences of using transactions on a "classical" RDBMS.
The consequence of using transactions on a classical RDBMS at scale without thinking carefully about it is deadlocks, which are just as bad as distributed transaction problems, and you still have to do conflict resolution in the case where a transaction gets aborted, which will happen sooner or later (most RDBMS fans just ignore it and lose data). If you want decent performance then eventually you'll have to model your data properly using CRDTs (this is true even in a classical RDBMS) and at that point you gain nothing from traditional RDBMS transactions.
> Therefore being RDBMSs like Postgres the best and safest, the default, approach that should be taken for data persistence layers.
RDBMSes are not safe. The way they achieve their integrity guarantees is by rejecting, and dropping, unexpected writes. This almost always loses data in practice.
The database doesn't ruin all good ideas. It makes all those ideas for other tiers possible.
Most SQL problems that most of us have to deal with stem from inadequate indexes and/or poorly written queries. No fancy active-active setup will ever solve those issues.
I don't think its foolish to keep asking for more. We still might find it.
I’ll replicate your clown ass todo list app with a few DB tables, a few queries and 1 HTML table dont make me fucking do it :D
It’s weird how developers get asked about OO design patterns in interviews, but not once have I ever been asked about database design (beyond some useless stuff about 2nd normal form).
But you are correct, interviews for those types of jobs for whatever reason don't tend to focus as much on DB stuff, though they really really should given that a quarter of your job will end up being unfucking the system away from the poor choices that your predecessors made (if you're guaranteed to be serving 1mm DAU and everyone will be generating thousands of records per month, maaaybe don't use normal ints for uids...).
The amount of times I've had to fix UUID4s stored as a PK string in my career... It's bloody infuriating, it's like your typical 'ORM developer' just doesn't understand that: a UUID is a 128bit number; doesn't understand that in Postgres there's an extension for it (which Django doesn't use); how much slower a string lookup is compared to a numeric lookup; and how that indexes are physically stored in order, meaning that a UUID4 will be inserted at random points on disk utterly killing insert performance.
Generally these kind of performance fuckups happen because you get developers who treat the database as a magical performance box that needs an occasional index. One of the worst systems I've seen was an accounting system that used triggers to update the balance, so you'd insert a ledger that would update a single row. Guess what that'd do? It'd block all other inserts until the update was committed. The previous developers kept on trying to improve the performance by adding more servers, but it was just burning money because they were inadvertently locking a row (rather than calculating balances on the fly).
Designing a good schema is, in my opinion, the single most important thing for creating useful software. I would much rather work on janky spaghetti code that operates on a clean and well-mapped schema than beautiful SOLID, commented, well-documented code that runs against a terrible, fragmented datastore.
(I seem to recall that many SAS drives now are dual-ported but designed as such to fulfil a desire for multi-path resiliency within storage arrays; this may have different characteristics to those needed of a quorum device)
If you need ACID compliance and you need a lot of it, everywhere, all the time, now there are better options than giant Sun/IBM boxes.
Databases are not the problem.
I was just trying to say that a stable, high performance SQL database with distributed writes is a fairly new thing.
Oh I dunno, I'd still say most claims to be "better" at data integrity with performance, and a lot of it than a good old zOS mainframe with a soundly designed DB2 database on it are [Citation needed]-class...
But yeah, databases certainly aren't the problem here.
Also, it seems unfair to say that the database ruins all good ideas when there continues to be significant impressive innovation in the space in cloud services. Like Amazon Aurora for one. If you really want to treat your RDBMS like an elastically scalable (with significant caveats) resource priced by the request, you can.
Also I think these issues arise at the database because that’s were the write concurrency is pushed.
https://docs.aws.amazon.com/AmazonRDS/latest/AuroraUserGuide...
As noted there and elsewhere though often you are better off just scaling up.
If you do want to really think about it then NewSQL (Cloud Spanner or others) could be the way to go but the cost, performance, compatibility/portability and other concerns make this a much more complex decision that isn’t easy to capture in a blog post (or HN comment ;).
The answer is to start siloing customers or application concerns out into separate clusters depending on your operating model.
If you spent the last decade writing shitty cross database SQL then you’ll have to silo by customer. That’s the only real constraint.
Of course, the next "good idea" is to sacrifice data integrity in the name of performance. That can work, but usually it's just a technical form of growth hacking. Sacrificing data integrity without understanding the problem space is a disaster waiting to happen (but luckily, it may not happen on your watch so you can continue hacking).
I have directly experienced the problems associated with blindly assuming that you can write your own consistency guarantees in Redis. I didn't do it but I had to write some tooling to rewrite the AOF to get some data back...
I don't feel that's what you mean, and it's probably my fault for misreading your comment.
I mean, one of my ideas requires many transactions and records that require data consistency. I have no idea how to make sure it can scale to Internet scale. Do you know any good resources? Or is it just "get mysql, pay thousands a month for either a cloud provider or colocated hardware"
The next from a design perspective is to leave natural places to shard/split the data up so you could represent the data consistently on multiple systems. Sometimes it's pretty bonehead things like your database tables related to login are affecting your critical path transactions and/or also tied to the referential integrity of your application data in ways that are hard to split apart. It goes back to the IO thing where if you are reading and writing login or UI traffic you sure as heck are not writing some other critical time sensitive transaction especially if it's to the same disk. The other obvious thing is to be trying to run analytic queries (aka look at many records) on a database tuned for mass INSERT/write traffic and having no plan to do that with a different resource.
If you are still worried about it I would suggest just prototyping it using a currently modern API/RPC server in front of whatever datastore you come up with. As systems have gotten more distributed the straight line speed of connecting to the database directly makes less and less sense in my opinion and having an API server you control allows you to be clever later in ways that might not be obvious today (especially keeping application-aware counters so that you can see what your application is actually doing as query-style logs for databases are extremely expensive in performance). Either way if you sit down and consider the performance characteristics of the underlying system, think about each class of transactions are you considering from your application in a write/read context (hopefully from measured data), and using tools like EXPLAIN you'll quickly start finding what pieces of the puzzle might need to move to start scaling.
The answer also might just be some parts of your application would be better served something other than a database like a message queue, a more straightforward ISAM or key-value store, or depending on your situation just flat files written onto the fastest local storage you can get to then be later transformed and loaded into a database sometime later.
Did you read the article a little while back about Discord's challenges with storing tons of messages and how they upgraded their tech stack? Did it make sense to you?
They also had some clear characteristics of their system for their critical query that they could take advantage of: a new message is written to generally once, the messages table scales linearly, and updates and deletes are relatively rare. Using those things they were able to seek out a datastore that helped maximize the thing that's happening the most which is the write, and they were able to engineer a system accordingly. They also know that the messages would grow linearly without bound at a variable rate, so they could plan for that capacity. To me, that makes a lot of sense.
Especially if it is a side project you can't over-engineer capacity but you can do at least the measured math of what actions are actually happening in your system and what are your expectations for those queries. I would suggest that often times you'd want to focus on your critical queries that if slowed down would have the biggest net negative, but you'd have to find those. Every user action in a system should be attributable to some combination of reads or writes that can be reasoned about to try and figure out where there might be problems of contention, linear growth, exponential growth or growth up to a limit for a table, or something else. Either way, once you start classifying some of the critical paths in your system often either an answer or better questions start to come to the top. Without some of this information you can't back of the envelope whether or not your system can take a perceived load.
This is such an underrated solution (siloing or sharding your data in some way). I think people don't do it because:
1. The tooling doesn't make it super-easy (e.g. good luck sharding Postgres unless you're willing to pay for Citus)
2. "Trendy" companies in the past decade have been network-type products (social networks, etc), where the structure of the data makes it much harder to silo (people need to follow each other, interact with each other's content, etc)
3. We as an industry took a several-year detour over to NoSQL land as a promised solution to scalability.
Life would be a lot easier if you could say something like:
* I want a Postgres node.
* I'm happy to shard my data by some key (customerId, city, etc) and am willing to accept responsibility for thinking through that sharding key.
* My application has some logic that easily knows which DB to read/write from depending on shard key.
* There's some small amount of "global application data" that might need to live on a single node.
There’s YugabyteDB.
It radically reduces the infosec/screw up blast radius. For example it means however badly you screw up it is just about impossible to accidentally show data from customer A to a user from customer B.
You can let customers specify region/availability zone and bring their own keys for encryption making lots of enterprise security compliance things easier (for regulated use cases).
There is no one true way of scaling that suits all possible use cases and as engineers, we should think for ourselves and build a solution that works well in our situation, rather than thinking you have to do a certain thing because that's what google/facebook etc do. Their scaling problems are most likely really different to yours.
I have a team that won’t even look at the fact that we measure performance by response time, but we are lumping together write traffic and two classes of read traffic into a single monolithic chunk of code. We are SaaS where our customers write and their customers read. And of course search queries are by far the slowest traffic.
But there’s some sort of mental block about splitting things up that is only slowly changing, and our app is so “flexible” that a load balancer would struggle to tell which urls involve search functionality.
I think people just want to “know where to look” but if your code is on a cluster there already is no “there”. You’re probably looking at some log aggregator anyway, so who gives a tinker’s damn if they come from separate log files on the same box or the same log file?
Row based permissions are not great imo, either performance is crap or you have some weird bug.
Pointing out that database servers don’t scale horizontally like web servers had the same level of insight as pointing out cars don’t scale like kittens.
I'd be curious to know if there's a way to have databases talk to each other to just sync up primary keys for referential integrity. That could maximise the benefit of decoupled databases while still having good referential integrity. And a network disconnection would still mean existing PKs would be in place, it just wouldn't be aware of new ones. Not perfect, but not bad, perhaps.
I had that for a while. SQLite for the test suite and Postgres for actual runs. It was a shitshow, having to support two separate SQL dialects with separate feature sets. Then I realized that you can actually bootstrap and start a Postgres in tmpfs in about one second, so that's what my tests do nowadays.
But if you have multiple nodes that have to sync every transaction over a dedicated NIC, well, that's not a distributed system. It's just multiprocessing with extra latency.
CockroachDB and YugabyteDB would like to have a word.
Information can't be changed in two places simultaneously. That latency, no matter how small, requires us to either eat the latency, give up guarantees, or alter our system's behavior in some other small way.
Vitess [1], A database clustering system for horizontal scaling of MySQL, or Planetscale [2] which is the SaaS version. Of course everything is good on paper until you run into edge cases. But I am convinced within this decade scaling problem or hassle will be a thing of the past for 95% of us.
The most typical way to learn about splitbrain it seems. I don't think I've experienced more work trauma than split brain trauma.
For your main line-of-business database? Of course not. But a deployment of rqlite[0] for your service workers in a read-heavy workload? cuts out a round-trip out of the VM, mocking is trivial, there's a lot to like there.
I'd have liked to learn more about why it's only one db. For instance because they're often insanely efficient and one server can handle many queries.
I was also hoping that it'd suggest specific alternatives.
So, alas, you either go SQL in the early stages and then need to do considerable engineering to down-convert to say, Cassandra or DynamoDB.
Or you accept reduced database language sugar and complexity up-front (no joins, limited index/views, or architect with explicit sharding) with a more scalable database approach.
There's basically no magic sauce for scalable SQL. Frankly most of the people telling you otherwise are selling varying degrees of snake oil.
As people will point out, SQL databases on modern hardware scale pretty freaking large. So you can get a lot of mileage putting off the "true scalability refactor". The good news, by then you should know your queries and data that need full scaling, you aren't guessing ahead of time.
To distribute data, you have a data distribution scheme. Cassandra and various distributed hash maps (which is the typical approach) it is a consistent hash function. But even if you do distribution using natural ordering, the same problem exists:
The data you are joining is going to be on different nodes on a row-by-row case. Hashing will produce this by the design of the hash function. Natural ordering will do this because different key datatypes will order differently.
In the case of large scale distribution across a LOT of machines (which is what you invariably have to go to once you expend the options in big iron), that means a huge amount of network traffic, with each retrieval needing to be resolved for consistency due to if you want AP. If you rely on CP, then your join is dependent on SO MANY nodes correctly communicating that you become extremely exposed to network partitions, retries, etc.
Thus you either shard your data so all data is on the same machine (but your joins are necessarily subsets of the overall data: only the shard), or you don't and prepare for extremely bad performance, unreliable performance, or approximations of correctness.
Anyone know what these are?
Shipping multiple databases for multiple services does not necessarily mean multiple database servers. In the database world I saw large servers with many databases more often than multiple servers of one database each; in the second case it was the vendor laziness (cannot give the name), not than a reasonable business or technical reason. When I asked about it the answer was "we'll consolidate in the next release".
Fauna, Upstash, DynamoDB, CosmoDB, Firestore, Planetscale.
The reality is there are limited resources to work on any given project and "rewrite" is generally not the correct way to fix a given problem.
Especially since the first thing you are going to need to do is show that a network system isn't resilient to GC pauses or that GC pauses are frequent enough to impact throughput (very very few applications are actually materially impacted by sub second spikes in latency).
As for the subsecond spikes in latency, these tend to multiply in a distributed system. If serving a client request takes N internal requests, the likelihood of hitting a GC pause somewhere is much larger than if you did only one local request.
Not sure about Go, but none of the "free" Java's GC guarantees low pauses. There is STW fallback even in the most advanced ones like ZGC. So you never know when it stops for more than 1s.
Even if it is a hard guarantee, then, from the link you posted, it is not even generational, so it will scan the whole heap quite frequently, and that is going to influence average performance quite visibly - you definitely dont want a database system to access all its cached memory once in a while.
Database systems are really all about memory and I/O management. You shouldn't outsource those core features to a universal algorithm, unless you wish to forgo any competitive advantage (at least in performance department). So this pushes the devs into the off-heap manual memory management territory, where dragons live (at least in Java, again - maybe Go is better in that regard). I've been there, and I don't recommend.
False. All popular databases are written in C, and continue because of the GC issue. I would not use a general-purpose database written in Java because of GC, for example. We'll see how well Go works in practise.
> sub second spikes in latency
Go is supposed to be sub-second GC pause latency, but understand that most SQL queries are sub-millisecond, so GC latency is still a significant issue compared to query times.
Go might be acceptable now for niche databases like column-store for certain use cases, though.
Also, see the excellent comment above about distributed systems and server cache issues. You can't do application performance analysis with GC literally everywhere.
The puerile knee-jerk hatred for C on HN has to stop - almost every software you use is written in C, from scripting languages to operating systems to web servers to databases.
Source: DBA who's worked with current databases, as well as a custom database written in Java with significant (ie. brutal) GC problems that required a total refactor (rewrite) to "work" at all.
Our transactions implementation is our crown jewel. You might want to check your sources.
Also - anyone who says that x database must be rewritten “because GC” is just making an incredibly un nuanced argument about a nuanced problem. People have built production ready databases in both Java and Go. If you care about low/predictable tail latencies then you have a bunch of other more important problems to solve before you worry about the behaviour of a modern garbage collector. For example: how good is your cache hit ratio? How are your synchronous replication protocols affected by grey failures? That kind of thing.
For example, for https://corridorchat.com/, we have a relatively small number of business accounts (tenants), but with a many users per tenant. And new tenants are created relatively infrequently in the scheme of things.
So I have an architecture with a central 'corridorchat central' PostgreSQL database and a scalable number of shard clusters, all managed with Patroni + Consul, and fall-backs that are read-only until they need to be promoted. Consul DNS allows the application to look up either a read-only replica or a write one.
To know what shard a tenant is hosted on, it is necessary to read from the central database. This requires one of the read replicas - and I can create as many of these as I need. Many transactions then require writing to the shard database for that tenant - but since I balance tenants between shards (and have several shard databases per cluster to allow for future scaling if a shard becomes too hot) I can add more shard clusters as needed. Writing to the central database is constrained, but it is a very rare operation, so there is no expected scaling problem there.
I think for most workloads, this approach should work for a good long while.