12 ms·
How does database sharding work?
- darwinwhy 4y agoThe first link in the article [1], to a blogpost about the etymology of the definition of "shard" that we now think of as the primary meaning of the word, is super interesting. The release of Ultima Online doesn't really seem too close to any inflection point of the word on Google Ngram, but I'm not sure exactly how close we should expect it to be to 1980 or 2000. 1. https://www.raphkoster.com/2009/01/08/database-sharding-came-from-uo/ https://www.raphkoster.com/2009/01/08/database-sharding-came...
- dikei 4y agoOne thing I notice is the over-usage of sharding, especially hash-based, might turn your Relational Database into just another key-value store, with consistency constraints moving into application code, and you lose many advantages of a traditional RDBMS
- paulryanrogers 4y agoWhat's the alternative as things get too big? Sharding by date? By client? How do you prevent hotspots?
- preseinger 4y agothe only real answer here is to shard by customer optimistically, you can try to shard by read use case, but that's never gonna be stable over time if you need a true multi-tenant system you can only shard by individual entity and move all of the composition logic to the next layer up, there's no way to cheat
- AdieuToLogic 4y ago> the only real answer here is to shard by customer No. Pick a a stable, guaranteed-to-exist, shard key (composite or atomic properties) and use that. If composite, order properties used in most-to-least distinct value distribution. > if you need a true multi-tenant system you can only shard by individual entity and move all of the composition logic to the next layer up, there's no way to cheat This is incorrect. Sharding and multi-tenancy are orthogonal concepts.
- preseinger 4y ago> No. Pick a a stable, guaranteed-to-exist, shard key (composite or atomic properties) and use that. If composite, order properties used in most-to-least distinct value distribution. what?? if you shard users by user ID and orders by order ID, then a query that joins a bunch of user orders in the same tenant namespace will spread across multiple user shards and multiple order shards value distribution doesn't really have any impact here (shard keys are also guaranteed to exist by definition, not clear what you mean by that) if you don't care that specific queries cross sharding boundaries, okay, then no problem, but in that case sharding is not solving the problem that we are talking about here > Sharding and multi-tenancy are orthogonal concepts. sharding and multi-tenancy are only orthogonal if you don't care that a single tenant can have information on multiple shards
- AdieuToLogic 4y ago>> No. Pick a a stable, guaranteed-to-exist, shard key (composite or atomic properties) and use that. If composite, order properties used in most-to-least distinct value distribution. > what?? > if you shard users by user ID and orders by order ID, then a query that joins a bunch of user orders in the same tenant namespace will spread across multiple user shards and multiple order shards Note my recommendation of picking a "stable, guaranteed-to-exist, shard key." If there is a users table/collection sharded by its id and an orders table/collection sharded by its id, then there is no "guaranteed-to-exist shard key" between them, right? So, in that case, where the two are often accessed together, having a "guaranteed-to-exist shard key" of the "tenant namespace" would be the logical choice. > value distribution doesn't really have any impact here I mentioned value distribution strictly in the context of ordering composite shard keys (if applicable). My apologies for any confusion this might have introduced. > (shard keys are also guaranteed to exist by definition, not clear what you mean by that) My implication was in reference to what is always available across accessing sharded entities. In the scenario you describe, sharding by either user or order id would not be ideal. In situations where one or the other is not sharded, then identifying a common shard key likely is not needed. >> Sharding and multi-tenancy are orthogonal concepts. > sharding and multi-tenancy are only orthogonal if you don't care that a single tenant can have information on multiple shards My assertion was regarding theory, not specific scenarios. In practice, having multiple tables/collections needing shards with the tenant being the common entity strongly implies sharding on a tenant property and not of those unique to each table/collection, as implied in your reply.
- beoberha 4y agoYou’re exactly right. I work for a large cloud database service and the vast majority of our top customers shard by customer. This also gives you the benefit of using higher levels of sharping abstractions that map to performance SKUs for more demanding customers, allowing much more efficient allocation of COGs.
- AdieuToLogic 4y agoPick a a stable, guaranteed-to-exist, shard key (composite or atomic properties) and use that. If composite, order properties used in most-to-least distinct value distribution.
- dilyevsky 4y agoNot so easy at scale because there’s additional requirement for it to also not be susceptible to hotspotting
- esafak 4y agoDistributed databases like Yugabyte, Cockroach, TiDB, etc.
- paulryanrogers 4y agoSorry if I was unclear. Hoping to hear what strategies were practical within a given DB. A whole different platform is usually a much larger undertaking.
- dikei 4y agoPass a certain point, you ought to think about whether to keep using a RDBMS as a K-V store or switch to a real distributed K-V store like Cassandra, ScyllaDB, DynamoDB and the like About hot spots, it has always been an issue with K-V stores, and the only real solution is a good key design, though there are some tricks: * Use a uniformly distributed but deterministic key prefix. For example, instead of using raw user_id as key, attach a small hash before it: (hash(<user_id>), <user_id>) This can help with load distribution if your <user_id> is not uniformly distributed by itself such as a phone or id number. * Add more data to your key to increase cardinality. For example, with time series data, instead of using object_id as partition key, use (user_id, time_bucket) so the data for a busy object will get split into different partition over time.
- charcircuit 4y agoYes, this is where caching layers come in. Now your cache acts as the RDBMS.
- dehrmann 4y agoOnline use cases always(?) scale out to key-value stores. Offline use cases almost always scale to distributed, date-partitioned columnar stores.
- samsquire 4y agoThanks for this article planetscale. I've been enamoured with sharding recently but more for multithreaded performance. I want multimaster postgres. I started trying to write a postgres synchronizer, by sorting every data by row and column and hashing the data of every column and row, then doing a rolling hash of the data. In theory, two databases can synchronize by sending the final hash of their data and then doing a binary search backwards until the hashes match. This way you can work out which parts of the databases differ and need to be transmitted. If the databases are identical, very little data gets transferred. One problem I've not worked out is how to decide which database is the winning database without having to change application queries. If you synchronize two multimaster Postgres databases that have had independent writes to different sections, how do you identify which database is the source of truth for a column/row combination?
- preseinger 4y ago"a binary search backwards" is not well defined and doesn't guarantee any upper bound on consistency sharding data based on individual rows within tables is tricky, you won't get reliable guarantees for queries in this way "which database is the winning database" is a function of individual rows, a query that reads data from N different row-owners needs to query N different instances, or else accept that it will work against stale data
- AdieuToLogic 4y ago> One problem I've not worked out is how to decide which database is the winning database without having to change application queries. I believe this is a case of the "Two Generals' Problem"[0], which implies that there is no provably correct solution to achieve this. > If you synchronize two multimaster Postgres databases that have had independent writes to different sections, how do you identify which database is the source of truth for a column/row combination? You can't without a quorum[1] and even that does not guarantee success. 0 - https://en.wikipedia.org/wiki/Two_Generals%27_Problem https://en.wikipedia.org/wiki/Two_Generals%27_Problem 1 - https://en.wikipedia.org/wiki/Split-brain_(computing) https://en.wikipedia.org/wiki/Split-brain_(computing)
- 4y ago
- phamilton 4y agoA favorite resource: https://learn.microsoft.com/en-us/azure/architecture/patterns/sharding https://learn.microsoft.com/en-us/azure/architecture/pattern... Microsofts Azure Cloud Patterns is some of the best documentation out there. It's not Azure centric. It focuses on why you may want to do something and describes the techniques that are commonly used to solve it. It's pretty great.
- shortrounddev 4y agoMSDN is like the Wikipedia of coding problems. You documentation for things that have nothing to do with MS there
- beebmam 4y agoI also couldn't recommend this higher. If you are designing a new application, read these docs. You will almost certainly learn something deeply useful that you'll carry me with you for the rest of your career.
- Douger 4y agoJust want to say thanks for pointing out this resource. Will make for some great long weekend reading!
- andirk 4y agoSometimes I wonder how many other industries allow a worker to stumble upon a document that they then read outside of work, despite the details often being kind of difficult to understand at first, frustrating even, but we do it literally as a pastime, for pleasure. I love that.
- revskill 4y agoThe main issue is, i can't stand C# or OOP syntax to illustrate patterns. Why class here ? For God sake, please use simple functions to prove the points.
- cpurdy 4y ago
- Dwedit 4y agoShards are the secret ingredient in the webscale sauce. They just work.
- creata 4y ago(source: https://youtu.be/b2F-DItXtZs?t=143 https://youtu.be/b2F-DItXtZs?t=143)
- mamcx 4y agoAnd the best sharding? Doing multi-tenant. Sharding at table level is very complex and expensive. Fully give a single DB per tenant is very practical, and the reasons to do sharding (like reporting) is where you do the other fancy things (like ship events in to kafka, etc, etc). Also, most issues of scalability are dominated by a few tenants that consume most resources, and distribute the loads is more easy per-tenant.
- canadianhacker 4y agoI think we need intuitive tooling to let developers continue just thinking in multi-tenant terms, with single-tenant behind the scenes.
- deleted 4y ago[deleted]
- stn_za 4y agoIssue is when your primary tenant(s) are much larger than the other...eventually vertical scaling is difficult even for only 1 tenant, depending on it's size.
- bsaul 4y agoI've been doing that as much as possible. however you're still left with availability issues, such as replicating / failover management etc. Which really are orthogonal issues to tenancy. How do you manage that ?
- mamcx 4y agoYou need at least 1 backup and/or replica. Availability is kinda good: You could lose one/few tenants but that not take down all the rest (if they are isolated properly). Is MUCH worse if all the db made the fancy way get down, or worse, you get a cascade of latency and/or crashed by the interlinked nature of the "scalable architecture that is fact share-everything global singleton". And you can get a bit fancy (not done it myself, my customers tolerate a bit of downtime) of fork the tenant and upload it another node.
- 4y ago
- GeneralMayhem 4y ago> I suppose the more fundamental question is: why are you not using a database that does sharding for you? Over the past few years the so-called “serverless” database has gotten a lot more traction. Starting with the infamous Spanner paper, many have been thinking about how running a distributed system should be native to the database itself, the foremost of which has been CockroachDB. You can even run cloud Spanner on GCP. This is the most important paragraph in the article. In a world where things like Dynamo/Cassandra or Spanner/Cockroach exist, manually-sharded DB solutions are pretty much entirely obsolete. Spanner exists because Google got so sick of people building and maintaining bespoke solutions for replication and resharding, which would inevitably have their own set of quirks, bugs, consistency gaps, scaling limits, and manual operations required to reshard or rebalance from time to time. When it's part of the database itself, all those problems just... disappear. It's like a single database that just happens to spread itself across multiple machines. It's an actual distributed cloud solution. My current employer uses sharded and replicated Postgres via RDS. Even basic things like deploying schema changes to every shard are an unbelievable pain in the ass, and changing the number of shards on a major database is a high-touch, multi-day operation. After having worked with Spanner in the past, it's like going back to the stone age, like we're only a step above babysitting individual machines in a closet. Nobody should do this.
- grogers 4y agoAll the "auto-sharding" DBs have their own quirks, especially around hot key throughput. You often end up having to add "sharding bits" to the beginning of your keys to get enough throughput. The size of one partition is usually tiny compared to one partition of an RDBMS too. So it ends up being not nearly the panacea that it would seem to be.
- bpicolo 4y agoSoftware side sharding keys seem significantly simpler than managing sharding infrastructure. That's the beauty of memcached right?
- varsketiz 4y ago> In a world where things like Dynamo/Cassandra or Spanner/Cockroach exist, manually-sharded DB solutions are pretty much entirely obsolete. Far from it. Plenty of companies that are now running mysql or postgress will shard manually when they will need to scale.
- vrglvrglvrgl 4y ago[dead]
- theredlancer 4y agoAre you going to make the joke or...
- winrid 4y agoIf you're using Mongo, the answer is "only with an Enterprise support contract and Ops Manager" :)
- elsadek 4y agoI haven't heard again about DB sharding technique since prior 2010. This technique was used to separate, in a given table, most involved fields from those less used by using different server for each shard. With the rise of memory database like Redis, sharding was abondonned.
- rvba 4y ago> For Amazon, that means the orders table and the products table containing the products in the orders table need to be physically colocated. Isnt this some wrong simplification? Product number can stay the same, but be a different thing over time. Product 12345 can be a book on first January 2023, while a comic book on first March 2023 . The idea is that product versions can change over time. (What is mostly abused by scammers who farm fake reviews and bait and switch products). So probably they have some product history table, or save the product information in the order. Also sharding in big data environment is when your data does not fit into Excel anymore, so you have 20 Acess files (1GB database limit in Access).
- vlovich123 4y agoIs the hash based approach scalable? A list operation would require contacting every single shard which seems outrageously expensive. It seems like you want some locality but I’m not aware of locality-preserving hash functions (seems like a contradiction but maybe people here have encountered one?)
- ilyt 4y agoIt works because most operations don't need to "list all of everything"
- vlovich123 4y agoAren’t most SQL queries a table scan at some point? I guess you’d shard the index on range and the actual data can be hash sharded, but I don’t know if that buys you much since now the index and data are on separate machines. + you probably need to completely disallow queries on unindexed tables (ie you’re ClickHouse not spanner). That being said there is also interesting work done in the auto-indexing field that might provide a way out of the problem (ie generating an index transparently for the range seeing your hottest queries) but I think you’re still left with the amplification problem of the machines that need to get hit to access the underlying value. Also, I’m not saying “list all”. I’m saying even list 10 or list 1000 is the same problem - you still have to contact all servers in your cluster to do a map/reduce to get the result. Sure, list operations may be less common but their cost seems exponentially more expensive.
- nullandvoid 4y agoDoes anyone else use planet scale but find it extremely slow? I'm using planetscale trial, simple queries can easily take 3-4s "cold", and then seem to get a little faster once presumably I hit cache.
- mscccc 4y agoHey, I work for PlanetScale. Definitely not normal. It’s hard to know why you’re seeing slow queries without more information. The most common causes are either: missing indexes or network latency between the app and database (are they in the same region?). We don’t have cold starts, but it is possible for queries to get faster once data is moved into memory. 3-4s is very slow though, I suspect it’s doing a full table scan and an index will solve it. If you check Insights you can get more info, would also help to run an explain on the query (https://planetscale.com/courses/mysql-for-developers/queries/explain-overview https://planetscale.com/courses/mysql-for-developers/queries...) to see what’s happening. Also, if you email support, they’ll help debug it for you. Hope that helps!