Solutions like Spanner are fantastic for some things (and within Google it's a no-brainer, since you're not paying an outrageous markup), but besides being expensive, they don't just let you drop them in as a scalable replacement for an existing usage pattern, and usually sharding starts first coming up when you're at a huge scale and already have a
ton of app code and database design but have grown past what a single db can support (that's a massive amount of traffic these days, well over 1M DAU for most apps). See for instance
https://cloud.google.com/spanner/docs/schema-design, a lot of what you have to plan for is very different than if you were designing a normal (e.g.) Rails app with typical patterns.
It's a great option, as long as you're aware of the tradeoffs and don't expect it to act (and cost) exactly like a single Postgres server.
My order of ops generally goes:
- Can you get away with pure (or mostly) key/value for your hottest table(s) and mostly avoid queries? Use Dynamo or Google Cloud Datastore, but deeply understand their quirks (e.g. never index a timestamp field on Datastore or you're worse off than a weak-ass SQL server for write-heavy loads). These can scale basically forever with zero effort if you fit their constraints, and are cheap enough
- Can you tolerate the price of Spanner and deal with it not being normal SQL? Go for it, expect some non-trivial migration cost
- If you have to shard/cluster yourself, can you make Redis work? It'll be easier operationally than managing N SQL DBs of any kind. I know, if you could you'd probably have been fine using Dynamo...
- Can you soft-shard, e.g. use different machines for different tables and spread load enough to get each table onto a single server? Do it.
- Can you minimize re-sharding ops? (If you're in a massive growth phase, the answer is no) Ok, fine, do SQL sharding, and make sure you have well-tested code to do it and you mostly use key/value access patterns.
- Consider a beefy AF bare metal SQL server, see how far that gets you.
- If none of this applies, re-evaluate the cost ($ and time) of Spanner, you're paying a big price one way or another
- Only bad ideas from this point on, just do sharding but it'll suck...
Nowhere on my list: running your own cluster for Cassandra, Mongo, or anything similar, it's such an ops nightmare and there are hosted services that are better in every way but cost.
If you do end up writing your own sharding layer, best of luck, I have yet to see any reasonably production-ready ones that work well with pretty much any web framework or ORM. And writing these is so tough, so many edge cases to consider...not something I hope to ever do again.