firas bouzazi

2026-09-23

Keeping the search space small

Most data we care about lives somewhere on disk, and most of the time we want to read it back quickly. Whenever I set out to make a query faster, SQL usually, sometimes NoSQL, I fall back on the same mental model: a query is a search through some space of rows, and its speed is mostly decided by how big that space is. So the question I keep asking is simple. How much data does the engine have to look at before it can answer me?

Once you frame it that way, most of the tricks below turn out to be the same move wearing different clothes. The table has 10⁶, 10⁹, 10¹², maybe 10¹⁵ rows, and the job is to shrink the slice the engine actually reads. Here are the tools I reach for, roughly in the order I reach for them.

One caveat before the numbers. The arithmetic here is a model for reasoning, not a description of what your database literally does. Real engines read pages, cache hot data in memory, keep statistics, and let a planner pick a strategy. Treat the figures as orders of magnitude, not benchmarks.

Change the schema

The size of the search space is not just the number of rows. It is also how wide each row is, because the engine pays for every byte it walks past. I will use SQL words (row, table), but swap in document and collection and the reasoning holds.

Say we have this table (Postgres syntax):

CREATE TABLE account (
  id          UUID PRIMARY KEY,
  name        TEXT NOT NULL,
  state       TEXT NOT NULL,
  description TEXT NOT NULL,
  created_at  TIMESTAMP NOT NULL
);

Suppose most of our queries only want id and name, while description is long in practice, easily a thousand characters or more. The hot query is:

SELECT id, name FROM account WHERE name = ?;

And suppose, for the sake of the example, that we cannot index name. Then this query has to do a full table scan.

With five columns, each row costs roughly 5n, where n is the average column size. It is a crude approximation, columns vary wildly, but it is good enough to reason with. At 10⁹ rows we walk past:

1,000,000,000 × 5n = 5,000,000,000n
(n = average column size)

The fix is to stop dragging the wide, cold columns through the hot query. Split the table in two:

CREATE TABLE account (
  id   UUID PRIMARY KEY,
  name TEXT NOT NULL
);

CREATE TABLE account_details (
  id          UUID PRIMARY KEY REFERENCES account(id),
  state       TEXT NOT NULL,
  description TEXT NOT NULL,
  created_at  TIMESTAMP NOT NULL
);

The hot query now touches a table that is two columns wide instead of five:

1,000,000,000 × 2n = 2,000,000,000n

That is about 40% of the original, so roughly 2.5 times less to scan (5n / 2n). This is vertical partitioning: we sliced the table down its columns, keeping every row in both halves.

The catch is the flip side of the win. Any query that does need the cold columns now pays for a join, and writes touch two tables instead of one. Reach for this when the hot path is narrow and the cold columns are genuinely heavy.

Indexing

An index is a separate data structure, usually a B-tree, that points at the rows of a table and is shaped so that searching it is cheap. It is the highest-leverage, lowest-effort tool in this whole list, so in practice it is the first thing I reach for, not the second.

Back to the original query:

SELECT id, name FROM account WHERE name = ?;

We add an index on the column we filter by:

CREATE INDEX account_name ON account(name);

A B-tree lookup is O(log_b n), where b is the branching factor. Its real value depends on the engine and is often in the hundreds or thousands, which only helps us. Take a conservative b = 10. Instead of scanning 5 × 10⁹ units of data, the lookup is on the order of:

log₁₀(1,000,000,000) = 9

From billions down to single digits. There is a constant cost on top, because the index and the table are two structures: the engine walks the index to find row pointers, then reads the actual columns from those locations. Even so, no other single change on this list comes close.

We can push it one step further with a composite index that also carries the columns we select:

CREATE INDEX account_name_id ON account(name, id);

Now name and id both live in the index, so the query can be answered straight from the index, skipping the second read into the table. This is a covering, or index-only, scan.

Indexes are not free, though, and it is worth being honest about the bill. Every index is a copy that has to be kept in sync, so it taxes writes and eats storage. Index the columns you actually filter or join on, not every column you have, and remember the planner is allowed to ignore an index when it judges a scan cheaper.

Partitioning

In a database, partitioning usually means table partitioning: taking one table and splitting it into several subtables by the value of one or more fields, the partition key. You also pick a strategy for mapping values to subtables. Unlike the vertical split above, every partition has the same schema as its parent. The common strategies are range, list, and hash:

// Range
partition_key in (0, 10)  -> table_0
partition_key in (10, 20) -> table_1
partition_key in (20, 30) -> table_2

