Scalable OLTP in the Cloud: What's the Big Deal?
muratbuffalo.blogspot.com
muratbuffalo.blogspot.com
From my experience, the vast majority of complications in systems is people not realizing they are asking an OLAP question while wanting parts of OLTP semantics.
I'm also curious how much of this extra complication comes from having a central database that tries to be a source of truth for everything? As an easy example, inventory systems that try and have a cloud source of truth have obvious problems when on the ground inventory issues crop up. That is, a lot of the complication is between the distributed nature of the application and database, sure. But, another large source is the non-distributed abstraction that a central database can represent not centralized activity.
If you could elaborate on this further, I and others are probably very interested in reading more about it.
- OLAP: most queries need most records (aggregations span large swaths of data)
- OLTP: most queries access just a few records
Also, in OLAP, in many cases, you can live with a single-updater model without much trouble, where OLTP, the strength is to have many concurrent updaters (but mostly non-overlapping).
OLTP - You are dealing with the ins and outs of people using stuff to do things. You buy something, you check out something, you some way perturb the state of things. This should not require large amount of row lookups in 99.9% of cases.
- OLTP: write-mostly, index-seek heavy, ~all queries pre-defined up front
Briefly looking, I can't find the books that I thought covered a lot of this in a good way. Will keep looking, apologies.
I can't recommend it highly enough.
[1] https://www.amazon.com/Designing-Data-Intensive-Applications...
However, at global/twitter scale, just basic CRUD operations become pretty hard for OLTP, so scalable OLTP in the cloud is a pretty big deal, especially when you get fun things like phantom read/writes as writes may take time to get to distributed read replicas, so a user may send a tweet and then not see it on their profile for a while.
But global/twitter/reddit etc. scale does matter in that I don't really remember a PHPBB or slashdot or HN swallowing my comments, ever. Even if you get an error you can always just go back in the browser and your input box is there with your comment.
But reddit has been effing atrocious over the last few months. You get "Internal Server Error" when trying to do any voting or commenting a lot and it won't go away until you refresh the page like 17 times (until presumably you round robin get to some BE server that actually works. It's also been swallowing comments, where about every third or fourth comment it accepts the comment instead of throwing an Internal Server Error at you (which would be preferable) and instead it "accepts" it but it will never show up. Ever. Only chance is to copy every comment before submitting, in case you have to re-submit it from scratch.
I think I agree across the board. HN and Stack Overflow are both amusing examples of much simpler architectures that kind of get to the point, though? Such that I'm not clear if we are on the same page, there.
(And I, sadly, have zero experience with Reddit.)
Really? I remember it happening a lot.
> Even if you get an error you can always just go back in the browser and your input box is there with your comment.
That's a relatively new feature and doesn't always work. (Old, I can't speak to new) Reddit's "push a button and it submits without moving off the page" is much nicer.
The problematic case is that very often now it accepts your comment without any errors and... Then nothing. Your comment will never ever show up. The system swallowed it. Whole. Without chewing.
Huh. I've never seen that. Does it not show on your user page even? I've known it to take a few seconds for comments to show up, and mods on certain subs love to silently delete comments (sometimes automatically), but I've never seen them just disappear.
I do also have the case where it disappears but a refresh / going to the user page will show it, i.e. something async on their end was just a few milliseconds too slow or something.
And posting the exact same comment again (I c+p them now) works. So it wouldn't be any automatic removals and it's too fast for manual removals.
https://www.infoq.com/presentations/reddit-architecture-evol...
That and the Instagram architecture video are such eye-opening examples that I always send to juniors when they have doubts about their own abilities. In most cases, we're astronomically no where near Reddit/Twitter/Instagram scale, and yet we avoid a lot of the common pitfalls we see presented in those videos. Side note, I'm shocked at the quality, but hey "Time To Market!"
The Instagram video: https://www.youtube.com/watch?v=hnpzNAPiC0E
I had the same eye opening experience when I first "met" Jira. I was pretty fresh out of uni myself, yet I knew you should never use a business key that people might want to change and always use surrogate keys for your database.
And here is Jira "happily" using Issue Key as both the business key that signifies which project an issue belongs to, which can very obviously change as soon as you move a ticket to another project and as the primary key in their database and thus for everything related to an issue like say issue links.
I was laughing so hard that you couldn't tell I was also crying and cursing at the same time coz I had to deal with that fact.
They've since introduced ids for things but ... duuuudes! WTF! lol
You can reverse this as 'how much extra complication comes from having many disparate databases, so now you push all ACID complications at the company level instead of the software(rdbms)'
It is a tricky problem!
Having a single source of truth is a tremendous simplification, especially for the people.
It should be ideal that the other databases are pure materializations on top of it, but then not matter what you have the need to input more stuff that is not part of the central one, and now you have both problems.
Oddly, I'd argue that the better way to view this is that the central database is the materialization of the data sources scattered around where work actually happens. As such, you have to have constructs that account for provenance and timing introduced there that aren't as necessary at the edges.
Using the language of 'software architecture: thee hard parts' this forces your entire system into a single 'architectural quanta', basically it is a monolith.
There are situations where monoliths are acceptable or even the least worst option.
The 'Fallacies of distributed computing' covers most of these.
It can work better if you use a hierarchical model, but OLTP is almost exclusively the relational model.
ACID transactions become pretty brittle, software becomes set in stone etc...
There are use cases, but the value proposition is typically much weaker
Simply, the mad crushing dash to get the last bit of committed inventory.
Ticketmaster has 50,000 General Admission Taylor Swift tickets and 1M fans eager to hoover them up.
This is a crushing load on a shared resource.
I don't know if there's any reasonable outcome from this besides the data center not catching on fire.
I've been to events that used a lottery. You order a ticket any time in, say, a month window. At the end, they tell you if you actually got a ticket. They have a process for linking orders together so you can choose to get a ticket if and only if your friends do.
I've also been to C3, which (in the second phase) knowingly used a first-come-first-serve three times, putting some kind of lightweight proxy in front of the actual shop that would only allow a certain number of people to access it at a time (this is important because you don't know how many tickets each user is going to order). In the first phase, they use a system of "self-replicating vouchers" for various trusted C3-affiliated groups: one voucher is given to each group, allowing an order; at the end of each day until this portion of the ticket pool runs out, a new voucher is given to whoever made an order the previous day. I don't know the reasons why self-replicating vouchers are designed exactly that way, but it means each group gets to run down their own self-determined priority order and gets punished for inefficiency.
The capitalist approach is, of course, raise the price to twenty thousand dollars or whatever level it takes for only 50,000 people to want to buy a ticket.
On a business side, they drastically lower the request load and make scalping unprofitable by holding a Dutch auction.
I think something like trivial postgres setup can handle this thing..
A good async setup can easily handle 100k+ TPS
If you want to go the synchronous route, it's more complicated but amounts to partitioning and creating separate swim-lanes (copies of the system, both at the compute and data layers)
That's very much an option when it's something this popular - the Olympics I went to did an even more extreme version of that ("Thank you for putting in which events you wanted to see, your card may be charged up to x some time within the next month").
Or you can do it like plane seats: allocate 50k provisional tickets during the initial release (async but on a small timescale), and then if a provisional ticket isn't paid for within e.g. 3 days you put it back on sale.
Ultimately if it takes you x minutes to confirm payment details then you have to either take payment details from some people who then don't get tickets, or put some tickets back on sale when payment for them fails. But that's not really a scaling issue - you have the same problem trying to sell 1 thing to 5 people on an online shop.
Taylor Swift tickets cost $500 each [1]
An AWS u7in-32tb.224xlarge with 896 vCPUs, 32 TB of memory and 200 Gbps network bandwidth costs $407/hour even at list price [2].
The question is whether they feel like going to the expense and engineering effort.
[1] https://www.businessinsider.com/guides/streaming/how-to-buy-... [2] https://instances.vantage.sh/aws/ec2/u7in-32tb.224xlarge
I think perhaps the view of a database being a single server codebase like this is a bit naive. When you read about how Meta deploy MySQL for example, it's a whole service ecosystem that provides different levels of caching, replication, etc, to create a database at a higher level of abstraction that provides the necessary properties. FoundationDB is similarly better viewed as a set of microservices. When you architect a database like that it is possible to achieve these things, but that doesn't seem to be a new idea, that seems to be just how it has been done in the industry for a while now. The article isn't entirely clear on whether they realise this or are proposing something new.
"Scalable"? Needs more context please. Scalable how? And what are you willing to compromise on for that scale?