If this has improved lately, I'd be very interested to hear about the changes!
Since our complexity requirements were low we settled for a home-grown solution (basically memcached with write-through) and so far didn't regret.
I, too, remain curious if anyone is running this at scale and has bumped into the various corner cases (exceeding capacity, hardware dying, etc.).
Automatic sharding and memcached integration are pretty awesome features, and could definitely ease code at the application level (sharding code is a particularly special pain in the ass, not so much getting it working, but allowing for re-sharding migrations if you decide you need more shards, especially trying to do so without downtime which involves all sorts of nasty tradeoffs).
But bang-for-the-buck and reliability wise, I'm still unsure, I've heard very little about this in the wild. Is this more suited for high write v. read ratio situations, or is it aiming more at the Vertica/Greenplum big-data uses, or something else?
NDB was originally created by and for telecoms. So very high write/read rates with very fast response times and very high availability required. Generally not extremely large datasets.
If you know what you are doing, NDB can work extremely well. It supports all of the highend cluster goodies including online software upgrade, online node addition, automatic handling of node failures, geographic async replication, etc...
However, it certainly has a pretty steep learning curve from the admin point of view and it is a bit easy to mess things up. It is a bit brittle due to this, but once it is setup properly and running, it can deliver on the promises.
Check out Postgres or Cassandra.
I know I must have been doing something horribly wrong, but I never could really figure out what it was.