How we built our Real-time Analytics Platform
blog.maxcdn.com
blog.maxcdn.com
The biggest take-away for me was that someone is selling ($2500/server/year for hot backup + support) an improved ('accelerated speed for inserting data as well as compression') GPL2 MongoDB.
Also, seems a bit odd to use MongoDB for write optimized work loads, since it essentially has a global write lock. Doesn't Cassandra perform better at writes?
They use it as well for similar purposes. So it isn't a strange choice.
That said, I'd use Cassandra as well. ;) I found when testing the write performance was significantly better. However, I didn't test the special Mongo variant they are using.
As for why we chose the database we did, while TokuMX is basically a drop in replacement for MongoDB, it has features that make it much more usable for our situation. Specifically, a compression rate of over 80% without negatively impacting our insert speed and document level locking. And because TokuMX works as a drop in replacement for MongoDB, we were able to use the MongoDB driver (mgo http://labix.org/v2/mgo) for Go which is just a really great tool. Additionally, we are now able to leverage the MongoDB aggregation framework which allows us to use the data in a way that builds very helpful aggregations of the data very quickly.
That's a fair argument: when written intelligently, Go has great performance.
The other reason is performance v. flexibility. Any generic log collector (Logstash, Flunetd, etc.) needs to support multiple data sources and outputs and is probably less optimized for the specific situation at hand.
Just curious: what's the buffering strategy for your custom agent written in Go? Does it support file-/memory-based buffering?
Working from the end backwards, the first (or last) buffer is a Go buffered channel running within the process on the router which feeds into the database. These channels work as a sort of queue between concurrent Go-routines, or "workers", which have a set amount of allocated space in memory. These are empty most of the time, but if there is a failure with pushing to the database they can start to fill up in order to not block the process before it in the workflow until the system recovers. Before those Go workers on each of the "Router" servers that push to the database is a Redis queue which basically serves as a holding pin in between the CDN servers and the database cluster, which is the data's first stop after leaving it's originating servers.
On those originating servers is the Go process which reads the logs and pushes the data to Redis on the routers, these also utilize buffered channels. So all these layers of buffers work to prevent a block in the workflow during momentary downtimes while the system recovers. However, each of these layers of buffering do have a limit and if they fill up will begin to block on the process before them. In the event of a major failure (such as all members of a replica set being down or the network between the CDN server and the Router being down) the Go process running on the CDN server will stop reading the logs at a position where it lands when it is unable to feed more logs into the channel and hold that position until the channel starts to be drained from the other side and space in the channel becomes available.
Recently, they advertised they now provide server logs to prove usage, but I'm done with them.