Lessons learned from sharding Postgres at Notion
notion.so
notion.so
I looked at a bunch of options but ultimately managing a sharded database seemed like too much overhead for our small (but expanding!) dev team - so we decided to move our heavily loaded tables to CockroachDB instead, as since it's mostly Postgres compatible our transition would be easier, and it would be much easier to manage as it automatically balances load and heals. We're running cockroach on GKE and it's a really nice fit for running on kubernetes.
Ended up working well for us - we still have our smaller tables on PG but we want to move it all to Cockroach over time. There were some teething issues as we learnt what to monitor on it (read amplification is the big one here, though pebble has helped a lot vs the previous rocksdb), and it does need a lot of hardware at high scale (we're running 50 16-vcpu nodes for it) but overall I'm happy with it.
I'm also curious how you have found dealing with legacy schemas and queries that might cause a lot of data shuffling due to poor data locality. Was it necessary to put a lot of time into the sharding strategy? Have you experienced new query hotspots?
We did have to make some schema changes in order to improve the data locality and avoid hotspots - CockroachDB did _work_ without doing those, but the changes improved performance massively. It wasn't super time consuming though, and the admin UI/monitoring in cockroach is nice for showing you any hot ranges.
The biggest one was our table which stores annotations on a document - on PG it's primary key was only (annotation_id UUID), and we had a secondary index on (document_id). This meant that on Cockroach the data for a document was spread across many nodes and querying all the annotations for a document would take ~40ms. We changed the primary key to (document_id, annotation_id) which co-located the data and that came down to ~2ms. Composite primary keys aren't ideal for ActiveRecord but the performance win was worth it.
We also had some tables with hotspots which was mostly just a matter of making sure they used uuid keys. It would be a problem though if we had a table with individual rows which were super hot or where the table didn't have enough data to get split into multiple ranges - though cockroach does have load based sharding to help here. It was only really necessary to worry about the higher load tables here, if it's below ~2000 qps or so it won't matter.
This is a difference from PG, where having all the writes be concentrated in the key range is helpful as the Btree nodes covering that will be more likely to be highly cached
In a past company we managed to run a bank on Cassandra with eventual consistency. It was ... well, it was good CV experience for the DevOps guys. (I also didn't need much encouragement to take the 'no getting drunk while on-call' rule very seriously.)
That said some of the things you can do from day 1 make things way easier when the time comes later. Workspace ID/Tenant ID/Customer ID, having your data grouped that way makes it very shardable and saves you a lot of work later on in the case you didn't have that in your data model.
I'm not sure I buy that Citus/Vitess are magic, both are reasonably clear how they work and you can dig into it. At the same time I'd weigh the downsides of Citus (can't speak to Vitess) in that the online rebalancer isn't open source so at that point it's a proprietary product.
Dumb question, if someone starts off by grouping their data only by Customer ID and then later needs to shard. Couldn’t all of the sharding problems go away if they simply created a new Customer_Tenant table to map customer_id into tenant_group_id
If you have that tenant discriminator on all tables then it's easy to route it to the appropriate physical and logical shard right away vs. having to do some DB join first in your request to figure out where it goes.
Is there a specific term for the trade off of denormalized vs more easily query/shard/etc?
No, it's pretty much just called that. A good DBA will be able to strike the right balance between normalization and performance using their intuition and experience, which is one of the reasons they're often paid very well despite working in relatively ancient ecosystems.
I wonder if this could or should be a built-in feature in databases. It is "meta-data" meaning data about data, who owns it.
The big problem everyone I have talked to about sharding runs into is managing the shards as you expand. In this case it looks like notion over sharded so they can spin out up to 480 physical nodes, but when they need the 481 it is going to be a nightmare, thats what Vitess gets you for free, expand to any number of shards and just never worry about it again
It's been a while since the last attempt so I forget the details but I think the reshard replication lag was growing no matter what resources we threw at it. We've since scaled vertically instead and are working on architecture changes to reduce our MySQL load (which admittedly is extremely high).
Our dataset size is actually reasonably small (10s of TB), but our transaction throughput is very high. Glancing at our dashboard, baseline is 250k/s, with sustained daytime load in the 500k/s region.
If you wait until you have to, then isn't it too late to plan and test your solution? So you suffer from degraded performance and/or outages while you rush out a solution.
This is just a specific case of premature optimization.
That's a good point. Sharding has an interesting trade-off. Early on, you don't need it and there's an overhead to distributing work across multiple machines (you're taking an additional network hop). Later when/if you need it, sharding becomes painful to introduce. You may need to change your data model for performance and move data.
For years, I also cautioned against introducing sharding as a premature optimization. Recently, we changed our approach and also introduced sharding on a single VM. I think that offers a pretty good trade-off. If you're interested in this topic, I’d be curious to hear what you think.
https://www.citusdata.com/blog/2021/03/20/sharding-postgres-...
Off hand this seems like an almost worst case for PG. Since the updates to blocks could contain large data (causing them to be moved often) and there is one big table; it seems likely that the blocks for a single notion document will end up being non-continuous on disk and thus require a lot of IO/memory trashing to read them back out. PG doesn't have a way to tell it how to organize data on disk so there is no good way around this (CLUSTER doesn't count, it's unusable in most use cases).
Arm chair engineering of course - but my first thought would be to find another storage system for blocks that better fits the use case and leave the rest in PG. This does introduce other problems, but it just feels like storing data like this in PG is a bad fit. Maybe storing an entire doc's worth of block entities in a jsonb column would avoid a lot of this?
for instance, in SQL:
User table, that has columns: id, email
Article table, that has columns: id, text, user_id
KV/noSQL equivalent:
Article document, that has properties: id, text, user_email
Aren't tablespaces (https://www.postgresql.org/docs/10/manage-ag-tablespaces.htm...) supposed to help with that?
Haven't used them, I'm honestly curious
What I was talking about is controlling the ordering of the rows within a table on disk. If you are going to be reading some group of rows together often, ideally you want those rows to be contiguous on disk as a sequential read of a range is much faster than bouncing around to dozens of locations to collect the needed rows. This becomes more important for very large tables. Imagine a 5TB `blocks` table and you need to read 50 blocks to render a given notion doc but those blocks could be scattered all over the the place on disk, it's a lot more work and it thrashes the page cache.
PG doesn't normally make any guarantees about how rows are ordered on disk and it may move rows around when updates are made. It does has a CLUSTER operation, which re-orders rows based on the order of an index you give it, but this is a one time operation and locks the table while running. This makes it functionally useless for large tables that are accessed and updated frequently.
Some other databases do give you control over disk ordering, SQL Server for example has `CLUSTERED INDEX` which you can apply to a table and it'll order data on disk based on the index order, even for new insertions / updates. It does cost a bit more on the write side to manage this, but it can be worth it in some cases.
This is an important point, doing things divisible by 12 gives you a lot of flexibility. It’s not a coincidence both time (clocks) and degrees (360) are multiples of 12.
There is an oeis sequence of them that starts "1, 2, 4, 6, 12, 24, 36, 48, 60, 120, 180, 240, 360, 720, 840...", which notably does not include 480. https://oeis.org/A002182
Edit: Actually, this was (also) already done by Ramanujan and has a better name than any proposed here: http://oeis.org/A067128 - "largely composite numbers" which does in fact contain 480. Perhaps every sharded system should choose a shard count from this list?
Obviously not this thing just by itself, but this sort of knowledge, applied to everyday decisions that get compound over time.
7! = 5040 shows up as well in the OEIS sequence. Buuut, 8! = 40320 does not...too many useless 2's in the prime factorization?
The only downside imo is not having ten(s) be a balanced deploy number.
I’m happy to answer questions about the project here, feel free to reply below.
Could you have selected some workspaces with lower traffic to migrate first? That would have decreased the load on the primary, potentially speeding up the replication, which is a flywheel to enable more customers to migrate to shards.
At the end of the day, this is something we could have explored in more depth, but we were ultimately comfortable with the risk tradeoff of migrating all users at once vs. the consequences of depending on the monolith for longer, largely thanks to the effort we put into validating our migration strategy.
[0] https://blog.sentry.io/2015/07/23/transaction-id-wraparound-...
- Can you share details on the routing? I.e. how does the app know which database + schema it needs to go to for given workspace?
- Did you consider using several databases on the same Postgres host (instead of schemas within a single database)? Not sure what's better really, curious whether you have any thoughts about it.
Thanks!
All in the application layer! All of our server code runs from the same repo, and every Postgres query gets routed through the same module. This means that it was relatively easy to add a required "shard key" argument to all of our existing queries, and then within our Postgres module consult an in-app mapping between shard key range and DB+schema.
Plumbing that shard key argument through the application was more difficult, but luckily possible due to the hierarchical nature[0] of our data model.
> Did you consider using several databases on the same Postgres host
If I recall correctly, you cannot use a single client connection to connect to multiple databases on the same host, and so this could have ballooned our connection counts across the application. This is not something we explored too deeply though, would love to hear about potential benefits of splitting tables in this way.
Luckily we did not have many join queries prior to sharding. The few cross-database joins were trivial to implement as separate queries in application logic.
Our main concern here was referential consistency: if you need to write to multiple databases at once, what happens if one of your writes fails? This was not a problem in practice for us since (1) our unsharded data is relatively static and (2) there are very few bidirectional pointers between data across databases.
Long term there are more interesting problems to solve when distributing unsharded data globally. However, given that our unsharded data is less dynamic and consistency is less critical, we have many levers we can pull here.
1) Can you please talk more about the risks of using DynamoDB or a similar NoSQL solution?
2) Did you consider Spanner, which is the SQL DB of choice within Google and is available on Google Cloud.
Thanks for the wonderful engineering blog posts!
We were on a tight timeline due to impending TXID wraparound. Switching database technologies would have required us to rewrite our queries, reexamine all of our indexes, and then validate that queries were both correct and performant on the new database. Even if we had the time to derisk those concerns, we'd be moving our most critical data from a system with scaling thresholds and failure modes we understand to a system we are less familiar with. Generally, we were already familiar with the performance characteristics of Postgres, and could leverage a decent amount of monitoring and tooling we built atop it.
There is nothing inherent about non-relational DBs that make them unsuitable for our workload. If we were designing our data storage architecture from scratch, we'd consider databases that are optimized for our heaviest queries (vs. the flexibility afforded by Postgres today). A large number of those are simple key-value lookups, and a plain key/value store like DynamoDB is great for that. We're considering these alternatives going forward, especially as we optimize specific user workloads.
Re: Cloud Spanner: we didn't consider a cross-cloud migration at the time due to the same time constraints. Still sounds like a wonderful product, we were just not ready at the time.
At the cost of potentially introducing more cross-partition queries, you might benefit from splitting up high-throughput workspaces. See strategy in https://d0.awsstatic.com/whitepapers/Multi_Tenant_SaaS_Stora..., pages 17-20.
You can subscribe via email (no thanks), and I realize there are some good ways to turn an email subscription into a feed (Feedbin.com handles this nicely). Still, I thought some public shaming here might encourage the to build this basic blog feature (that would help people follow their company!).
Notion chose to do manual sharding (aka application-level sharding). That's what you end up doing if you didn't choose a database that has built-in sharding from the get go, because it is extremely hard to switch to a different database technology. (Larry Ellison compares it to a catholic marriage -- there is no divorce!)
I skimmed through the article to find the critical piece of info I was looking for: rationale for doing manual sharding. The rationale supplied in this article is "we wanted control over the distribution of our data." That's a weak explanation. It's the kind of thing you say to justify the bad choice made earlier on: failure to choose a database that supports automatic sharding from the get go.
> By mid-2020, it was clear that product usage would surpass the abilities of our trusty Postgres monolith, which had served us dutifully through five years and four orders of magnitude of growth.
This should answer your question. Just use a standard database and get on with coding features instead.
If you are successful enough to get to the point that sharding PostgreSQL becomes your bottleneck, you've won.
In addition, you really can't solve scaling problems up front. Where your bottleneck actually occurs will differ from where you think it will.
Better still, choose a "standard database" that supports sharding.
We will be happily building out features while your engineering team wastes time solving for scaling problems that you will never have because you don't have features your customers want.
Choosing database A instead of database B doesn't mean you're suddenly "solving for scaling problems". It just means you're better prepared to one day solve scaling problems, should your product take off.
I would see an argument that Notion's structure should have made this problem more important for them to resolve earlier in the process (like some document-based DB), but Postgres works well!
But if you're developing an app that scales to internet users consider Couchbase or some other database that matches your requirements.
If you had some behemoth with 32TB ram and 1PB storage could all of motion fit?
I’m obviously ignoring the obvious single point of failure here, but sometimes the simplicity could be worth it if you’re willing to handle that.
I’d be curious to hear about a site that’s in the Alexa 1000 architected how I describe.
Postgres has a bunch of low level locks and buffers (protected by locks) which are essentially single threaded. So even if you had a 500 CPU instance at some point you'd not be able to get more throughput out of it.
Of course it also depends on what you are doing. Large tables with high update rate are harder to handle than large tables which are insert only. My personal opinion is that Postgres tables with >100GB data (without indexes) are a starting to be a pain to work with, no matter how much RAM, CPU or IOPS you have.
There are fast database engines specifically designed for single servers with 1PB of storage, and these servers commonly have hundreds of cores. This is much more efficient for some use cases. It is the kind of thing you find on edge platforms designed to manage and operate on sensor data models. You build database engine internals very differently when working at this storage density and number of cores; many good architectural ideas in databases at much smaller scales become severe bottlenecks even on a single server.
Also, are you referring to a single-instance of Postgres, or that Postgres would struggle with 10TB even with something like Citus or TimescaleDB?
Extreme scale-up database engines are relatively rare and usually bespoke, I am not aware of anything close in open source. They only have one major advantage -- they are easily deployable outside the data center, e.g. at the edge, because it is a single self-contained box and relatively power efficient. Most people opt for scale-out because 1) there are many to choose from in open source and 2) it will work just as well as extreme scale-up for most applications, though you'll need more hardware for the same workload.
I've started to appreciate the utility of the extreme scale-up engines in real-world applications, though I've mostly worked on scale-out engines. They solve a real problem at the edge for data intensive applications but most currently people deploy in the data center so there isn't a pressing need.
I'm founding engineer of a new startup now and I advocated for Guru. https://www.getguru.com/
I like how information is marked stale/verified and the permission are much tighter. The data ladder is also better structured. I'll see how it works as our headcount increases.
Notion has a lock feature [1] on a document-level.
[1] - https://www.notion.so/Lock-page-content-d2b995727c0b483f9f35...
Doing big surgery takes time, meanwhile your workload continues to grow and grow. It’s not uncommon for a major win to only reset the clock by 6-12 months, and then you either have to run to the next one, or you have two teams working on separate angles at the same time. All the while people are learning as they go because the business settled on catch-up instead of capacity planning.
Being a little successful can be tough.
Search is still pretty slow.
Im surprised at the 500gb table / 10tb db sizes listed.
I start to worry when my db hits 1tb and any single "hot" table is over 100gb.
I know there isnt a one size fits all answer but I'm genuinely curious if I'm just too behind living in lessons learned ~5 years ago and and we are in a better place now.
Does anyone with experience have any thoughts in favour or against such implementation?
It's way more powerful than their official API. The official API can only interact with databases, which (as anyone who uses Notion knows) is a tiny subset of the overall things you want to use Notion for.
The unofficial Python API lets you have complete programmatic access to all notion pages.
The tradeoff is that certain things break. However, after reading the code, I found it quite easy to fix the problems I ran into (https://github.com/jamalex/notion-py/pull/345) and I suspect you'll be able to do the same. (I empathize with maintainers not prioritizing their open source work ahead of family life, business, etc.)
Was surprised that it was even possible to have a nice Python class for every possible Notion object, let alone control them and update them on the fly. I wish their official API would be as flexible. Maybe one day.
It seems like you made the database sharding decision as a result of vacuum taking too long. Partitioning the problematic table (assuming you're storing "blocks" in 1 giant table) would have enabled per-partition vacuuming, which will avoid long running vacuum processes.
May have saved months of planning that may never have been required.
You can't base a huge decision like that on "well it worked for somebody else". How did it work for them? What was their platform? What was their data profile? What were their requirements? What was their acceptable level of risk? It could be that they have so many caching layers and lose so little money that they are fine with the database being down for 5 days while they recover. Or that it doesn't actually work that fine for them, because who wants to broadcast that their system is shitty? Or that they've just been lucky.
It's not worth finding out the hard way. Build the best system you can with the time and money and expertise you have. Don't cheap out just because you think you can get away with it - especially on the critical stuff.
- Many concurrent connections to a single database which Postgres has traditionally not handled well (though improved in recent versions).
- You're now on the hook for writing a database control plane.
- Backup and restore is much flakier since the data volume is so much larger. Lots of weird shit starts happening when you download tens of TB.
- In general, everything is a harder once you start nearing machine limits.
* How do you setup hot stand-by for each database ?
* Do you have coordinated backup for all the databases, or are they backed up individually ?