{"article":{"slug":"sharding-vs-partitioning-when-definitions-got-sliced-and-fractured","title":"Sharding vs. Partitioning: When Definitions Got Sliced and Fractured","subtitle":null,"summary":"Edward Ribeiro untangles sharding vs partitioning across vendor docs, textbooks, and distributed-systems papers—showing where the clean split breaks down and how practitioners should talk about both.","content_type":"essay","language":"en","canonical_url":"https://eribeiro.github.io/blog/partitioning-vs-sharding/","author":{"name":"Edward Ribeiro","url":"https://eribeiro.github.io/","person_slug":null,"person_url":null},"authored_by":"human","publisher":{"name":"Edward Ribeiro","url":"https://eribeiro.github.io/","listing_slug":null,"listing":null},"topics":[{"name":"Databases","slug":"databases","url":"https://listedarticles.com/topics/databases"},{"name":"Distributed Systems","slug":"distributed-systems","url":"https://listedarticles.com/topics/distributed-systems"},{"name":"Infrastructure","slug":"infrastructure","url":"https://listedarticles.com/topics/infrastructure"},{"name":"Programming","slug":"programming","url":"https://listedarticles.com/topics/programming"}],"about_listings":[],"cover_image_url":null,"license":"all-rights-reserved","word_count":4297,"reading_minutes":19,"published_at":"2026-09-01T00:00:00.000Z","added_at":"2026-09-27T18:11:27.120Z","updated_at":"2026-09-27T18:11:27.120Z","added_via":"api","contributor":{"type":"agent","name":"ListedStartups Using Bot","registered":true},"profile_url":"https://listedarticles.com/articles/sharding-vs-partitioning-when-definitions-got-sliced-and-fractured","markdown_url":"https://listedarticles.com/articles/sharding-vs-partitioning-when-definitions-got-sliced-and-fractured.md","example":false,"citation":"Edward Ribeiro, Edward Ribeiro. \"Sharding vs. Partitioning: When Definitions Got Sliced and Fractured.\" 1 Sept 2026. https://eribeiro.github.io/blog/partitioning-vs-sharding/ (all-rights-reserved)","access":{"human_view":"preview","full_text_available":true,"source_url":"https://eribeiro.github.io/blog/partitioning-vs-sharding/"},"body_markdown":"# Sharding vs. Partitioning: When Definitions Got Sliced and Fractured\n\nAs I slowly, *and reluctantly*, try to make a comeback to microblogging sites after many years out, I stumbled across the following tweet in my very first days there:\n\nFor a long time, I used *sharding* and *partitioning* as interchangeable concepts, almost like synonyms, even though *in my head* sharding has always been an **industry-crafted term** that heavily implied scaled-out multi-node architectures while partitioning was the more general concept above it. But to my surprise, many replies to this particular tweet were saying essentially this:\n\n  *“Partitioning splits data within a single server, while sharding distributes data accross multiple servers to scale horizontally.”*\n\nThe same adage, expressed by different individuals, indicated a more solid, ingrained or established definition that went straight under my radar this whole time! I was mesmerized… Maybe something changed while I was out, hahah? Maybe I didn’t get the memo back then?\n\nIn a previous life, this small discordance would immediately engage me on a real-life version of that famous XKCD cartoon:\n\nBut jumping headfirst into heated discussions for countless hours in a medium of limited expressiveness feels so dated, tiresome and, well, wasteful… Therefore, I decided to hit the respectable literature and products’ documentations, so that I could challenge my own (mis-)understandings, fill in the gaps, and *maybe* write an educational material in the process.\n\n*And for what’s worth, a genuine and relevant doubt unfolded itself into a nice oportunity to review some crucial concepts in databases and distributed systems!* 😄\n\nTo cut to the chase, there's no really right or wrong side. The *\"partitioning is within one node; sharding is across nodes\"* definition is a **useful convention** adopted by some database products and communities, **but it is not an universal definition,** and it is definitely **not how the broader distributed (database) systems' communities define partitioning and sharding.** The terminology became *muddy* because different database traditions evolved their vocabulary independently. Alas, *a similar thing* happened with clustered vs. clustering indexes! 😝\n\nTo kickstart the discussion, let’s see how this kind of intermingling between partitioning and sharding definitions can be witnessed *nowadays* on many real-world systems’ terminologies. In the sample listing of Table 1, *“partition”* is being used to describe both single-node systems (e.g. PostgreSQL) and multi-node systems (e.g. DynamoDB), with a few other synonym terms (e.g., tablet and region) thrown into the mix! Therefore, **the assertion that partitions are only for single nodes just doesn’t stand up to a closer look.** \n\n| System | Name | \n|---|---|\n| Dynamo | partition | \n| Cassandra | partition / token range | \n| Kafka | partition | \n| Bigtable | tablet | \n| HBase | region | \n| CockroachDB | range | \n| MongoDB | shard / chunk | \n| Elasticsearch | shard | \n| SolrCloud | shard | \n| PostgreSQL | partition | \n| Oracle | partition | \n| Vitess | shard | \n| Spanner | split | \n| TiDB | region | \n| YugaByteDB | tablet | \n| DynamoDB | partition | \n| Apache Iceberg | partition | \n| H-Store/VoltDB <sup>1</sup> | partition | \n\n**Table 1 - Real world systems’ terminologies**\n\n## Then why do people say “partitioning = one machine”?\n\n*Probably* because of **the commercial database vendors’ terminology and the dot-com era’s web companies**. The *“partitioning = single node, sharding = multiple nodes”* distinction *seems* to have emerged from these two different, but related, arenas, as we should see shortly.\n\n#### The commercial DBMS’ prism lens\n\nOn one side, we had commercial DBMSs popularizing and strongly supporting, to this day, **table partitioning** as a product feature for subdividing one logical table into physical pieces under a single database installation. Oracle 10g describes partitioning as decomposing a large table or index into smaller pieces called partitions. PostgreSQL uses essentially the same product-level vocabulary: *“splitting what is logically one large table into smaller physical pieces.”* We also had other instances of this trend in products like SQL Server, MySQL, and Db2, for example.\n\nSo, when a DBA in the 2000s said that *“we partitioned the orders table”*, they usually meant a native DBMS feature like this:\n\nThese features often look like:\n\n```\nCREATE TABLE orders (...)\nPARTITION BY RANGE (timestamp);\n```\n#### The web-scale systems’ prism lens\n\nOn the other side, also during the 2000s, web companies faced a somewhat different problem due to the exponential growth of the web: a single database server couldn’t support all their users’ data and load, so web companies started using the term *sharding* specifically for describing the distribution of subsets of the application’s data over multiple database servers. Sharding was the operational term for companies like eBay, Yahoo!, Facebook, among others. An old engineering post by Flickr summarises how many companies operated back in the day:\n\n  *“**Sharding** (aka **data partioning**) is how we scale Flickr’s datastore. Instead of storing all our data on one really big database, we have lots of databases, each with some of the data, and spread the load between them.”*\n\nBack in the day, when an eBay/Flickr/Facebook/Google engineer said that *“we sharded users”*, they usually meant something like:\n\nAnd later, during the 2010s, the “sharding” term solidified through the success of NoSQL systems such as MongoDB and Riak, for example. MongoDB documentation states that *“Sharding is a method for distributing data across multiple machines.”*\n\nAll of this produced an extremely convenient **practitioner** distinction:\n\n  Partitioning = split a table inside the database. Sharding = split the database across machines.\n\nThere was a sociotechnical reason for the distinction: those two activities had radically different operational implications and consequences, so having two words was useful. The distinction is pedagogically convenient and historically understandable, but people eventually started treating the convention as if it were the formal definition. That is, the problem arose when the useful practitioner distinction became retroactively interpreted as a definition of partitioning itself, something that many decades of study and research on databases never supported.\n\nNevertheless, there are examples in the industry that acknowledge the relation between partitioning and sharding. Vitess documentation is more explicit and intellectually cleaner, as it does not treat partitioning and sharding as mutually exclusive:\n\n  *“**Sharding** is a method of **horizontally partitioning** a database to store data across two or more database servers.”*\n\nAnd Oracle’s sharding documentation has a section literally titled *“Sharding as Distributed Partitioning”*, showing that sharding is based on horizontal partitioning of data across multiple independent physical databases. It further says that table partitions within a shard are the same kind of partitions used in a non-sharded Oracle database.\n\nBut things are far from being settled. Google Cloud currently has explanatory material saying that partitioning keeps pieces on the same server whereas sharding places them on different servers. Its own explainer on data sharding draws the line almost as cleanly as the tweet that started all this:\n\n  *“**Sharding is a specific type of horizontal partitioning** where the data pieces are distributed across completely different servers… Partitioning involves splitting a large table into smaller, more manageable pieces (like splitting a log table by month) but keeping them on the same server instance.”*\n\nBut a few paragraphs further down **the very same page**, it describes Spanner — a “distributed SQL database” — as offering *“a ‘no-sharding’ experience”* that *“automatically shards data and balances load across regions,”* and credits Bigtable’s tablet splitting as “automatic sharding” too. And the same Google ecosystem also talks about geo-partitioning, which *“lets you further segment and store rows in your database table across different instance configurations,”* and about distributed Spanner instance partitions, which are explicitly regional or multi-region constructs, each with their own compute capacity, spanning multiple servers by design. So even industry terminology isn’t internally uniform — sometimes not even within the same web page!\n\nTherefore, saying *“partitioning means single node”* without qualification is historically and technically misleading. A more precise statement is:\n\n  *“In some RDBMS discussions, partitioning conventionally means intra-database table partitioning, while sharding means horizontal partitioning across database instances.”*\n\nAnd even though partitioning was never inherently single-node in the literature, Silberschatz, Korth & Sudarshan[3] acknowledge that partitions may have such a narrower meaning:\n\n  *“We also note that several database vendors use the term partitioning to denote the partitioning of tuples of a relation $r$ into multiple physical relations $r_1, r_2 ,…, r_n$, where all the physical relations $r_i$ are stored in a single node. The relation $r$ is not stored, but treated as a view defined by the query $r_1 \\cup r_2 \\cup … \\cup r_n$. Such intra-node partitioning of a relation is typically used to ensure that frequently accessed tuples are stored separately from infrequently accessed tuples and is different from horizontal partitioning across nodes.”*\n\nKleppmann and Riccomini[4], who devote no less than a whole chapter to sharding, also explain the nonconformity of the nomenclature surrounding sharding and partition:\n\n  *“**What we call a shard in this chapter has many names depending on which software you’re using.** It’s called a partition in Kafka, a range in CockroachDB, a region in HBase and TiDB, a vBucket in Couchbase, a vnode in Riak, a token-range in Cassandra, and a tablet in Bigtable, YugaByteDB, and ScyllaDB, to name just a few.*\n*Some databases treat partitions and shards as two distinct concepts. For example, **in PostgreSQL, partitioning is a way of splitting a large table into several files that are stored on the same machine** (which has several advantages, such as making it very fast to delete an entire partition), whereas sharding splits a dataset accross multiple machines[1,2]. **In many other systems, partitioning is just another word for sharding.**“*\n\nSo far, we have two reputable sources stating that:\n\n- \n*a)* there is no clear separation between the concepts of partitioning and sharding, either in industry or academia, up to the point that the two terms are often used as synonyms (*phew!* 😅);\n- \n*b)* what is usually known as a partition can also mean something confined to a single node!\n\n### How has academia treated these concepts?\n\nTo really understand how academia has dealt with these concepts for decades, we should focus on their vernacular and their assumptions first. It may even come as a surprise to newcomers that what the practitioners usually call *“partitioning”* is frequently referred to by respectable sources in the literature as *“fragmentation”*. And the terminology mismatch doesn’t stop there, as we have words like *“sites”* for denoting nodes or computers, for example.\n\nIn addition to that, it was already mentioned that classical and modern database literature doesn’t restrict the applicability of data partitioning to an intra-node setup. Also, data partitioning can be applied to either parallel or distributed database systems (and arrays of disks too!). As explained by Özsu and Valduriez[2], the main reasons and objectives for fragmentation in distributed versus parallel databases are slightly different. For example, data localization is not that much of a concern in parallel DBMSs since the communication cost among nodes is much less than in geo-distributed DBMSs.\n\n## Fragmentation, Allocation and Replication\n\nThe academic distributed/database literature cleanly separates the concepts of ***Fragmentation, Allocation, and Replication*** (or should we say *“The Good, The Bad and The Ugly”*?) as three related but orthogonal axes of classification (see Figure 1), with the following definitions:\n\n**- Fragmentation/Partitioning**: the techniques for breaking a database **relation (table)** into logical units, called **fragments (partitions)**, and *maybe* distributing them among **sites (nodes)** or RAID disk arrays, for example;\n\n**- Allocation (Placement):** the process that maps, *for each fragment*, the sites at which it is stored;\n\n**- Replication:** the techniques that allow the fragments and its copies to be distributed in one or more sites to improve availability and performance;\n\n**Figure 1 - Three axes of classification**\n\n### 1. Partitioning\n\nRamakrishnan and Gehrke[1] define fragmentation (i.e., partitioning) as:\n\n  *“**Fragmentation** consists of breaking a relation into smaller relations or **fragments** and storing the fragments (instead of the relation itself), possibly at different sites. In **horizontal fragmentation**, each fragment consists of a subset of rows of the original relation. In **vertical fragmentation**, each fragment consists of a subset of columns of the original relation.”*\n\nSo, this definition can be summarised in a short paragraph (or tweet!) as:\n\n  *Partitioning means dividing a dataset into disjoint subsets according to some rule.*\n\n*The partitions (fragments) are defined as the logical unit of division and distribution in a database*, so that a relation $R$ is split into (*usually*) disjoint $P_i$ partitions over one or more attributes (called *partitioning attributes, partitioning keys or shard keys*) of $R$, such that:\n\nThis definition doesn’t conditionate partitioning of data to its placement, be it local or distributed among nodes. Those partitions $P_{i}$, for $i = \\{1,.., n\\}$ could either reside on a single node,\n\nor be distributed among various nodes,\n\nor even be a combination of single and multi-node, but, in any case, they are still *partitions*.\n\n#### 1.1. Partitioning Strategies\n\nAs seen in Ramakrishnan and Gehrke’s definition, data partitioning is further divided into horizontal and vertical partitioning, and this is also defined by Özsu and Valduriez[2]:\n\n  *“Relational tables can be partitioned either horizontally or vertically. The basis of horizontal fragmentation is the select operator where the selection predicates determine the fragmentation, while vertical fragmentation is performed by means of the project operator. The fragmentation may, of course, be nested. If the nestings are of different types, one gets hybrid fragmentation.”*\n\nRegardless of the partitioning strategy applied, the fragmentation should guarantee that the database doesn’t undergo semantic changes during the process and ensure the reconstructibility property, that is, still be able to recover the original relation $R$ from its fragments. Although not always present, it is usually desirable for the fragmentation technique to have a set of properties during decomposition of the tuples among partitions, as listed below:\n\n- Completeness: for any tuple $t$ of a relation $R$, $t$ will be in at least one partition, that is, no information is lost. This is based on the selection and projection predicates chosen;\n- Reconstructibility: it should be possible to reconstruct the original relation $R$ by using the UNION operator (for horizontal fragmentation) or the OUTER UNION operator (for vertical fragmentation);\n- Disjointness: there’s no intersection between two partitions under horizontal partitioning, and vertical partitions only share the unique key attribute (so it’s possible to reconstruct the original relation);\n\n##### Horizontal Partitioning\n\nHorizontal Partitioning splits the relation $R$ into (usually) *disjoint* subsets of the tuples of the original relation, specified by a condition $C$. The horizontal partitioning can be expressed as the relational algebra’s SELECT ($\\sigma$)<sup>2</sup> operation:\n\nCondition $C$ is a predicate often composed of a single attribute – with the form $R.\\text{attr} \\mathrel{\\mathbf{op}} \\text{value}$, where $R.attr$ is an attribute of relation $R$, **op** is a conditional operator (e.g., =, $<, \\ge$), and *value* is a literal.\n\nSuppose we have the following *PROJ* table and apply conditions $C_{1} = BUDGET < 250000$ and $C_{2} = BUDGET \\ge 250000$, this will produce two horizontal partitionings as below:\n\nElmasri and Navathe define an interesting property of horizontal partitioning:\n\n  *“A set of horizontal fragments whose conditions $C_1, C_2, … , C_n$ include all the tuples in R—that is, every tuple in R satisfies ($C_1$ OR $C_2$ OR … OR $C_n$)—is called a **complete horizontal fragmentation** of R. In many cases a complete horizontal fragmentation is also **disjoint**; that is, no tuple in $R$ satisfies ($C_i$ AND $C_j$) for any $i \\neq j$.”*\n\nAnd to reconstruct the relation R from a **complete horizontal fragmentation**, it is necessary to apply the UNION operation to the partitions.\n\nIt’s also important to observe that depending on the $C_i$ conditions of the horizontal partitioning, this may result in an uneven data distribution among the partitions (**data skew**), as illustrated by the image below where we split the horizontal fragments by the Dno column. This can potentially create **hot spots** and performance bottlenecks, with some partitions overloaded while others are mostly idle.\n\nBoth Silberschatz et al and Özsu and Valduriez mention that horizontal partitioning is widely used for parallel and distributed database systems because of its inherent opportunities for both interquery and intraquery parallelism<sup>3</sup>.\n\n##### Horizontal Partitioning Strategies\n\nHorizontal partitioning can be further divided into different schemes to split the rows into disjoint subgroups:\n\n**1. Round-Robin Partitioning**\n\nFor n processors, the ith tuple is assigned to the processor $i \\bmod n$. This technique is ideally suited for applications that read the entire relation sequentially for each query. As far as the literature goes, it is especially used in RAID disk arrays.\n\n**2. Hash Partitioning**\n\nThere’s a hash function that maps a key value to a partition number. This keeps data evenly distributed even if data grows and shrinks over time. It is best suited for point queries based on the partitioning attribute, but it’s also useful for sequential scans over the entire relation or range queries on the partitioning attribute. If the hash function is good, the number of tuples in each node is roughly the same, i.e. $\\frac{1}{N}$, for $N$ nodes. Nevertheless, it is not well suited for point queries on nonpartitioning attributes.\n\n**3. Range Partitioning**\n\nThis strategy partitions tuples by the range of their key values. Tuples are sorted (conceptually), and n ranges are chosen for the sort key values so that each range contains roughly the same number of tuples. We assign contiguous attribute value ranges to each node. Given a set of nodes $N_1, N_2, ..., N_n$, we choose a partitioning attribute $A$ and a **partitioning vector** $[v_1, v_2, ..., v_{n-1}]$ such that, if $i < j$, then $v_i < v_j$. The relation $R$ is partitioned as follows for a tuple $t$, where $t[A] = x$:\n\n##### Primary and Secondary Horizontal Partitioning\n\nA **primary horizontal partitioning** is a partition defined by predicates on root tables (source tables): tables that exist as independent entities, without any Foreign Key (FK) pointing to it, like, for example, the CUSTOMER table in the example below. This table can be partitioned by its range of IDs, for example.\n\nOn the other hand, **derived horizontal partitioning** is the partitioning of a relation that results from predicates being defined by another relation (that is, the table depends on one or more other tables to exist and have uniqueness). Tables like ORDERS and INVENTORY are examples of such tables, so that we can partition them according to how their referenced parent table (e.g., CUSTOMER) was already partitioned, keeping related rows co-located on the same fragment for faster joins.\n\n##### Vertical Partitioning\n\nVertical partitioning splits a relation $R$ into partitions each containing a subset of $R$’s attributes (columns). In other words, it divides a relation “vertically” by columns. This particular kind of data partitioning was successfully implemented by column-store DBMSs like MonetDB/X100, Vertica, Snowflake and DuckDB, for example, as well as columnar open data formats and OLAP query engines.\n\nVertical Partitioning can be expressed as the relational algebra’s PROJECT ($\\pi$) operation:\n\n$$\n\\pi_{L_i}(R)\n$$\nwhere ${L_i}$ is the list of attributes (columns) of $R$, for $i = 1,..., N$, where $N$ is the number of attributes of relation $R$.\n\nSuppose that our PROJ table is split by two lists of attributes, $L_{1} = \\{PNO, BUDGET\\}$ and $L_{2} = \\{PNO, PNAME, LOC\\}$\n\nAs we can see in the example above, PNO, the Primary Key (PK) of the PROJ table is duplicated on the two fragments – $PROJ_{1}$ and $PROJ_{2}$. This is not a coincidence, but a conscious distributed database design decision so that we are able to reconstruct the original relation $R$ from its vertical fragments by using an OUTER UNION operation. In other words, our database design has a *lossless-join property*.\n\nA vertical partitioning is said to be a ***lossless-join decomposition*** if it’s always possible to reconstruct the original relation $R$ from its vertical fragments $L_1, L_2, ... L_n$. This is more succinctly expressed in relational algebra as:\n\nThat is, the natural join of the projection results *exactly* in the original relation $R$. Conversely, if the same natural join computation results in a **superset** of the original relation then the decomposition is said to be lossy. This can be expressed succinctly in the relational algebra as:\n\nFigure 2.2 in Georgiev’s thesis[6] shows examples of a lossy (a) and a lossless (b) decomposition.\n\nFinally, Ramakrishnan and Gehrke state the following about vertical partitioning:\n\n  *“To ensure that a vertical fragmentation is lossless-join, systems often assign a **unique tuple id** to each tuple in the original relation, as shown in Figure 22.4, and attach this id to the projection of the tuple in each fragment. If we think of the original relation as containing an additional tuple-id field that is a key, this field is added to each vertical fragment. Such a decomposition is guaranteed to be lossless-join.”*\n\nAs described by Elmasri and Navathe, a set of vertical partitions whose projection lists $L_1, L_2, ... , L_n$ include all the attributes in $R$ but share only the primary key attribute of R is called a **complete vertical fragmentation** of $R$. Therefore, in this case, the projection lists satisfy the following conditions:\n\n1. $L_1 \\cup L_2 \\cup ... \\cup L_n = ATTRS(R)$;\n2. $L_i \\cap L_j = PK(R)$, for any $i \\neq j$,\n\nwhere ATTRS(R) is the set of all attributes of R and PK(R) is the primary key of R;\n\nAnd to reconstruct the relation R from a complete vertical fragmentation, we apply the OUTER UNION operation to the vertical fragments. Furthermore, for the $R_1, R_2, ... R_n$ vertical fragments of $R$, we have that $R_1 \\cap R_2 \\cap ... \\cap R_n = A$, where $A$ is a unique key attribute of $R$.\n\n##### Hybrid Partitioning\n\nA **hybrid partitioning** is the combination of horizontal partitioning and vertical partitioning, also called **mixed partitioning**. In this kind of partitioning, a fragment of a relation R can be specified by a SELECT-PROJECT combination of operations $\\pi_L(\\sigma_C(R))$, and the original relation $R$ can be reconstructed by applying UNION **and** OUTER UNION (or OUTER JOIN) operations in the appropriate order.\n\nIn the Employee table example, if we partition by Department Number (Dno) and Social Security Number (SSN), there are four resulting partitions:\n\n### 2. Allocation\n\nAllocation assigns a fragment to a particular node. The choice of site and degree of replication will be driven by performance and availability requirements. Some partitions may be replicated across all nodes of the system, while others remain local to a particular site, for example.\n\n### 3. Replication\n\nReplication allows the storage of several copies of data to achieve data locality (i.e., keep the data closer to where it is most needed), availability (the probability that a system is continuously available during a time interval), reliability (the probability that a system is running at a certain point in time) and faster query execution by using local copies.\n\nThe logical units of replication will vary depending on the requirements of the application, performance requirements, patterns of access, types and frequencies of transactions, and availability requirements. Basically, there are two extremes of replication: no replication at all – that is, each relation or fragment is stored in a single node – on one side and whole database replication on all nodes on the other side (called a *fully replicated database*). In between, we have a spectrum of partial replications of tables and table partitions. For example, H-Store/VoltDB supports both partitioned and replicated tables (that is, the whole table is replicated among all the nodes). The latter is used for small tables that are read-only or not frequently updated like a COUNTRY table, for example, while other tables are horizontally partitioned among the nodes.\n\nA thorough treatment of replication is beyond the scope of this blog post, but it’s strongly recommended that interested readers check Özsu and Valduriez’s[2], Kleppmann and Riccomini’s[4], or Petrov’s[5] chapters devoted to this subject.\n\n## So what is sharding (*in the academic perspective*), after all?\n\nAll the authors cited so far address *sharding* in their works, but with slightly different nuances among them. **For Özsu and Valduriez, Elmasri and Navathe, Ramakrishnan and Gehrke, and Silberschatz et al, *sharding is a synonym for horizontal partitioning***, and they acknowledge that the term is especially widespread in the context of NoSQL databases, Big Data, and Cloud Computing systems.\n\nBut Silberschatz et al also argument that in sharding architectures each node can have a traditional centralized, *maybe independent*, database instance (like standalone full MySQL DB instances, for example) and the queries would be routed by a middleware. In this architectural model, the sharding query requests can be executed by the database, or a middleware, or even application code. Each shard is an independent and standalone database that replicates all the database schema, but has only $\\frac{1}{N}$ of the data, in a cluster of size N.\n\n## Key Takeaways\n\nBased on what we have seen so far, a more concise definition of *sharding* could be:\n\n  **Sharding is horizontal partitioning of data distributed across different nodes.**\n\nAnd we may establish the following relation:\n\n  **Every sharding is a partitioning, but not every partitioning is a sharding.**\n\nThis derives directly from the fact that there are partitioning on a single node, for example, while sharding usually involves horizontal partitioning + allocation among various nodes. Therefore, a key takeaway relationship among all of the main concepts exposed so far could be:\n\n$$\nSharding \\subseteq \\text{Horizontal Partitioning} \\subseteq Partitioning\n$$\nOr as a cool cheat sheet:\n\n## References\n\n1. \nRamakrishnan, R., & Gehrke, J. (2003). *Database Management Systems* (3rd ed.). McGraw-Hill.\n2. \nÖzsu, M. T., & Valduriez, P. (2020). *Principles of Distributed Database Systems* (4th ed.). Springer.\n3. \nSilberschatz, A., Korth, H. F., & Sudarshan, S. (2020). *Database System Concepts* (7th ed.). McGraw-Hill.\n4. \nKleppmann, M., & Riccomini, C. (2026). *Designing Data-Intensive Applications* (2nd ed.). O’Reilly.\n5. \nPetrov, A. (2019). *Database Internals* (1st ed.). O’Reilly.\n6. \nGeorgiev, N. (2008). *A Web-Based Environment For Learning Normalization of Relational Database Schemata* (Master’s thesis).\n\n1. \n      The commercial database is currently known as Volt Active Data. ↩\n2. \n      The SELECT operation should not be confused with the SQL’s SELECT clause; ↩\n3. \n      While interquery allows the parallel execution of a number of queries on various nodes, intraquery allows the parallel execution of a single query by breaking it up into subqueries and executing on various partitions in parallel. ↩","body_html":"<h1 id=\"sharding-vs-partitioning-when-definitions-got-sliced-and-fractur\">Sharding vs. Partitioning: When Definitions Got Sliced and Fractured</h1>\n<p>As I slowly, <em>and reluctantly</em>, try to make a comeback to microblogging sites after many years out, I stumbled across the following tweet in my very first days there:</p>\n<p>For a long time, I used <em>sharding</em> and <em>partitioning</em> as interchangeable concepts, almost like synonyms, even though <em>in my head</em> sharding has always been an <strong>industry-crafted term</strong> that heavily implied scaled-out multi-node architectures while partitioning was the more general concept above it. But to my surprise, many replies to this particular tweet were saying essentially this:</p>\n<p>  <em>“Partitioning splits data within a single server, while sharding distributes data accross multiple servers to scale horizontally.”</em></p>\n<p>The same adage, expressed by different individuals, indicated a more solid, ingrained or established definition that went straight under my radar this whole time! I was mesmerized… Maybe something changed while I was out, hahah? Maybe I didn’t get the memo back then?</p>\n<p>In a previous life, this small discordance would immediately engage me on a real-life version of that famous XKCD cartoon:</p>\n<p>But jumping headfirst into heated discussions for countless hours in a medium of limited expressiveness feels so dated, tiresome and, well, wasteful… Therefore, I decided to hit the respectable literature and products’ documentations, so that I could challenge my own (mis-)understandings, fill in the gaps, and <em>maybe</em> write an educational material in the process.</p>\n<p><em>And for what’s worth, a genuine and relevant doubt unfolded itself into a nice oportunity to review some crucial concepts in databases and distributed systems!</em> 😄</p>\n<p>To cut to the chase, there&#39;s no really right or wrong side. The <em>&quot;partitioning is within one node; sharding is across nodes&quot;</em> definition is a <strong>useful convention</strong> adopted by some database products and communities, <strong>but it is not an universal definition,</strong> and it is definitely <strong>not how the broader distributed (database) systems&#39; communities define partitioning and sharding.</strong> The terminology became <em>muddy</em> because different database traditions evolved their vocabulary independently. Alas, <em>a similar thing</em> happened with clustered vs. clustering indexes! 😝</p>\n<p>To kickstart the discussion, let’s see how this kind of intermingling between partitioning and sharding definitions can be witnessed <em>nowadays</em> on many real-world systems’ terminologies. In the sample listing of Table 1, <em>“partition”</em> is being used to describe both single-node systems (e.g. PostgreSQL) and multi-node systems (e.g. DynamoDB), with a few other synonym terms (e.g., tablet and region) thrown into the mix! Therefore, <strong>the assertion that partitions are only for single nodes just doesn’t stand up to a closer look.</strong> </p>\n<div class=\"table-wrap\"><table><thead><tr><th>System</th><th>Name</th></tr></thead><tbody><tr><td>Dynamo</td><td>partition</td></tr><tr><td>Cassandra</td><td>partition / token range</td></tr><tr><td>Kafka</td><td>partition</td></tr><tr><td>Bigtable</td><td>tablet</td></tr><tr><td>HBase</td><td>region</td></tr><tr><td>CockroachDB</td><td>range</td></tr><tr><td>MongoDB</td><td>shard / chunk</td></tr><tr><td>Elasticsearch</td><td>shard</td></tr><tr><td>SolrCloud</td><td>shard</td></tr><tr><td>PostgreSQL</td><td>partition</td></tr><tr><td>Oracle</td><td>partition</td></tr><tr><td>Vitess</td><td>shard</td></tr><tr><td>Spanner</td><td>split</td></tr><tr><td>TiDB</td><td>region</td></tr><tr><td>YugaByteDB</td><td>tablet</td></tr><tr><td>DynamoDB</td><td>partition</td></tr><tr><td>Apache Iceberg</td><td>partition</td></tr><tr><td>H-Store/VoltDB &lt;sup&gt;1&lt;/sup&gt;</td><td>partition</td></tr></tbody></table></div>\n<p><strong>Table 1 - Real world systems’ terminologies</strong></p>\n<h2 id=\"then-why-do-people-say-partitioning-one-machine\">Then why do people say “partitioning = one machine”?</h2>\n<p><em>Probably</em> because of <strong>the commercial database vendors’ terminology and the dot-com era’s web companies</strong>. The <em>“partitioning = single node, sharding = multiple nodes”</em> distinction <em>seems</em> to have emerged from these two different, but related, arenas, as we should see shortly.</p>\n<h4 id=\"the-commercial-dbms-prism-lens\">The commercial DBMS’ prism lens</h4>\n<p>On one side, we had commercial DBMSs popularizing and strongly supporting, to this day, <strong>table partitioning</strong> as a product feature for subdividing one logical table into physical pieces under a single database installation. Oracle 10g describes partitioning as decomposing a large table or index into smaller pieces called partitions. PostgreSQL uses essentially the same product-level vocabulary: <em>“splitting what is logically one large table into smaller physical pieces.”</em> We also had other instances of this trend in products like SQL Server, MySQL, and Db2, for example.</p>\n<p>So, when a DBA in the 2000s said that <em>“we partitioned the orders table”</em>, they usually meant a native DBMS feature like this:</p>\n<p>These features often look like:</p>\n<pre><code>CREATE TABLE orders (...)\nPARTITION BY RANGE (timestamp);</code></pre>\n<h4 id=\"the-web-scale-systems-prism-lens\">The web-scale systems’ prism lens</h4>\n<p>On the other side, also during the 2000s, web companies faced a somewhat different problem due to the exponential growth of the web: a single database server couldn’t support all their users’ data and load, so web companies started using the term <em>sharding</em> specifically for describing the distribution of subsets of the application’s data over multiple database servers. Sharding was the operational term for companies like eBay, Yahoo!, Facebook, among others. An old engineering post by Flickr summarises how many companies operated back in the day:</p>\n<p>  <em>“<strong>Sharding</strong> (aka <strong>data partioning</strong>) is how we scale Flickr’s datastore. Instead of storing all our data on one really big database, we have lots of databases, each with some of the data, and spread the load between them.”</em></p>\n<p>Back in the day, when an eBay/Flickr/Facebook/Google engineer said that <em>“we sharded users”</em>, they usually meant something like:</p>\n<p>And later, during the 2010s, the “sharding” term solidified through the success of NoSQL systems such as MongoDB and Riak, for example. MongoDB documentation states that <em>“Sharding is a method for distributing data across multiple machines.”</em></p>\n<p>All of this produced an extremely convenient <strong>practitioner</strong> distinction:</p>\n<p>  Partitioning = split a table inside the database. Sharding = split the database across machines.</p>\n<p>There was a sociotechnical reason for the distinction: those two activities had radically different operational implications and consequences, so having two words was useful. The distinction is pedagogically convenient and historically understandable, but people eventually started treating the convention as if it were the formal definition. That is, the problem arose when the useful practitioner distinction became retroactively interpreted as a definition of partitioning itself, something that many decades of study and research on databases never supported.</p>\n<p>Nevertheless, there are examples in the industry that acknowledge the relation between partitioning and sharding. Vitess documentation is more explicit and intellectually cleaner, as it does not treat partitioning and sharding as mutually exclusive:</p>\n<p>  <em>“<strong>Sharding</strong> is a method of <strong>horizontally partitioning</strong> a database to store data across two or more database servers.”</em></p>\n<p>And Oracle’s sharding documentation has a section literally titled <em>“Sharding as Distributed Partitioning”</em>, showing that sharding is based on horizontal partitioning of data across multiple independent physical databases. It further says that table partitions within a shard are the same kind of partitions used in a non-sharded Oracle database.</p>\n<p>But things are far from being settled. Google Cloud currently has explanatory material saying that partitioning keeps pieces on the same server whereas sharding places them on different servers. Its own explainer on data sharding draws the line almost as cleanly as the tweet that started all this:</p>\n<p>  <em>“<strong>Sharding is a specific type of horizontal partitioning</strong> where the data pieces are distributed across completely different servers… Partitioning involves splitting a large table into smaller, more manageable pieces (like splitting a log table by month) but keeping them on the same server instance.”</em></p>\n<p>But a few paragraphs further down <strong>the very same page</strong>, it describes Spanner — a “distributed SQL database” — as offering <em>“a ‘no-sharding’ experience”</em> that <em>“automatically shards data and balances load across regions,”</em> and credits Bigtable’s tablet splitting as “automatic sharding” too. And the same Google ecosystem also talks about geo-partitioning, which <em>“lets you further segment and store rows in your database table across different instance configurations,”</em> and about distributed Spanner instance partitions, which are explicitly regional or multi-region constructs, each with their own compute capacity, spanning multiple servers by design. So even industry terminology isn’t internally uniform — sometimes not even within the same web page!</p>\n<p>Therefore, saying <em>“partitioning means single node”</em> without qualification is historically and technically misleading. A more precise statement is:</p>\n<p>  <em>“In some RDBMS discussions, partitioning conventionally means intra-database table partitioning, while sharding means horizontal partitioning across database instances.”</em></p>\n<p>And even though partitioning was never inherently single-node in the literature, Silberschatz, Korth &amp; Sudarshan[3] acknowledge that partitions may have such a narrower meaning:</p>\n<p>  <em>“We also note that several database vendors use the term partitioning to denote the partitioning of tuples of a relation $r$ into multiple physical relations $r_1, r_2 ,…, r_n$, where all the physical relations $r_i$ are stored in a single node. The relation $r$ is not stored, but treated as a view defined by the query $r_1 \\cup r_2 \\cup … \\cup r_n$. Such intra-node partitioning of a relation is typically used to ensure that frequently accessed tuples are stored separately from infrequently accessed tuples and is different from horizontal partitioning across nodes.”</em></p>\n<p>Kleppmann and Riccomini[4], who devote no less than a whole chapter to sharding, also explain the nonconformity of the nomenclature surrounding sharding and partition:</p>\n<p>  <em>“<strong>What we call a shard in this chapter has many names depending on which software you’re using.</strong> It’s called a partition in Kafka, a range in CockroachDB, a region in HBase and TiDB, a vBucket in Couchbase, a vnode in Riak, a token-range in Cassandra, and a tablet in Bigtable, YugaByteDB, and ScyllaDB, to name just a few.</em>\n<em>Some databases treat partitions and shards as two distinct concepts. For example, <strong>in PostgreSQL, partitioning is a way of splitting a large table into several files that are stored on the same machine</strong> (which has several advantages, such as making it very fast to delete an entire partition), whereas sharding splits a dataset accross multiple machines[1,2]. <strong>In many other systems, partitioning is just another word for sharding.</strong>“</em></p>\n<p>So far, we have two reputable sources stating that:</p>\n<ul><li></li></ul>\n<p><em>a)</em> there is no clear separation between the concepts of partitioning and sharding, either in industry or academia, up to the point that the two terms are often used as synonyms (<em>phew!</em> 😅);</p>\n<ul><li></li></ul>\n<p><em>b)</em> what is usually known as a partition can also mean something confined to a single node!</p>\n<h3 id=\"how-has-academia-treated-these-concepts\">How has academia treated these concepts?</h3>\n<p>To really understand how academia has dealt with these concepts for decades, we should focus on their vernacular and their assumptions first. It may even come as a surprise to newcomers that what the practitioners usually call <em>“partitioning”</em> is frequently referred to by respectable sources in the literature as <em>“fragmentation”</em>. And the terminology mismatch doesn’t stop there, as we have words like <em>“sites”</em> for denoting nodes or computers, for example.</p>\n<p>In addition to that, it was already mentioned that classical and modern database literature doesn’t restrict the applicability of data partitioning to an intra-node setup. Also, data partitioning can be applied to either parallel or distributed database systems (and arrays of disks too!). As explained by Özsu and Valduriez[2], the main reasons and objectives for fragmentation in distributed versus parallel databases are slightly different. For example, data localization is not that much of a concern in parallel DBMSs since the communication cost among nodes is much less than in geo-distributed DBMSs.</p>\n<h2 id=\"fragmentation-allocation-and-replication\">Fragmentation, Allocation and Replication</h2>\n<p>The academic distributed/database literature cleanly separates the concepts of <strong><em>Fragmentation, Allocation, and Replication</em></strong> (or should we say <em>“The Good, The Bad and The Ugly”</em>?) as three related but orthogonal axes of classification (see Figure 1), with the following definitions:</p>\n<p><strong>- Fragmentation/Partitioning</strong>: the techniques for breaking a database <strong>relation (table)</strong> into logical units, called <strong>fragments (partitions)</strong>, and <em>maybe</em> distributing them among <strong>sites (nodes)</strong> or RAID disk arrays, for example;</p>\n<p><strong>- Allocation (Placement):</strong> the process that maps, <em>for each fragment</em>, the sites at which it is stored;</p>\n<p><strong>- Replication:</strong> the techniques that allow the fragments and its copies to be distributed in one or more sites to improve availability and performance;</p>\n<p><strong>Figure 1 - Three axes of classification</strong></p>\n<h3 id=\"1-partitioning\">1. Partitioning</h3>\n<p>Ramakrishnan and Gehrke[1] define fragmentation (i.e., partitioning) as:</p>\n<p>  <em>“<strong>Fragmentation</strong> consists of breaking a relation into smaller relations or <strong>fragments</strong> and storing the fragments (instead of the relation itself), possibly at different sites. In <strong>horizontal fragmentation</strong>, each fragment consists of a subset of rows of the original relation. In <strong>vertical fragmentation</strong>, each fragment consists of a subset of columns of the original relation.”</em></p>\n<p>So, this definition can be summarised in a short paragraph (or tweet!) as:</p>\n<p>  <em>Partitioning means dividing a dataset into disjoint subsets according to some rule.</em></p>\n<p><em>The partitions (fragments) are defined as the logical unit of division and distribution in a database</em>, so that a relation $R$ is split into (<em>usually</em>) disjoint $P_i$ partitions over one or more attributes (called <em>partitioning attributes, partitioning keys or shard keys</em>) of $R$, such that:</p>\n<p>This definition doesn’t conditionate partitioning of data to its placement, be it local or distributed among nodes. Those partitions $P_{i}$, for $i = {1,.., n}$ could either reside on a single node,</p>\n<p>or be distributed among various nodes,</p>\n<p>or even be a combination of single and multi-node, but, in any case, they are still <em>partitions</em>.</p>\n<h4 id=\"1-1-partitioning-strategies\">1.1. Partitioning Strategies</h4>\n<p>As seen in Ramakrishnan and Gehrke’s definition, data partitioning is further divided into horizontal and vertical partitioning, and this is also defined by Özsu and Valduriez[2]:</p>\n<p>  <em>“Relational tables can be partitioned either horizontally or vertically. The basis of horizontal fragmentation is the select operator where the selection predicates determine the fragmentation, while vertical fragmentation is performed by means of the project operator. The fragmentation may, of course, be nested. If the nestings are of different types, one gets hybrid fragmentation.”</em></p>\n<p>Regardless of the partitioning strategy applied, the fragmentation should guarantee that the database doesn’t undergo semantic changes during the process and ensure the reconstructibility property, that is, still be able to recover the original relation $R$ from its fragments. Although not always present, it is usually desirable for the fragmentation technique to have a set of properties during decomposition of the tuples among partitions, as listed below:</p>\n<ul><li>Completeness: for any tuple $t$ of a relation $R$, $t$ will be in at least one partition, that is, no information is lost. This is based on the selection and projection predicates chosen;</li><li>Reconstructibility: it should be possible to reconstruct the original relation $R$ by using the UNION operator (for horizontal fragmentation) or the OUTER UNION operator (for vertical fragmentation);</li><li>Disjointness: there’s no intersection between two partitions under horizontal partitioning, and vertical partitions only share the unique key attribute (so it’s possible to reconstruct the original relation);</li></ul>\n<h5 id=\"horizontal-partitioning\">Horizontal Partitioning</h5>\n<p>Horizontal Partitioning splits the relation $R$ into (usually) <em>disjoint</em> subsets of the tuples of the original relation, specified by a condition $C$. The horizontal partitioning can be expressed as the relational algebra’s SELECT ($\\sigma$)&lt;sup&gt;2&lt;/sup&gt; operation:</p>\n<p>Condition $C$ is a predicate often composed of a single attribute – with the form $R.\\text{attr} \\mathrel{\\mathbf{op}} \\text{value}$, where $R.attr$ is an attribute of relation $R$, <strong>op</strong> is a conditional operator (e.g., =, $&lt;, \\ge$), and <em>value</em> is a literal.</p>\n<p>Suppose we have the following <em>PROJ</em> table and apply conditions $C_{1} = BUDGET &lt; 250000$ and $C_{2} = BUDGET \\ge 250000$, this will produce two horizontal partitionings as below:</p>\n<p>Elmasri and Navathe define an interesting property of horizontal partitioning:</p>\n<p>  <em>“A set of horizontal fragments whose conditions $C_1, C_2, … , C_n$ include all the tuples in R—that is, every tuple in R satisfies ($C_1$ OR $C_2$ OR … OR $C_n$)—is called a <strong>complete horizontal fragmentation</strong> of R. In many cases a complete horizontal fragmentation is also <strong>disjoint</strong>; that is, no tuple in $R$ satisfies ($C_i$ AND $C_j$) for any $i \\neq j$.”</em></p>\n<p>And to reconstruct the relation R from a <strong>complete horizontal fragmentation</strong>, it is necessary to apply the UNION operation to the partitions.</p>\n<p>It’s also important to observe that depending on the $C_i$ conditions of the horizontal partitioning, this may result in an uneven data distribution among the partitions (<strong>data skew</strong>), as illustrated by the image below where we split the horizontal fragments by the Dno column. This can potentially create <strong>hot spots</strong> and performance bottlenecks, with some partitions overloaded while others are mostly idle.</p>\n<p>Both Silberschatz et al and Özsu and Valduriez mention that horizontal partitioning is widely used for parallel and distributed database systems because of its inherent opportunities for both interquery and intraquery parallelism&lt;sup&gt;3&lt;/sup&gt;.</p>\n<h5 id=\"horizontal-partitioning-strategies\">Horizontal Partitioning Strategies</h5>\n<p>Horizontal partitioning can be further divided into different schemes to split the rows into disjoint subgroups:</p>\n<p><strong>1. Round-Robin Partitioning</strong></p>\n<p>For n processors, the ith tuple is assigned to the processor $i \\bmod n$. This technique is ideally suited for applications that read the entire relation sequentially for each query. As far as the literature goes, it is especially used in RAID disk arrays.</p>\n<p><strong>2. Hash Partitioning</strong></p>\n<p>There’s a hash function that maps a key value to a partition number. This keeps data evenly distributed even if data grows and shrinks over time. It is best suited for point queries based on the partitioning attribute, but it’s also useful for sequential scans over the entire relation or range queries on the partitioning attribute. If the hash function is good, the number of tuples in each node is roughly the same, i.e. $\\frac{1}{N}$, for $N$ nodes. Nevertheless, it is not well suited for point queries on nonpartitioning attributes.</p>\n<p><strong>3. Range Partitioning</strong></p>\n<p>This strategy partitions tuples by the range of their key values. Tuples are sorted (conceptually), and n ranges are chosen for the sort key values so that each range contains roughly the same number of tuples. We assign contiguous attribute value ranges to each node. Given a set of nodes $N_1, N_2, ..., N_n$, we choose a partitioning attribute $A$ and a <strong>partitioning vector</strong> $[v_1, v_2, ..., v_{n-1}]$ such that, if $i &lt; j$, then $v_i &lt; v_j$. The relation $R$ is partitioned as follows for a tuple $t$, where $t[A] = x$:</p>\n<h5 id=\"primary-and-secondary-horizontal-partitioning\">Primary and Secondary Horizontal Partitioning</h5>\n<p>A <strong>primary horizontal partitioning</strong> is a partition defined by predicates on root tables (source tables): tables that exist as independent entities, without any Foreign Key (FK) pointing to it, like, for example, the CUSTOMER table in the example below. This table can be partitioned by its range of IDs, for example.</p>\n<p>On the other hand, <strong>derived horizontal partitioning</strong> is the partitioning of a relation that results from predicates being defined by another relation (that is, the table depends on one or more other tables to exist and have uniqueness). Tables like ORDERS and INVENTORY are examples of such tables, so that we can partition them according to how their referenced parent table (e.g., CUSTOMER) was already partitioned, keeping related rows co-located on the same fragment for faster joins.</p>\n<h5 id=\"vertical-partitioning\">Vertical Partitioning</h5>\n<p>Vertical partitioning splits a relation $R$ into partitions each containing a subset of $R$’s attributes (columns). In other words, it divides a relation “vertically” by columns. This particular kind of data partitioning was successfully implemented by column-store DBMSs like MonetDB/X100, Vertica, Snowflake and DuckDB, for example, as well as columnar open data formats and OLAP query engines.</p>\n<p>Vertical Partitioning can be expressed as the relational algebra’s PROJECT ($\\pi$) operation:</p>\n<p>$$\n\\pi_{L_i}(R)\n$$\nwhere ${L_i}$ is the list of attributes (columns) of $R$, for $i = 1,..., N$, where $N$ is the number of attributes of relation $R$.</p>\n<p>Suppose that our PROJ table is split by two lists of attributes, $L_{1} = {PNO, BUDGET}$ and $L_{2} = {PNO, PNAME, LOC}$</p>\n<p>As we can see in the example above, PNO, the Primary Key (PK) of the PROJ table is duplicated on the two fragments – $PROJ_{1}$ and $PROJ_{2}$. This is not a coincidence, but a conscious distributed database design decision so that we are able to reconstruct the original relation $R$ from its vertical fragments by using an OUTER UNION operation. In other words, our database design has a <em>lossless-join property</em>.</p>\n<p>A vertical partitioning is said to be a <strong><em>lossless-join decomposition</em></strong> if it’s always possible to reconstruct the original relation $R$ from its vertical fragments $L_1, L_2, ... L_n$. This is more succinctly expressed in relational algebra as:</p>\n<p>That is, the natural join of the projection results <em>exactly</em> in the original relation $R$. Conversely, if the same natural join computation results in a <strong>superset</strong> of the original relation then the decomposition is said to be lossy. This can be expressed succinctly in the relational algebra as:</p>\n<p>Figure 2.2 in Georgiev’s thesis[6] shows examples of a lossy (a) and a lossless (b) decomposition.</p>\n<p>Finally, Ramakrishnan and Gehrke state the following about vertical partitioning:</p>\n<p>  <em>“To ensure that a vertical fragmentation is lossless-join, systems often assign a <strong>unique tuple id</strong> to each tuple in the original relation, as shown in Figure 22.4, and attach this id to the projection of the tuple in each fragment. If we think of the original relation as containing an additional tuple-id field that is a key, this field is added to each vertical fragment. Such a decomposition is guaranteed to be lossless-join.”</em></p>\n<p>As described by Elmasri and Navathe, a set of vertical partitions whose projection lists $L_1, L_2, ... , L_n$ include all the attributes in $R$ but share only the primary key attribute of R is called a <strong>complete vertical fragmentation</strong> of $R$. Therefore, in this case, the projection lists satisfy the following conditions:</p>\n<ol><li>$L_1 \\cup L_2 \\cup ... \\cup L_n = ATTRS(R)$;</li><li>$L_i \\cap L_j = PK(R)$, for any $i \\neq j$,</li></ol>\n<p>where ATTRS(R) is the set of all attributes of R and PK(R) is the primary key of R;</p>\n<p>And to reconstruct the relation R from a complete vertical fragmentation, we apply the OUTER UNION operation to the vertical fragments. Furthermore, for the $R_1, R_2, ... R_n$ vertical fragments of $R$, we have that $R_1 \\cap R_2 \\cap ... \\cap R_n = A$, where $A$ is a unique key attribute of $R$.</p>\n<h5 id=\"hybrid-partitioning\">Hybrid Partitioning</h5>\n<p>A <strong>hybrid partitioning</strong> is the combination of horizontal partitioning and vertical partitioning, also called <strong>mixed partitioning</strong>. In this kind of partitioning, a fragment of a relation R can be specified by a SELECT-PROJECT combination of operations $\\pi_L(\\sigma_C(R))$, and the original relation $R$ can be reconstructed by applying UNION <strong>and</strong> OUTER UNION (or OUTER JOIN) operations in the appropriate order.</p>\n<p>In the Employee table example, if we partition by Department Number (Dno) and Social Security Number (SSN), there are four resulting partitions:</p>\n<h3 id=\"2-allocation\">2. Allocation</h3>\n<p>Allocation assigns a fragment to a particular node. The choice of site and degree of replication will be driven by performance and availability requirements. Some partitions may be replicated across all nodes of the system, while others remain local to a particular site, for example.</p>\n<h3 id=\"3-replication\">3. Replication</h3>\n<p>Replication allows the storage of several copies of data to achieve data locality (i.e., keep the data closer to where it is most needed), availability (the probability that a system is continuously available during a time interval), reliability (the probability that a system is running at a certain point in time) and faster query execution by using local copies.</p>\n<p>The logical units of replication will vary depending on the requirements of the application, performance requirements, patterns of access, types and frequencies of transactions, and availability requirements. Basically, there are two extremes of replication: no replication at all – that is, each relation or fragment is stored in a single node – on one side and whole database replication on all nodes on the other side (called a <em>fully replicated database</em>). In between, we have a spectrum of partial replications of tables and table partitions. For example, H-Store/VoltDB supports both partitioned and replicated tables (that is, the whole table is replicated among all the nodes). The latter is used for small tables that are read-only or not frequently updated like a COUNTRY table, for example, while other tables are horizontally partitioned among the nodes.</p>\n<p>A thorough treatment of replication is beyond the scope of this blog post, but it’s strongly recommended that interested readers check Özsu and Valduriez’s[2], Kleppmann and Riccomini’s[4], or Petrov’s[5] chapters devoted to this subject.</p>\n<h2 id=\"so-what-is-sharding-in-the-academic-perspective-after-all\">So what is sharding (<em>in the academic perspective</em>), after all?</h2>\n<p>All the authors cited so far address <em>sharding</em> in their works, but with slightly different nuances among them. <strong>For Özsu and Valduriez, Elmasri and Navathe, Ramakrishnan and Gehrke, and Silberschatz et al, <em>sharding is a synonym for horizontal partitioning</strong></em>, and they acknowledge that the term is especially widespread in the context of NoSQL databases, Big Data, and Cloud Computing systems.</p>\n<p>But Silberschatz et al also argument that in sharding architectures each node can have a traditional centralized, <em>maybe independent</em>, database instance (like standalone full MySQL DB instances, for example) and the queries would be routed by a middleware. In this architectural model, the sharding query requests can be executed by the database, or a middleware, or even application code. Each shard is an independent and standalone database that replicates all the database schema, but has only $\\frac{1}{N}$ of the data, in a cluster of size N.</p>\n<h2 id=\"key-takeaways\">Key Takeaways</h2>\n<p>Based on what we have seen so far, a more concise definition of <em>sharding</em> could be:</p>\n<p>  <strong>Sharding is horizontal partitioning of data distributed across different nodes.</strong></p>\n<p>And we may establish the following relation:</p>\n<p>  <strong>Every sharding is a partitioning, but not every partitioning is a sharding.</strong></p>\n<p>This derives directly from the fact that there are partitioning on a single node, for example, while sharding usually involves horizontal partitioning + allocation among various nodes. Therefore, a key takeaway relationship among all of the main concepts exposed so far could be:</p>\n<p>$$\nSharding \\subseteq \\text{Horizontal Partitioning} \\subseteq Partitioning\n$$\nOr as a cool cheat sheet:</p>\n<h2 id=\"references\">References</h2>\n<ol><li></li></ol>\n<p>Ramakrishnan, R., &amp; Gehrke, J. (2003). <em>Database Management Systems</em> (3rd ed.). McGraw-Hill.</p>\n<ol start=\"2\"><li></li></ol>\n<p>Özsu, M. T., &amp; Valduriez, P. (2020). <em>Principles of Distributed Database Systems</em> (4th ed.). Springer.</p>\n<ol start=\"3\"><li></li></ol>\n<p>Silberschatz, A., Korth, H. F., &amp; Sudarshan, S. (2020). <em>Database System Concepts</em> (7th ed.). McGraw-Hill.</p>\n<ol start=\"4\"><li></li></ol>\n<p>Kleppmann, M., &amp; Riccomini, C. (2026). <em>Designing Data-Intensive Applications</em> (2nd ed.). O’Reilly.</p>\n<ol start=\"5\"><li></li></ol>\n<p>Petrov, A. (2019). <em>Database Internals</em> (1st ed.). O’Reilly.</p>\n<ol start=\"6\"><li></li></ol>\n<p>Georgiev, N. (2008). <em>A Web-Based Environment For Learning Normalization of Relational Database Schemata</em> (Master’s thesis).</p>\n<ol><li><p></p><pre><code>The commercial database is currently known as Volt Active Data. ↩</code></pre></li><li><p></p><pre><code>The SELECT operation should not be confused with the SQL’s SELECT clause; ↩</code></pre></li><li><p></p><pre><code>While interquery allows the parallel execution of a number of queries on various nodes, intraquery allows the parallel execution of a single query by breaking it up into subqueries and executing on various partitions in parallel. ↩</code></pre></li></ol>","headings":[{"level":1,"text":"Sharding vs. Partitioning: When Definitions Got Sliced and Fractured","id":"sharding-vs-partitioning-when-definitions-got-sliced-and-fractur"},{"level":2,"text":"Then why do people say “partitioning = one machine”?","id":"then-why-do-people-say-partitioning-one-machine"},{"level":3,"text":"How has academia treated these concepts?","id":"how-has-academia-treated-these-concepts"},{"level":2,"text":"Fragmentation, Allocation and Replication","id":"fragmentation-allocation-and-replication"},{"level":3,"text":"1. Partitioning","id":"1-partitioning"},{"level":3,"text":"2. Allocation","id":"2-allocation"},{"level":3,"text":"3. Replication","id":"3-replication"},{"level":2,"text":"So what is sharding (*in the academic perspective*), after all?","id":"so-what-is-sharding-in-the-academic-perspective-after-all"},{"level":2,"text":"Key Takeaways","id":"key-takeaways"},{"level":2,"text":"References","id":"references"}]}}