Database sharding vs partitioning

Viewed 137562

I have been reading about scalable architectures recently. In that context, two words that keep on showing up with regards to databases are sharding and partitioning. I looked up descriptions but still ended up confused.

Could the experts at stackoverflow help me get the basics right?

  • What is the difference between sharding and partitioning ?
  • Is it true that 'all sharded databases are essentially partitioned (over different nodes), but all partitioned databases are not necessarily sharded' ?
8 Answers

Consider a Table in database with 1 Million rows and 100 columns In Partitioning you can divide the table into 2 or more table having property like:

  1. 0.4 Million rows(table1), 0.6 million rows(table2)

  2. 1 Million rows & 60 columns(table1) and 1 Million rows & 40 columns(table2)

    There could be multiple cases like that

This is general partitioning

But Sharding refer to 1st case only where we are dividing the data on the basis of rows. If we are dividing the table into multiple table we need to maintain multiple similar copies of schemas as now we have multiple tables.

When talking about partitioning please do not use term replicate or replication. Replication is a different concept and out of scope of this page. When we talk about partitioning then better word is divide and when we talk about sharding then better word is distribute. In partition (normally and in common understanding not always) the rows of large data set table are divided into two or more disjoint (not sharing any row) groups. You can call each group a partition. These groups or all the partitions remain under the control of once RDMB instance and this is all logical. The base of each group can be a hash or range or etc. If you have ten years data in a table then you can store each of the year's data in a separate partition and this can be achieved by setting partition boundaries on the basis of a non-null column CREATE_DATE. Once you query the db then if you specify a create date between 01-01-1999 and 31-12-2000 then only two partitions will be hit and it will be sequential. I did similar on DB for billion + records and sql time came to 50 millis from 30 seconds using indices etc all. Sharding is that you host each partition on a different node/machine. Now searching inside the partitions/shards can happen in parallel.

Sharding in a special case of horizontal partitioning, when partitions spans across multiple database instances. If a database is sharded, it means that it's partitioned by definition.

Horizontal partition when moved to another database instance* becomes a database shard.

Database instance can be on the same machine or on another machine.

Another thing to consider beyond other answers is that you will either partition or shard your database depending on the limitations you want to solve.

For example, you may want to partition because your database doesn't work well with huge tables.

However, you could also face server limitations, where you've done everything you could to optimise your server, but now you have to go for more servers/nodes, and then you'll be sharding.

Related