We had a quartet of legacy-ish databases with sharded customer data (had some interconnection, but we sharded in the application-layer, so the database didn't have to move any data internally). It was a bitch to maintain, but it performed OK, by throwing fairly beefy machines at it.
We began a quest to replace the database with MongoDB/HDFS/Riak/CouchBase/Cassandra/whatever. We wrote a script that converted our data to JSON for easier ingestion and began hammering away. Until we discovered that by converting it to JSON (and simplifying some parts of the data-model), our dataset had actually become small enough to fit on a single server in memory (still ~100GB).
So we ditched the cluster-thingie and moved everything to CouchDB. Pricy hardware, but the saved man-hours recouped that within a month or two.
(And in half a year's time, the database will have outgrown whatever hardware we can reasonably throw at it, so then we'll have to start over for real...)