// List
partition_key = a -> table_a
partition_key = b -> table_b
partition_key = c -> table_c

// Hash, with 3 partitions
hash(partition_key) % 3 = 0 -> table_0
hash(partition_key) % 3 = 1 -> table_1
hash(partition_key) % 3 = 2 -> table_2

Take the account table again, with one extra column:

ALTER TABLE account
  ADD COLUMN country_code INTEGER NOT NULL;

Assume we still have around 10⁹ rows, and that we are almost always after accounts from a specific country. The queries look like:

SELECT (...) FROM account
WHERE country_code = ? (AND ...);

Since country_code is nearly always in the WHERE clause, it is a natural partition key. We partition by list:

CREATE TABLE account (
  id           UUID PRIMARY KEY,
  name         TEXT NOT NULL,
  state        TEXT NOT NULL,
  description  TEXT NOT NULL,
  created_at   TIMESTAMP NOT NULL,
  country_code INTEGER NOT NULL
) PARTITION BY LIST(country_code);

Each country code gets its own subtable, sharing the parent schema:

country_code = 0 -> account_0
country_code = 1 -> account_1
country_code = 2 -> account_2
...
country_code = n -> account_n

Say there are 10 country codes, spread more or less evenly, and the filter is almost always present. Then only one partition of the ten has to be read:

1,000,000,000 × 0.1 = 100,000,000

A tenfold cut, and the gain scales with the number of partitions:

10 partitions   -> ~10% of the data each   -> 10x smaller
100 partitions  -> ~1% of the data each     -> 100x smaller
1000 partitions -> ~0.1% of the data each   -> 1000x smaller

Each partition also keeps its own indexes, so combining partitioning with indexing (which you almost always do) gives you smaller indexes too, faster again for the same reason.

The important limit: partitioning only helps queries that name the partition key. A query without country_code in its predicate has to visit every partition, and if you overdo the partition count the planning overhead starts to bite. It buys you nothing for access patterns that ignore the key you chose.

Sharding

Sharding is the last one to reach for, and it looks like partitioning with one decisive difference: the pieces live on separate physical nodes. Partitioning splits a table into subtables inside a single database; sharding spreads that data across multiple databases. The dividing strategies are the same, but each slice now belongs to its own machine.

Suppose we have the full account table again and decide on 5 shards, meaning 5 databases. Shard count and country count are independent, so we can route by:

country_code % 5 = 0 -> db_0 (shard_0)
country_code % 5 = 1 -> db_1
country_code % 5 = 2 -> db_2
country_code % 5 = 3 -> db_3
country_code % 5 = 4 -> db_4
(5 = number of shards)

Because the shards are separate machines, the application has to know which one to ask. Some databases, MongoDB for instance, handle routing for you; with Postgres or MySQL you usually shard at the application level, holding connections to every shard and deciding per query where to go, then assembling the results when more than one shard is involved.

With 5 evenly loaded shards you get roughly a fivefold win, and like partitioning it only applies to queries that fit the scheme:

SELECT * FROM account WHERE country_code = 1;
// one shard, 20% of the data

SELECT * FROM account
WHERE country_code = 1 OR country_code = 2;
// two shards, 40% of the data

SELECT * FROM account WHERE name = ?;
// every shard, since that name could live anywhere

One nice property: even when a query hits several shards, each shard is an independent database with its own CPU, memory, and disk, so the work happens in parallel.

The reason sharding comes last is cost, and it is a real one. Changing a schema, adding an index, or partitioning all happen inside a single database that holds all your data. Sharding does not. Even with a system that routes and assembles for you, you now own a distributed system, with the routing, rebalancing, cross-shard joins, and operational weight that come with it. Consider it only when everything simpler has genuinely run out of room.

An order of operations

In practice, start simple. For the vast majority of cases, a schema shaped around your access patterns plus the right indexes (and queries actually written to use them) is enough. Indexing in particular does more for less than anything else here.

When data keeps growing, it usually plays out like this. You tune the schema and add indexes, and it is fast until it is not. You scale the box vertically, more CPU and memory and a faster disk, and it is fast again until it is not. You keep the schema and indexes and partition the hot tables, tens or hundreds of partitions, maybe thousands. If reads dwarf writes, as they often do, you add a read replica or two, and a cache here and there. Only when all of that has run out, and not before, do you shard.

Closing thought

Every technique here is one idea in different forms: make the space the engine has to search through small. Do that, and it almost does not matter how big the data gets underneath. A small set of data is always fast to work with.


Further reading

Back home