We had a 3 node system across a wan, where local clients would read from a local node, and if they needed to write, would write to the master across the wan. The issue we had was there was a significant amount of data being written to the master that one or two of the nodes weren't interested in 100% of the time, but it was still being replicated across the wan (and at cost)
so I investigated how the replication mechanism worked - to see if I could gain greater control over what data was replicated. As it turns out, there isn't, but you can emulate the replication yourself if you're interested. Well this is how we did it:
- master node has a capped collection (we call db.messages) - slave nodes have a bit of mongo console javascript that execute tailed cursor querying for messages being inserted - when a message that matches the query inputs is inserted into the messages collection, it gets inserted into a local db.files collection, which clients then read from
The added bonus is that occasionally we do need to replication additional data to particular nodes, so we just craft up a tailed cursor query that finds any messages we can, and pulls them across the wire
we make a fair bit of use of the adhoc querying in mongodb, so that was a massive selling point for us
speed was also a major issue - we have a average write/very high read requirement, and it's really really fast.
finally the platform was a consideration. We're mostly a windows shop, but we'll use linux where needed, and we were prepared to have riak running on linux if thats what we thought was bets, but it just didn't really fit for us.
the only downside was we use Delphi, and the existing Delphi drivers were.. not good, so we wrote our own, which I'm trying to negotiate with my boss so we can push onto github.
good luck!