Stream and Go: News feeds for 300M users, built on RocksDB and Raft
stackshare.io
stackshare.io
Cassandra is very scalable, but it's not very efficient. The hosting costs for our cassandra cluster were so big that it was infeasible to run it in another region as well.
Appart from that we've had (partial) downtime a couple of times because one node just started going crazy because of unclear reasons.
Keevo solves this by not trying to be cassandra, but be much simpler. It doesn't do schema's or indexes. All it is is a very fast ordered key value store that is stored to disk and replicated automatically to multiple servers (using raft). Any other features we need, we build on top of this usually outside of Keevo itself. This simplicity saves us a lot of hosting costs and makes its performance much more predictable and easier to debug.
Last but not least a very important advantage is control and understanding of the database internals. Because we build Keevo ourselves, we know the performance and consistency tradeoffs it has and can change/improve them when needed.
I hope this helps in understanding our choice. It's definitely not something I would recommend for most companies, but since our product is storage at its core it makes sense for us.
Do you have GC logs? GC lockups are the most common case. Have you used G1GC + Java8?
Sounds awfully similar to Riak. :)
Even if this feature was/is available in Riak, I still think this was the better choice for us. Bringing this core component in house has been a real boon in gaining good and more importantly predictable performance for our API.
See yugabyte that does cassandra+keevo+rocksdb+raft
Was Riak KV ever considered?
- it is not ACID the worst possible way
"Cassandra is not row level consistent,[21] meaning that inserts and updates into the table that affect the same row that are processed at approximately the same time may affect the non-key columns in inconsistent ways. One update may affect one column while another affects the other, resulting in sets of values within the row that were never specified or intended."
"This is true, to a point. I'm firmly convinced that AP is a better way to build distributed systems for fault tolerance, performance, and simplicity. But it's incredibly useful to be able to "opt in" to CP for pieces of the application as needed. That's what Cassandra's lightweight transactions (LWT) are for, and that's what the authors of this piece used. However! Fundamentally, mixing serializable (LWT) and non-serializable (plain UPDATE) ops will produce unpredictable results and that's what bit them here.Basically the same as if you marked half the accesses to a concurrently-updated Java variable with "synchronized" and left it off of the other half as an "optimization."Don't take shortcuts and you won't get burned."
- there are several things that needs to be tuning (like GC) that is not trivial to do
- modeling is challenging say the least, easy to create hotspots that the developers are not aware of
"Cassandra is not row level consistent,[21] meaning that inserts and updates into the table that affect the same row that are processed at approximately the same time may affect the non-key columns in inconsistent ways. One update may affect one column while another affects the other, resulting in sets of values within the row that were never specified or intended." This is not true. Cassandra is row level atomic but I guess it also depends on the version you use maybe. Can you talk about which version and provide a test that satisfies your assertion ?
"there are several things that needs to be tuning (like GC) that is not trivial to do". Do you know that Go's GC is way less advanced than Java GC as of today? This is admitted by Google's Go team lead.
I guess these days people can make up whatever they want without providing any valid tests that prove their assertions. It all comes down to is either love & hate of a programming language or someone wants to put some fancy sounding tools in their resume!!
You can go ahead and try to convince people who experienced this in production.
http://datanerds.io/post/cassandra-no-row-consistency/
I could provide you all of the versions but it is irrelevant from the angle that the version my clients have in production are affected.
> Do you know that Go's GC is way less advanced than Java GC as of today?
Not only I know, I work as a consultant who actually configures it so that it provides the best performance for the workload a customer has. I made a large sum by changing the settings that most Cassandra installations have. In fact, all of my clients were running Cassandra with the default settings in production and failed miserably to have a stable service. In one occasion there was almost 1 minute GC time on average, causing nodes marked as down, read and write tied up as well. Facebook and Netflix has production engineers with these skills so they can make it happen for them but clients not having such engineering resources are either reaching out consultancies or try to hire somebody with the skills. What does your question have to do with the topic?
> I guess these days people can make up whatever they want without providing any valid tests that prove their assertions.
I guess these days people can go on HN and write down some of their experiences that they encountered in their professional life.
I won't comment on the other points - but I managed a medium sized Cassandra cluster for a couple of years, and the GC point is valid. It has nothing to do with Java's GC being more advanced. It's easier to bypass the Go GC with stack allocations (not possible in Java), and many of Cassandra's processes (compaction, repair) end up being very GC heavy. GC tuning ends up being a function of your workload and if you ignore it, background processes like repair can adversely affect projection nodes, or throw nodes in a loop. - and adversely giving the node more memory can make things worse. Cassandra has left a bad taste in my mouth for Java-based databases.
GC of death:
Write latency:
Typical Cassandra user has these sort of problems. I have one client where they identified the issue (data modeling) and fixed it by redesigning their tables before I got there. Some of the Cassandra users are not even aware. The company where the pictures are from engaged me because the system did not meet with business requirements anymore. Could not insert data into the cluster, their ETL jobs were running for 20 hours and if one failed they did not have data for their business. We could speed it up to run it for 4 hours without remodeling the data. With remodeling it Cassandra was not the bottleneck anymore.
Cassandra is usually competitive in open benchmarks.
JVM does stack allocation whenever it is possible, you can learn about this by googling "escape analysis".
http://cassandra.apache.org/doc/latest/faq/index.html#what-h...
"What happens if two updates are made with the same timestamp? Updates must be commutative, since they may arrive in different orders on different replicas. As long as Cassandra has a deterministic way to pick the winner (in a timestamp tie), the one selected is as valid as any other, and the specifics should be treated as an implementation detail. That said, in the case of a timestamp tie, Cassandra follows two rules: first, deletes take precedence over inserts/updates. Second, if there are two updates, the one with the lexically larger value is selected."
Does this sound to you as atomic? My expectation of atomic is that 2 things changing the same value are ordered and both executed. "lexically larger value is selected" does not qualify for me as atomic.
This is yet another Python success story. You created a viable product with Python and moved critical parts to Go when there was real need for doing so.
Considering Go has only been your primary language for 9 months, as the beginner's style code base is optimized for all that Go has to offer you'll continue to see more performance gains, albeit probably not as significant as you saw between v1 and v2. Good on you for not being tied down with optimizations and actually releasing.
"With Python we often found ourselves delegating logic to the database layer purely for performance reasons. " -- you mean you wrote postgres functions.. or.. views!? gasp! heresy! Why would you migrate something working perfectly well in postgres away to Go, though? You probably didn't achieve any performance gains with that. Were postgres functions really gobbling that much memory up? Or.. maybe my assumptions are wrong here.
SendGrid engineer here.
IIRC, each of the three founders wrote components in the language that they were most productive in at the time. So there were services in Python, Perl, or PHP depending on who wrote it. Seems like a solid strategy to take the fastest path to MVP when you're still trying to prove there's a market for what you're building.
As the company grew, Python emerged as the most commonly chosen language for new projects.
As you said, the initial appeal of Go was its performance. An occasional side-effect of "fast" is needing less hardware to do the same work, so its efficiency was also attractive.
I think a key moment that got us interested in Go came at an internal hackathon. One of our engineers built a prototype replacement for a core piece of our MTA. It was something like 6x faster than the existing Python version. The ability to produce that in such a short period of time demonstrated that Go was also a surprisingly productive language.
These days, most of our engineering teams choose Go for new backend services -- even those that aren't performance-sensitive. However, we still have a lot of active development in Python (and Ruby).
Disclaimer: I've been at SendGrid about 4 years, so some of this is second-hand. I'll ask the founders if they're up for doing an article on our engineering blog about it.
I only use Python for the same reason I use Java: all the NLP / ML frameworks built by people who don't know Go yet.
On the Go side, since they mentioned that rocksdb is written in C++, I am interested to know the cgo issues they might had when using Go and rocksdb together, but of course they didn't mention that. Same goes for the GC stw pauses.
I strongly believe that articles like this largely used as soft ads materials without much detailed technical meat should be banned from Hack News. If it is all about "we built system X using building blocks Y and Z in language T!", what is the real value for readers?
Using Rocksdb with Go has indeed its own challenges because CGO calls are far from cheap and have an important consequences on Go applications. For performance reasons, we ended up moving some logic from Go to C in a few hot-spots (eg. moving the retrieval of large amount of keys from Go to C).
Disclaimer: I work for Stream and I am one the two co-founders
On that subject, is there a Badger[0] based distributed key-value store, equivalent to etcd, ready for production?
I would love to know also the "business history" of the product: how you started, the first sale, the marketing approach and how you move from there...
Starter, $59/month, 5 million updates
Growth, $269/month, 9 million updates
Wait, what? To increase the updates by 80%, the price jumps by 356%.
I understand they're throwing in higher processing and custom ranking. That's still approximately the most absurd pricing variance on one plan to the next on a scaling service that I've ever seen.
What you maybe should have said is (speaking from a sales perspective): the reason there's a 356% jump in price, is because of an incredible value proposition step up of 1,000% in what you're getting!
I think that price gap is a classic, and very large, opening for a competitor. There's a ~$89-$129 style plan missing in there, that is a logical step up from $59. That's obviously just my opinion. I also happen to be a big fan of charging more rather than less, in the service space; that under-charging is a common mistake young companies make. And I still can't make sense of $59 to $269 as a step.
I launched a new business in the last quarter. It has a user feed system that is fairly standard. I've built several feed-like features in the last decade. I looked at Stream a while back, having run across it while trying to avoid building my own again. The good news is, I like your product, it got my attention. That price jump between plans, for a very modest increase in service volume, made it a non-starter. It's too easy and cheap to build my own, which then gives me tight control over its evolution over time as well. It's basically impossible at your price points to compete with what I can just do for myself relatively quickly; that $269 product can be run on a low-cost Digital Ocean droplet in terms of its actual resource demands. You're charging for the quality of your software product primarily, and some support, which makes perfect sense. The problem, is that I can build that and scale it to the $899 plan and save myself $10k a year. If you could do 10 million updates at ~$50-$75 / month, with normal plan stepping from there, it'd be interesting. 250k to 350k updates per day is not an immense sum, just 20,000 daily active users can saturate that pretty easily.
However, I can tell you that our features that are hardest to implement in a scalabale way are aggregated feeds (e.g. one message with "tom watched 10 videos") and custom ranking (reddit/HN style ordering with time decay and points). If you're feed system doesn't have either of these features it might indeed be possible to do it cheaper than our more expensive plans.
However keep in mind, a lot of the cost is in the availability and reliability of our service. Basically for every type of server we run (including databases), we need to run at least two to be able to stay online in case one dies. Finally building and maintaining a feed system is not cheap in terms of developer cost (try finding a dev for $10k a year).
- 6.3 million feed updates per day per server
- 111 API requests per minute per server
If you guess 4 cores per server it even sounds worse.
(Ps, will HN ever support better formatting?)
Since we can only speak in broad terms, those requests/updates aren't going to be evenly distributed throughout the day and they might not be consistent from day to day so some over-provisioning is reasonable.
Since they're on EC2 a) they shouldn't be that over-provisioned though and b) are getting comparatively (to bare metal) awful performance (say 2x-8x depending on exactly what they're doing) and paying for the privilege.
I'd like to better understand why "Fanout-on-write" wasn't their go-to solution. All they said was that it was "expensive".
- High Availability requires 2x or 3x more servers
- We run our own monitoring, log collection and tracing infrastructure
- ML/Analytics features
- Distribution of traffic
- Disaster Recovery
Edit: fix list formatting
A news feed for 300M is child's play when you consider that the content is pretty much static.
"Go and RocksDB and Raft" hit all the keywords but what makes this post relevant for HN?