Why Database Sharding Is No Longer the Default Choice
The article argues that manual database sharding is shifting from a default scaling strategy to a fallback option, as native distributed databases like TiDB and OceanBase internalize sharding, routing, and distributed transactions, reducing application complexity, though they introduce higher latency, operational overhead, and cost, making single-node MySQL with proper indexing preferable for smaller datasets.
Why We Partitioned That Way
The author recounts sharding an 80-million-row order table into 8 databases and 64 tables. At the time, single-node MySQL on mechanical disks with 32 GB RAM could not handle 20–30 million rows per table; indexes would thrash. Business grew 10×, and the database would have died without sharding.
The Cost: Permanently Embedding a Routing Layer in Business Code
Middleware like ShardingSphere performs six steps — parse SQL, optimize executor, route, rewrite, execute, merge — that are natively database responsibilities but are pushed into the application layer. The configuration below shows sharding rules hard-coded on the business side.
# ShardingSphere-JDBC typical configuration: sharding rules hard-coded on the business side
rules:
- !SHARDING
tables:
t_order:
actualDataNodes: ds_${0..1}.t_order_${0..1} # Database-table mapping, entirely manually maintained
tableStrategy:
standard:
shardingColumn: user_id # Wrong shard key choice can only be fixed later by data migration
shardingAlgorithmName: order_inlineThis means sharding is "configured" into existence, not grown organically by the database.
The Trouble Comes in Every Little Thing After That
Queries without the shard key force a fan-out across all shards with application-side aggregation, multiplying latency. Cross-shard JOINs, deep pagination, global unique IDs, and distributed transactions all require manual implementation. The most painful operation is scaling — expanding from 8 databases/64 tables to 16 databases/128 tables demands rehashing, double-writes, data migration, consistency verification, and business pause windows at every step.
Native Distributed Databases Bring That Layer Back Into the Database
TiDB, OceanBase, and PolarDB have emerged. Their logic is simple: sharding, data migration, and cross-node transactions are the database's job. The application sees a single logical database. OceanBase uses partitioned tables; partitions auto-distribute across nodes and queries benefit from partition pruning. TiDB separates compute and storage, scaling each independently — SQL execution becomes a scheduling concern, not a business concern.
The Most Obvious Difference Is in the Code
The SQL text is identical, but in the sharding era you must configure hints, split the query into multiple parts, and merge results yourself. In the native distributed era, the execution plan already shows distributed parallel computation.
-- Sharding era: no shard key, full-cluster scan awaits
SELECT COUNT(*), status FROM t_order WHERE create_time > '2026-01-01' GROUP BY status;
-- Native distributed era: one SQL runs across nodes, cross-node aggregation handled by the database
SELECT COUNT(*), status FROM t_order WHERE create_time > '2026-01-01' GROUP BY status;But Don't Rush to Declare Sharding Dead
Native distributed databases are not a silver bullet. Single-node query latency is higher due to multi-node scheduling. Minimum deployment is not cheap: TiDB requires compute, storage, and PD components; OceanBase needs enough replicas for quorum. Operating such a cluster demands distributed-systems expertise. For data under a terabyte with simple query patterns, the most practical solution remains a large-memory MySQL instance with a covering index.
Sharding's True Historical Shift: From Default to Fallback
Sharding has moved from "default option" to "fallback option." Previously, data growth triggered immediate sharding. Now the first question is: "What problem am I actually facing? Storage bottleneck, poor index design, or untreated slow SQL?"
My Verdict
New projects should not design shard keys upfront. Implement SQL standards and data governance with structures that the database can support natively. Only when a single node truly hits its ceiling should you consider upgrading to a distributed database or vertical scaling. For existing sharded systems that run well, don't rush migration — the benefits may not outweigh the risks. The real anxiety belongs to systems with wrong shard keys or routing logic baked into business code; they will eventually pay technical debt. Early preparation is far easier than late rescue. Sharding won't vanish tomorrow, but it has exited the technology-selection conversation. That is progress: we no longer need to embed database concerns into business code.
Signed-in readers can open the original source through BestHub's protected redirect.
This article has been distilled and summarized from source material, then republished for learning and reference. If you believe it infringes your rights, please contactand we will review it promptly.
Architect's Tech Stack
Java backend, microservices, distributed systems, containerized programming, and more.
How this landed with the community
Was this worth your time?
0 Comments
Thoughtful readers leave field notes, pushback, and hard-won operational detail here.
