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...)
You need a rdbms. Field names won't be repeated for each query and you will push later the time to 'start over'.
I expect that slimming down the data model was the main factor, or you simply didn't want to use SQL at all (or your data wasn't really structured enough).