Partitioning GitHub’s relational databases to handle scale
github.blog
github.blog
What I did not fully grasp - isn’t Vitess it’s own database system built on top of a K/V store similar to e.g. CockroachDB? Here it sounds like they only use pieces of Vitess in front of regular MySQL primaries with sharding?
It’s also curious that they sharded based on domains. I assume at some point the QPS on specific domains (e.g. Gist) will be so high that partitioning no longer is effective?
edit:// Looks like Vitess is indeed a set of tools on top of MySQL that allow for these type of scaling operations. It is not re-implementation of MySQL using different underlying technologies.
refreshing to see large scale providers still relying on „old-school“ tech!
Somewhat OT but Github is really "old school" - they don't have an SPA frontend and render most things server side the classic Rails way, with lots of UX enhancements coming from judicious use of JS and websockets for interactions. That's one of the things I find very endearing about their platform.SPA frameworks aren’t a silver bullet here, and you can definitely fix these without using one, but these bugs are definitely much harder to make since the DOM is rendered from a single source of truth instead of being patched piecemeal.
They make a clear distinction between "vertical partitioning" and "horizontal partitioning", I think there are probably bigger wins to be had, and probably bigger challenges in the horizontal partitioning than the vertical.
All of this work only got them a 2 factor reduction in load, whereas with an ideal horizontal solution, they could scale it infinitely.
"In addition, FPGA hardware is used to accelerate the compaction process, further maximizing the performance of the system. This marks the first time that hardware acceleration was applied to the storage engine of an OLTP database."
[1] https://www.alibabacloud.com/blog/new-insights-into-x-engine... [2] https://dl.acm.org/doi/pdf/10.1145/3299869.3314041
Centralized vtgates proxy to vttablet services which run on the same host as MySQL, the tablet then queries the local MySQL server. On top of this, a lot of magic can be built
https://docs.microsoft.com/en-us/sql/relational-databases/pa...
They are currently contributing to both Vitess and Rails. I am guessing they are also helping / a paying customer of PlanetScale.
Which means all the improvement and edgecase are battle tested on Github and upstreamed. I much rather they continue their current path. So others could enjoy the Github Stack.
I know HN hate Oracle and MySQL. But generally speaking I think Oracle has been doing a great job in Java and MySQL development.
Do you still think so after reading this: https://news.ycombinator.com/item?id=18442941
I think this is reading far too much into nothing. He writes 'bug' correctly about 5 times. I'm not an expert in psycholinguistics, but I suspect there's a phenomenon where you can mangle your internal pronunciation of the word and hit the wrong vowel.
Could be. I could have left that somewhat tenuous side note out, my argument doesn't rest on it.
Is it? Can you give me an example of the top notch engineering from "modern" Oracle?
I mean this as a genuine question, I am not that familiar with recent work coming out of Oracle.
The database harkens back to when computers were new. There's a ton of money that goes into continued development of it, and is vertically integrated, including custom hardware and the software to run it. It's extremely expensive hardware - they still sell SPARC servers, running Solaris if that's what you want (but they do also support a custom Linux kernel for their hardware).
It's such a high end niche that there's only a handful of companies that can even run the benchmark competitively because it just costs so much in hardware to play at that level, which makes it very opaque unless you're fluent in a lot of terms, some of them proprietary, others not. Eg https://blogs.oracle.com/exadata/post/exadata-uses-persisten... it's an absolutely fascinating journey into getting better SQL performance that involves some really high end shit, and (like kubernetes) most of the people out there just don't play in the same league. Which isn't a judgement against them and their needs, but it costs a lot of resources to wring microseconds more performance from a multi-million dollar machine. An AWS EC2 cloud VM, this ain't.
Alternatively, you don’t see benchmarks because Oracle’s licensing bans posting benchmarks - https://www.brentozar.com/archive/2018/05/the-dewitt-clause-...
The ability to tolerate machine failure with zero downtime and zero data loss.
Sure, Shenandoah (started at RedHat) is even more brutally mindblowing, but it has a constant overhead (and maybe obviously, maybe not, but it builds on the already existing pretty good GC infrastructure in the JVM).
So at least it seems the Hotspot Group is left mostly alone to do their high quality work.
What github did instead was to separate out whole sub-systems from the main db cluster to separate clusters so that requests could hit completely separate engines so that the number of connections to each would be reduced and the chance of breaking something by a mistake in a cross-schema query is reduced by physical separation.
I think you could still do this with SQL Server but why bother if you already use MySQL and the tools exist there.
> Building on top of schema domains, two new SQL linters enforce virtual boundaries between domains. They identify any violating queries and transactions that span schema domains by adding a query annotation and treating them as exemptions. If a domain has no violations, it is virtually partitioned and ready to be physically moved to another database cluster
Table-level partitioning doesn't help with this (AFAIU), as queries access multiple tables anyway (without app-level changes).
The standand db-level feature closest to what they're doing, if somebody "wants to try this at home", is probably tablespacing (or separate dbs, in the next step).
I see some inspiration from microservices (separation of models/storages), except that they're (I suppose) keeping the monolith approach.
I am glad it isn't me trying to split it up :-)
Vitess for example basically follows this principle. 2PC/Cross-shard in Vitess can be done and you could make it faster, but essentially all users and the developers instead take the view that applications designed to scale with Vitess should instead avoid 2PC at the design stage. Why do they take this view? Because their real-world experience is that sharding scales forever and cross-shard transactions don't, no matter the effort. It might scale far, but never "enough." This kind of experience is what drives such a design philosophy.
There's of course a feedback loop to all this -- nobody will use 2PC if it's slow, so therefore it's avoided, and nobody will patch it to be fast, because you can just avoid it instead of spending time on that, etc etc. Also, designing for such a system up front is maybe more difficult, or less familiar. But I think it should be noted historically and contextually that this whole push for what is effectively strongly-consistent databases on the WAN where strong transactions won't obliterate performance is a relatively recent trend. You can build something like this bespoke if you're careful but in terms of COTS tools, well, that's been very limited until recently. So the institutional knowledge that can be carried around for engineers is all built around a different set of assumptions, like "don't use 2PC."
My own personal experience talking to people "in the wild" (whatever that means) is that the Vitess philosophy is relatively popular at a lot of high-scale places, and probably isn't going away anytime soon. It helps that you can, you know, just use tools like Vitess. But the push for strongly consistent WAN databases (Yugabyte, CockroachDB, etc) isn't slowing down either, so...
Does anyone know how a service like Github would store & backup the Git repo itself? Not in MySQL, I imagine?
I don't see anywhere that we've talked publicly about how backups work. If there's interest, I can see if someone wants to write a blog post.
Disclosure: I'm the GitHub product manager for Git systems including storage and protocols.
There is! Storage, backup & restore at that scale are always interesting :)
Especially considering Gitlab developed Gitaly due to scale and redundancy issues with a regular filesystem.
For a brand new project without legacy code, newer generation databases become an easier sell.
The area I really want to chew on is partial indexes, now that it's no longer a niche feature. I'd like to have a more informed opinion of how far you can scale a system without resorting to sharding, just by use of more efficient indexes. Especially on read-dominated problem domains.
Glad folks from GitHub shared this story. As a database practitioner managing couple world largest database fleets in the past and see the scale a single db can support, I am 100% convinced that only a very very tiny fraction of application needs sharding.
Using "boring" technology to build an innovative product, not the other way around.
I wonder if people see it as ripping the bandaid off to go straight to sharding. Or perhaps being multitenant introduces that idea. Seems like everyone, or at least the louder ones, discover that they have one 'whale' who gets to be too big even for a single database, or buys one of your other customers, and then where are you?
> Joining data in the application instead of in the database is another common solution. [...] In some cases, this leads to surprising performance improvements.
So, maybe such approach is not impossible for a more heavily connected dataset, although probably prohibitively expensive if you have to re-write most queries into separate-queries+in-app-join.
See https://api.rubyonrails.org/classes/ActiveSupport/Notificati... for the API to subscribe to these notifications.
What did the virtual partitioning add that those other techniques could not provide?
If we used separate db users as you're suggesting, any query that we didn't catch beforehand (e.g. via our CI builds) would cause noticeable problems for our users, which is something that we want to avoid.
Additionally, switching to a separate user account would require holding open twice the amount of connections to each database server (old db user plus new db user), which probably would be fine but is still a lot of additional connections at our scale.
1. wonder what happens on the edge of the boundaries when a table does need data from another domain. And what if that domain/cluster is down?
2. How do they physically connect to the cluster? A seperate db connection?
Impressive numbers. Busiest master I worked on has been 15,000 qps.
(SSD can do 1+ million IOPS and support 2 mysqld server processes at 100,000 qps.)
Source: DBA.
1. Migrating is a huge, difficult, expensive project at Github's scale -- almost as bad as a rewrite
2. "Can handle more data and load" doesn't make sense as a blanket statement
How much data and load Postgres can handle vs MySQL is determined by lots of variables. It's neither absolute nor consistent. It depends more on database and infrastructure design than on the underlying RDBMS.
The best reason to choose Postgres over MySQL is not hardware efficiency but rather developer efficiency, although some projects get both.