ClickHouse Keeper: A ZooKeeper alternative written in C++
clickhouse.com
clickhouse.com
Looks like folks at StreamNative did as well, with their Oxia project: https://github.com/streamnative/oxia. They were just talking about this yesterday at Confluent Current ("Introducing Oxia: A Scalable Zookeeper Alternative" was the title of their talk). https://streamnative.io/blog/introducing-oxia-scalable-metad...
Seems to be a trend :)
I hadn't seen Oxia before but the idea, for their implementation, of making Zookeeper more like Bookkeeper was an interesting one.
Not right for ClickHouse needs but, IMO, a novel approach.
I mean, I have worked and, and been guilty of tooling driven development (RiiR anyone?) .
But, also, in a comment below Alexey shares many of the reasons other than language. I think Oxia does a good job of sharing their approach in - https://github.com/streamnative/oxia/blob/main/docs/design-g...
(Alexey's comment, FYI, https://news.ycombinator.com/item?id=37677324)
My guess if anything is that people always complain about ZK being annoying as a second piece of infra to distribute for simple setups, so my guess this is just a prelude to them embedding their keeper into the DB deployment itself which is the same general strategy Kafka is taking (a verrrry loong time) to rollout.
I mean in a typical use case of it you would
- run it on long running nodes (e.g. not lambda or spot instances)
- run more or less exactly 3 nodes up to quite a cluster size, I guess some use-cases which involve a lot of serverless might need more
- configs tend to not change "that" much nor are they "that" big
what this means is
- java needing time to run-hot (JIT optimize) is not an issue for it
- GC isn't an issue
and if you look at how much memory (RAM) typical minimal nodes in the cloud have it in context of typical config sizes and cluster sizes is also not an issue
through I guess depending what you want to do there could be issues if you use it
- for squeezing through analytics or similar
- setups with very very very large constantly changing clusters, e.g. in some serverless context with ton of ad-hoc spawned up instances, maybe using WASM and certain snapshot tricks allowing insane fast startup time
- you want to bundle it directly into other applications, running on the same hardware and that applications need more memory
but all of it are cases it wasn't designed for so I wouldn't call it an "alternative" but a ZooKeeper like service for different use-cases, I guess
If you're running a Lambda function in which startup time is extremely important, or an embedded application where size and resources are paramount, or even just a short-lived process where you don't care either way... then AOT makes a lot of sense.
But for long-running server processes, just-in-time compilation almost always results is better performance than AOT compilation that cannot optimize at runtime based on what's actually happening.
HN should be full of people who know better, but these discussions feel like piping information into /dev/null. Web devs, students and hobbyists, and other low-information voters just have it in their heads that AOT is always a superior model, and JIT always an inferior fallback, and there's nothing you can say to break through that. There aren't enough people from the business server side world who spend enough time in online discussion forums to correct the narrative.
With that said, this is not the case for both JVM and .NET. Stepping away from ought to is, both have JIT compilers which produce better optimized code than their AOT counterparts, due to a variety of reasons including R&D effort done for JIT throughout their history and JIT allowing to dynamically profile and recompile the code according to its execution characteristics (.NET's Tier 1 PGO Optimized and HotSpot JVM's C2).
Theres truth to a point I guess, but it’s also true that lots of “more information” is useless, and sometimes harmful to optimization, and additionally that it’s really diminishing returns.
Truth is as long as you’re not over relying on RAII, JVM will virtually never outperform C++.
That's why there is profile guided optimizations for C/C++.
Which instruments you C/C++ code to collect that needed additional information (on the cost of performance).
Then when recompiling you can feed that collected information into your system.
The problem with profile guided optimization is that it's way more annoying to deploy as on every update you have to deploy it twice once slower then naive and then once faster. And because the slower part might very well be to slow you might want to only deploy it to some nodes of a load balancer and deploy naive compiled versions to the other.
It also means you have a similar slower => faster startup time, excepts it's of your program as whole instead of for each restart of a node.
And while optimized Java likely will never outperform optimized C++ that is Java specific furthermore if we speak about common non manual optimized code which isn't implementing some tight math/CS algorithms (i.e. very common daily code in many companies) then the difference really isn't that big, small enough to make choosing java over C++ for server stuff a "in general" the right choice. (If you are not a company which only gets the "best" programmers like google).
I mean just to put it into context there are companies which had success with stuff like running high speed trading code on the JVM and it was a success. So if you can do that probably Java doesn't have a major performance problem.
> not over relying on RAII
you mean like throwing out all the major improvements of C++ which majorly reduced the probability of non highly expert programmers introducing bugs which could be turned into RCEs?
RAII is dogshit because it encourages and hides the fact of what’s really happening, leading to crazy performance gotchas. Just knowing the gotchas around it is often enough for any half way competent programmer to devise better code.
That doesn't mean that in practice the JVM is faster than well-written C++; the Java language semantics almost prevent it from doing so in the general case. But in principle, if all else were equal, it should be able to be.
For a long running, memory safe, server process, there's no other ecosystem (that I know of) that is quite like the JVM.
[0] https://www.usenix.org/system/files/osdi20-balakrishnan.pdf
I think there are plenty of other projects (e.g. FoundationDB, Kafka) that also replaced their usage of ZooKeeper as their systems matured. I guess I'm confused why anyone has been picking up new installations of ZooKeeper.
But: every such system is slightly different in the data model and the set of available primitives.
It's very hard to build a distributed system correctly, even relying on ZooKeeper/Etcd/FoundationDB. For example, when I hear "distributed lock," I know that there is 90% chance there is a bug (distributed lock can be safely used if every transaction made under a lock also atomically tests that the lock still holds).
So, if there is an existing system heavily relying on one distributed consensus implementation, it's very hard to switch to another. The main value of ClickHouse Keeper is its compatibility with ZooKeeper - it uses the same data model and wire protocol.
Keeper is a really interesting challenge and we're really open to any kind of feedback and thoughts.
If you tried it out and have some feedback for it, I encourage you to create an issue (https://github.com/ClickHouse/ClickHouse), ask on our Slack, ping me directly on Slack... (just don't call me on my phone)
And don't forget that it's completely open-source like ClickHouse so contributors are more than welcome.
In this case though, the blog outlines specific reasons why this had to be in C++ (interoperability w. their C++ codebase) as well as benefits that are separate from the language.
Possibly, a lot of us enjoy developing X in Y for the sake of doing so. Not everyone may end up caring about the value that the user gets.
:mug:
Built in s3 storage immediately sold me. I’ve used something called Exhibitor to manage ZK clusters in the past but it’s totally dead. Working with ZK is probably one of my least favorite things to do.
Do note the docs page...
https://clickhouse.com/docs/en/guides/sre/keeper/clickhouse-...
In particular, it is necessary to enable the `keeper_server.enable_reconfiguration` flag. It is pretty exhaustive coverage but if there is an important use case missing, let us know!
If anyone has any questions, I'll do my best to get them answered.
(Disclaimer: I work at ClickHouse)
Note: our stress tests have found a segmentation fault in Python's kazoo library.
We only wanted to test Keeper, but found every bug around it :) Let me find a link.
2. It could be configure to store - snapshots; - RAFT logs other than the latest log; in S3. It cannot use a stateless Kubernetes pod - the latest log has to be located on the filesystem.
Although I see you can make a multi-region setup with multiple independent Kubernetes clusters and store logs in tmpfs (which is not 100% wrong from a theoretical standpoint), it is too risky to be practical.
3. Only the snapshots and the previous logs could be on S3, so the PUT requests are done only on log rotation.
If one out of three nodes disappears, but two out of three nodes are shut down properly and written the latest snapshot to S3, it will restore correctly.
If two out of three nodes disappeared, but one out of three nodes is shut down properly and written the latest snapshot to S3, and you restore from its snapshot - it is equivalent to split-brain, and you could lose some of the transactions, that were acknowledged on the other two nodes.
If all three nodes suddenly disappear, and you restore from some previous snapshot on S3, you will lose the transactions acknowledged after the time of this snapshot - this is equivalent to restoring from a backup.
TLDR - Keeper writes the latest log on the filesystem. It does not continuously write data to S3 (it could be tempting, but if we do, it will give the latency around 100..500 ms, even in the same region, which is comparable to the latency between the most distant AWS regions), and it still requires a quorum, and the support of S3 gives no magic.
The primary motivation for such feature was to reduce the space needed on SSD/EBS disk.
Looking at the initial pull request, is it correct that ClickHouse Keeper is based on Ebay's NuRaft library? Or did the Clickhouse team fork and modified this library to accommodate for ClickHouse usage and performance needs?
edit: I'm not even saying this facetiously, it would be freaking awesome.
I've also written about it here: https://mrkaran.dev/posts/clickhouse-replication/
ClickHouse Keeper was released as feature complete in December of 2021.
It runs thousands of clusters, daily, both in CSP hosted offerings (including our own ClickHouse Cloud) and at customers running the OSS release.
Never accept any claims at face value and always test. But, in this case, it is quite battle-hardened (i.e. the Jepsen tests run 3x daily https://github.com/ClickHouse/ClickHouse/tree/master/tests/j...).
It is definitely opinionated and influenced by our work...but not designed solely for it.
But, also, we continue to improve. Most notably in the work on Multi-group Raft - https://github.com/ClickHouse/ClickHouse/issues/54172
Maybe neither of those two things matter for one's use case, but it's similar to someone rolling up on this blog post and saying "but what about etcd" -- they're just different, with wholly different operational and consumer concerns
Also scary how disciplined US society is, with everyone informally-required to actively cheer the proxy war (while no company "stands with Yemen" for example).
E.g. if we quote ZooKeeper:
> ZooKeeper is a centralized service for maintaining configuration information, naming, providing distributed synchronization, and providing group services.
and ClickHouse
> ClickHouse is the fastest and most resource-efficient open-source database for real-time applications and analytics.
Like this are completely different use cases with just a small overlap.
> ClickHouse Keeper is a drop-in replacement for ZooKeeper
That opening was about ClickHouse in general, but the article is about one particular application using the database.
Generally, ClickHouse Keeper provides the coordination system for data replication and distributed DDL query execution for ClickHouse clusters.
and now knowing more about it I would say the answer to my question is a clear yes, it used ZK for something ZK wasn't at all intended for
which means it makes a lot of sense that they replace it
yes I confused the quotes
but the question in general about weather ClickHouse Keeper and ZooKeeper are designed for completely different use-cases which happen to be similar and can work with the same interface still stands
and by now knowing more then my original comment I would answer the question with yes
It runs thousands of clusters, daily, both in CSP hosted offerings (including our own ClickHouse Cloud) and at customers running the OSS release.
Never accept any claims at face value and always test. But, in this case, it is quite battle-hardened (i.e. the Jepsen tests run 3x daily https://github.com/ClickHouse/ClickHouse/tree/master/tests/j...)
But yes, ZooKeeper is pretty amazing. We are building on the backs of giants.
I'd also argue the RAFT v. ZAB is an important production scale conversation. But, as the blog says, Zookeper is a better option when you require scalability with a read-heavy workload.
https://pradeepchhetri.xyz/clickhousekeeper/ talks about some experiments in exactly that vein.
I thought stuff were supposed to be rewritten Rust /s
I was waiting for that somewhere ;)
Thats a very short-sighted view IMO, software engineering is not just about technology choices.
From just this week: https://news.ycombinator.com/item?id=37600852
Just because someone built a fantastically functional building doesn't mean we can't criticize their choice of foundation. Case in point: Millennium Tower in SF: https://www.nbcbayarea.com/investigations/series/millennium-...
But more interesting, to me, is language adoption and familiarity by region.
I have a bookmarked dev.to article from 2020 that discussed programming language popularity by state - https://dev.to/eduecosystem/what-is-the-most-popular-program...
I'm uncertain if anyone has extrapolated that to more geographic regions. It would be interesting.
You are correct that bug free software does not exist. But choose a “memory safe” language does not prevent that. A seasoned C++ developer knows how to use memory sanitizers and other tools to guarantee the correctness of its code compared to an average Rust developer that just trust the compiler which, guess what, also may have bugs.
Any ideas on that?
I know, I read the blog:
> ClickHouse, Inc. is a Delaware company with headquarters in the San Francisco Bay Area. We have no operations in Russia, no Russian investors, and no Russian members of our Board of Directors. We do, however, have an incredibly talented team of Russian software engineers located in Amsterdam, and we could not be more proud to call them colleagues.
The FUD is really hard to overcome. This is coming from someone who advocated for Clickhouse, sent some PRs, and did a minor code audit.
As engineers we focus too much on the implementation details and not the benefits to the user.
How about:
- ZooKeeper alternative with lower latency
- ZooKeeper alternative with lower memory use
- ZooKeeper alternative with predictable overheads
(I don't know if these are true, just suggestions)
Also, we apply many different faults in our Jepsen tests which are run 3 times a day and we never had a problem with leader election. I know this doesn't confirm that there is no bug in it but it's pretty reassuring I would say.
I like your suggestions!
Some of the benefits we summarized in this summary page https://clickhouse.com/clickhouse/keeper include ease of setup and operation, no overflow issues, better compression, faster recovery, (dramatically) less memory used, etc..
There was actually a reason why C++ was important for us at ClickHouse, and it's because C++ is our main code base and managing a Java project as part of it was not natural, but you right - for standalone use of this alternative, that doesn't matter.
A more revealing question is, "How are you dealing with memory safety in this implementation?" There are ways to improve memory safety in C++ through tooling and idiomatic style. Are these things being used?
There are Keeper only tests, but we run ClickHouse with Keeper for all of our server tests.
For each test we try to use all useful tools for verifying safety and correctness like sanitizers.
E.g. an interesting tool we introduced in our codebase for thread safety https://clang.llvm.org/docs/ThreadSafetyAnalysis.html#
We found some issues using sanitizers in our codebase and NuRaft library itself which were instantly fixed.
And let's not forget about Jepsen which showed some really tricky bugs but were more related to the correctness.
I would suggest looking into CBMC and similar tools as well. Model checking is incredibly useful.
The real key for using it, in my opinion, is to isolate individual classes and functions. Avoid instrumenting code with recursion and loops, and focus on defining and verifying function contracts, class invariants, and resource / memory lifetimes.
It will require a significant amount of work to mock up standard library and third party library APIs, but the real beauty of CBMC is that once you define the interface contracts for these APIs and libraries, you can verify every use of them.
I used CBMC previously to verify proper usage rules with C / JNI integration. JNI can be one complicated beast, and CBMC handily managed rule checks for its use.
I'm an extremely careful developer who unit tests everything and strives for 99% coverage. CBMC was still able to detect a memory overwrite flaw in a networking library I wrote that was based on undefined behavior due to integer promotion and offset math. This passed the various sanitizers and unit tests I had in place, but CBMC was able to reduce it to an actual crash condition that was potentially exploitable.
I don't think I can over-emphasize the usefulness of this tool.
Several distributions to chose from.
You can use the OpenJDK distro shipped in your Linux distro (RedHat, Debian, etc.), you can use Microsoft's OpenJDK distro[1], you can use the Eclipse OpenJDK distro, you can use Amazon's OpenJDK distro [3] and there are a whole bunch more.
[1] https://www.microsoft.com/openjdk [2] https://adoptium.net/ [3] https://aws.amazon.com/corretto/
We'll look in to them and adding to some of our social promotion over the coming weeks. Will try to find a way to give you credit.
But the latest bugs found by ClickHouse continuous integration system in the related library were fixed about a year ago:
https://github.com/eBay/NuRaft/pull/373 https://github.com/eBay/NuRaft/pull/392
1. Snapshots and logs take much less amount of space on disk due to better compression.
2. No limit on the default packet and node data size (it is 1 MB in ZooKeeper)
3. No zxid overflow issue (it forces restart for every 2 bn transactions in ZooKeeper)
4. Faster recovery after network partitions due to the use of a different distributed consensus protocol.
5. It uses less amount of memory for the same volume of data.
6. It is easier to setup, as it does not require specifying the JVM heap size or a custom gc implementation.
7. A larger coverage by Jepsen tests. (This could be hard to believe, but true - ZooKeeper is tested by Jepsen, but Keeper takes the existing tests and adds more).
8. The possibility to store snapshots and previous logs on S3.
C++ isn't a key detail, just a consequence of the fact that the main ClickHouse code base is written in C++.
If you need a distributed consensus system but not necessarily compatible with ZooKeeper, there are plenty of options: Etcd, Consul, FoundationDB...