Sharding, federation, denormalisation: how to read a client's database before writing your first query
A JOIN that works on your laptop can still return the wrong number on a client's system if you don't yet know how they have split their data.

In brief
- The shard key is the most important design decision in a sharded system and the hardest to reverse, so it tells you which axis the client usually queries along.
- Federation makes JOINs between databases complicated. Denormalisation speeds up reads but slows down writes, and the copies can drift apart.
- Sharded systems often accept eventual consistency, so figures you aggregate across shards may not agree with each other at any given moment.
Picture your first week at a retail client. You have read-only access, and your job is to build an agent that answers the question “what was revenue by region last month?” You write SELECT ... JOIN customers ... GROUP BY region, run it, and get a figure so far below the finance team’s report that it makes no sense.
The query has no syntax errors and the database isn’t broken. You have only touched part of the data. The orders table you are reading might be a single shard. The customers table might live in a different database altogether. Or the region column in orders might have been copied in long ago and never updated.
For an FDE, working out how the client has scaled their database comes before the first query. The good news is that there is a clear process for it, and you can get good at it with practice.
Why is the client’s database no longer one block?
As data and users grow, a single server eventually hits its ceiling. Microsoft Azure’s architecture documentation says that scaling vertically, by adding disk, CPU, memory and network bandwidth, only temporarily postpones that limit. Sooner or later the client has to scale horizontally by adding nodes.
There are three common approaches, and each leaves its own kind of trace. Sharding splits the rows: every shard has the same schema but holds its own subset of the data and sits on its own storage node.
Federation, also called functional partitioning, splits databases by function: users in one place, products in another, orders in a third. Denormalisation is applied to an already normalised database to gain performance: reads get faster, writes get slower.
The three often appear together. Azure itself recommends that sharded designs denormalise so that entities usually queried together, such as a customer and their orders, sit on the same shard, which cuts down on separate reads. When you find one technique, look for the other two.
The shard key shows which axis the client queries along
Azure calls the shard key the most important design decision in a sharded system, and it is very hard to reverse. So when you read a client’s system, the shard key tells you which axis their engineers bet most queries would follow.
If the client shards by customer_id, every question like “orders for customer X” lands on exactly one shard and runs very fast. Your question, revenue by region across all customers, has to go through every shard.
Azure also warns that cross-shard queries and transactions are expensive, and that most sharded systems avoid distributed transactions in favour of eventual consistency.
The type of sharding is also worth asking about. Range sharding, for example by date range, suits range queries but spreads load unevenly and is hard to rebalance. Hash sharding avoids hotspots. But if the client uses hash(key) mod N, every time a shard is added or removed, most keys are reassigned, triggering a large data migration.
You can work it out yourself. Going from 4 to 5 shards, a key stays put only if h mod 4 == h mod 5. Of the 20 possible remainders, only 4 meet this condition, so about 80% of keys have to move.
That number affects you directly. If the client is rebalancing, an order may temporarily sit on the old shard, the new shard, or both. Ask before you trust a total.
Worked example: revenue by region
Back to the hypothetical retail client. After an afternoon of reading connection strings and asking the DBA, you can sketch the picture: orders is sharded by customer_id across 3 shards. customers and products live in two separate, federated databases. orders contains a customer_region column, copied from customers when the order is created.
Thanks to this duplicated column, you don’t need to JOIN to the customer database, which is exactly why people denormalise. What remains is to fan out to each shard and merge the results:
from collections import defaultdict
SHARDS = ["orders_shard_0", "orders_shard_1", "orders_shard_2"]
SQL = """
SELECT customer_region AS region,
SUM(total) AS revenue,
COUNT(*) AS n_orders
FROM orders
WHERE order_date >= %s AND order_date < %s
GROUP BY customer_region
"""
def revenue_by_region(start, end):
agg = defaultdict(lambda: {"revenue": 0, "n_orders": 0})
for dsn in SHARDS:
for region, revenue, n in run_query(dsn, SQL, [start, end]):
agg[region]["revenue"] += revenue
agg[region]["n_orders"] += n
for r in agg.values():
r["avg_order"] = r["revenue"] / r["n_orders"]
return dict(agg)
The detail to notice is the line that computes avg_order. Each shard returns only SUM and COUNT, and the average is computed after merging.
Suppose shard 0 has 10 orders averaging 100, and shard 1 has 1,000 orders averaging 50. The average of the two averages is 75, while the correct figure is (1,000 + 50,000) / 1,010, about 50.5.
Now the denormalised column. customer_region is the customer’s region at the time of the order. If a customer moves from Hanoi to Ho Chi Minh City, their old orders still carry the Hanoi label. That is not a technical bug but a business question: does finance want revenue by region at the time of purchase, or by current region? Ask before your agent answers for them.
Suppose they also want revenue by product category, and category exists only in the products database. This is where federation gets awkward, because JOINs between databases are complicated. A practical approach is to pull the product_id → category mapping table (usually small) into memory or a warehouse, then merge at the application layer instead of forcing a cross-database JOIN.
Five steps before your first query
First, draw a map: how many databases there are, what function each serves, which tables are sharded. Connection strings, ORM config files and database names like orders_shard_07 are the quickest clues. Second, find the shard key and the sharding type (hash, range or lookup table), and ask whether any rebalancing is under way.
Third, list the duplicated columns and trace where each one comes from: when it is updated, by which job, and whether it lags.
Fourth, put each business question into one of three groups: touches a single shard, needs a fan-out, or needs data combined from several functional databases.
Finally, agree with the client on acceptable latency, because an eventually consistent system does not promise that numbers match exactly at every moment.
| Technique | Common signs | What your query must do |
|---|---|---|
| Sharding | Several databases with the same schema, numbered names | Fan out, merge with SUM/COUNT, take care during rebalancing |
| Federation | One database per business domain | Combine at the application or warehouse layer, don’t JOIN directly |
| Denormalisation | Columns with the same name across several tables or databases | Identify the source column and the point in time the copy reflects |
With NoSQL systems such as MongoDB, the questions are the same; only the form changes. Relationships there are built by embedding or referencing data rather than with foreign keys. An embedded document is a denormalisation decision, so you still need to ask whether the embedded copy is updated when the original changes.
Mistakes FDEs often make
The most common mistake is thinking you are reading all the data when you are actually reading one shard. The second is taking an average of averages, or running COUNT(DISTINCT) on each shard and adding the results, when the same value may appear on several shards.
The third is harder to spot: running heavy fan-out queries against production at peak hours. Azure stresses that sharding brings long-term operational complexity, from monitoring and per-shard backups to applying DDL changes on every shard. The client’s operations team will not be pleased if your query slows down every shard at once.
Ask whether there is a read replica. Azure also advises that when the bottleneck is on the read side, read replicas and caching should come before sharding.
The last mistake is proposing to “consolidate everything into one database to keep it tidy” in your first week. A shard key is almost impossible to reverse. The client’s system reflects real trade-offs, and your job is to work with it.
How to show this skill on a CV
When reading job descriptions for FDE or solutions engineer roles, look for requirements such as integrating with a client’s existing data systems or working across many kinds of database. That is where this skill comes in.
On your CV, don’t just write “knows sharding”. Write a specific line, for example: built a fan-out query layer across N shards and agreed with the client how to handle figures delayed by eventual consistency.
Next time you get access to a client’s database, don’t start with SELECT. Spend an hour mapping the shards, the functional databases and the duplicated columns first, because that is where most of the right answers are.
Was this article useful?
Thanks for the feedback!