Kafka Removing Zookeeper Dependency
confluent.io
confluent.io
https://github.com/travisjeffery/jocko
I’ve always found things built on JVM are are PITA to deploy (especially when using SSL) so the single Golang binary is a welcome advancement.
It’s all the good things about Kafka (concept, API, and wire protocol) without all the crap (zookeeper dep, JVM foundation)
That's the power of having tunables.
https://blog.discord.com/why-discord-is-switching-from-go-to...
Not to mention: they resolved their Go GC problem, and rewrote the service in Rust as part of a wave of moving services to Rust because they just liked Rust.
Whilst I like the memory ballast as a clever solution, it does smell like a workaround to a solvable problem if a parameter existed to tune it directly.
https://www.elastic.co/guide/en/elasticsearch/reference/curr...
You also often want to stay below 32 gb to get compressed OOPs
Imagine launching every app and having to say ok you get 300mb but not more!
Of course there are ways to use all the available ram when you deploy your kubernetes pods, you just have to consult the documentation.
Honestly this criticism is ridiculous I can’t believe I even bothered to reply.
I keep wanting the same for JVM. It's probably nice that there are so many ways to adjust the JVM to best suit a very wide range of workloads, but from a practical perspective I rarely have spare time to go diving through the guts to experiment. Let me point something at a one-box, and it can tweak parameters to its hearts content, so I can export and apply to the rest of the fleet.
There are also plenty JVMs to pick and some are better than others working with their defaults.
If you don't, you just start the JVM with defaults, and it configures itself for optimal performance, using its idea of optimal.
10x faster. API compat. No jvm.
I also know of other private impls of it. Just makes sense.
W.r.t OSS doing research on it atm.
I had a VP at a major cloud say to me "we will try to use your tech" in person. So doing research before I make the decision
Look at stuff like LMAX. Java can be lightning fast.
Not to mention that a lot of libraries that are immediately avail for c++ (see io_uring) take a while to get ported to java. The cost of JNI for a libaio wrapper is also expensive. Last i bench(few years back), switching to a crc32 with JNI switch alone was 30 microseconds - an eternity - before doing the work.
In any case, we use seastar (seastar.io), which I'm not sure can actually be ported to java. The pinned thread per core makes a lot of sense for minimizing latency and cache pollution, etc. Externally, the feeling that apps in java are slow is real, less because JVM is slow per se, but because writing low latency apps in java is not the idiomatic way and those that do see to extract every ounce of performance of the hardware often look else where since the work is just about the same.
If you want to do this, you can gain a lot of performance with having custom allocators and pools on the JVM as well. E.g. frameworks like Netty have pooling strategies for ByteBuffers. If you go that route, you can also gain a lot of performance on the JVM, might really be competitive.
Unfortunately the JVM still enforces too many heap allocations since value-types are not a thing yet, but it still performs well.
One of interesting things I discovered about C++/Rust vs C#/Java is that the former language family wants you to do a lot of optimizations upfront, wich typically results in good performance. But sometimes you also spend too much time into optimizing something something that won't matter in practice.
Whereas the managed languages are a bit easier to work with by default, but will thereby only yield mediocre performance. However you still have the chance to look into the bottlenecks and improve them by large margins using the right approaches. In C# now even more so than in Java thanks to tools like Span and value types.
If you need really hard limits on collection then that's a tricky problem, but that's also tricky when you're managing memory yourself.
Simple tweaks can go a long way for a lot of developers, but GC performance has been a problem at the last 3 organizations I've been at - and I'm not in the valley or at a FAANG - so it isn't exactly an uncommon scenario for developers.
Unnecessary allocations will be your bottle neck at some point when shutting data around.
Perhaps open source engineers are not "paid" enough to spend countless hours carefully optimizing their programs on JVM. Properly paid engineers working on proprietary technology can surely bend and twist JVM and make it perform. It might be possible, but costly due to complexity and unpredictable nature of it.
It reminds me of world of SQL where you have to run your query with many small modifications and hope that the optimizer will generate a sensible query plan while the query is still readable. That's the cost of building on top unpredictable systems. You wonder why you can't simply get access to the lower level - physical plans - and program on that level directly, since you know what your desired outcome on that level is anyway.
I still wonder why new projects for performant data processing are not written for example in Rust since for this task it has all the upsides and none of the downsides of JVM.
Raft is really easy to parallelize and dispatch to multiple followers async. I measured recently on 3 i3.8xlarge instances which give you 1.2GB/s - and i got around 1.18GB/s sustained -https://twitter.com/emaxerrno/status/1260415381321084929
Also, what's nice about using raft is that if there is a bug, we know it's w/ our implementation and not w/ the protocol. so it gives users sound reasoning.
I consider myself a RAFT expert and worked in kafka for the past 5 years.
If that is the startup time, are Java AOT compilers taken into consideration into the said benchmarks?
No startup times of course. This is for pushing a simple 2 petabyte workload.
No AOT compilers, just download kafka bin distribution 2.4.1
Just launching 6 or 7 of these.
bin/kafka-run-class.sh org.apache.kafka.tools.ProducerPerformance \
--record-size 1024 \
--topic sfo \
--num-records $((1<<31)) \
--throughput $((1024*64)) \
--producer-props "acks=1" \
"client.id=alex.client" \
bootstrap.servers=172.31.31.28:9092 \
batch.size=81960 \
buffer.memory=$((1024 * 1024)) &> ~/nohup1.txt &Also I started to see a trend in books and blog posts regarding how to write Go code towards better performance, so it isn't a given that it excels at performance out of the box.
All of which comes back to the original point that many times isn't the language, rather how it is written and what tools one makes use of.
If Go is so much better than Java, Google would have replaced it already on Android with Go (battery life and such), instead they went with a mix of AOT/JIT with PGO, introducing Kotlin, while gomobile efforts were never given any serious consideration not even for the NDK.
Go will never replace Java for Android because (1) it doesn’t use a VM, so it would need to compile for every arch that android runs on (2) it would require a bug-compatible port of Android. Because Kotlin runs on the JVM and can run on top of Java codebases, it doesn’t have to leap over these hurdles.
There are other reasons you can come up with I’m sure, but those are the biggest that stick out to me. Notably it has nothing to do with “which language is better”.
They optimize for different use cases, and that’s OK. That being said, if you were writing Android from scratch and were only targeting ARM (an equally silly comparison meant to highlight the differences in the languages), you’d be hard pressed to justify Java over Go.
This isn't really true for non-trivial code bases. And even if it is, it's not a big difference in practice, especially with incremental compilation where I find it actually much quicker to change a couple of files and re-run unit tests compared to golang which has to spit out a multi-dozen MB binary each time, taking 5-6+ seconds.
> (2) compile into single binaries with no dynamic links
Java is getting AOT compilation which will do the same.
> (3) natively support concurrency and parallelism with M:N routines:threads.
Java is getting those as well: http://cr.openjdk.java.net/~rpressler/loom/loom/sol1_part1.h...
Google is writing Android from scratch, is it called Fuchsia and Go also doesn't get to play there.
The few parts that were written in Go are scheduled to be rewritten in C++ or Rust, with Dart being the main userspace language.
I entirely admit the possibility, but I see no plan. Hopes and aspirations are not plans.
Instead they went with Rust, C++ and Dart.
Now what have all those languages in common that Go lacks?
What sort of performance measurement uses default configurations? What is even the point of not tuning the GC to your application workload?
I have spent the last 15 years on running Java apps in production and to optimize for the p99 latency is really not that hard. Optimize p99.99999 is a whole different subject though. I don't understand what is the point of comparing a default GC setting that is for hello world applications to a software that is optimized every way possible. Apples to oranges. It is a grat marketing gimmick though. Look ma no performance! Look here, so much faster. We live in a single dimension word, yay!
Seastar is the foundation of Scylla, which shows that rewriting in C++ can deliver magnitudes more performance which is not possible by just tuning Cassandra on the JVM. In fact, Datastax has now copied the Scylla approach in Cassandra but still lags behind drastically with performance.
For any decent SRE out there.
>> magnitudes more performance which is not possible by just tuning Cassandra
Magnitudes?? Are you talking about the order of magnitude? You should read the ScyllaDB performance report first.
https://www.scylladb.com/product/benchmarks/aws-i3-metal-ben...
Avg. 99.9% Latency (ms): 9.9
vs.
Avg. 99.9% Latency (ms): 474.4
While there is no significant latency difference in lover percentile tiers. Where Scylla really shines is TCO. Some companies trade SRE time for license cost, some other companies tune GC.
See Gil Tene's talk on how not to measure latency.
https://www.youtube.com/watch?v=lJ8ydIuPFeU
What matters is what most of your customers will get - p99, p999, p9999, p100.
How is 10ms vs 475 not a major improvement? How is 4 nodes vs 40 not a major improvement? If you're an SRE than how is managing 4 servers with far less tuning and maintenance not a major improvement? Also 99.9% percentile still matters. They're testing with 300k ops/sec which means 300/sec are facing extreme latency spikes that can be enough to fail and/or cause cascading issues through out the application.
There's no metric where Cassandra is better here and you can't tune your way to the same performance in the first place which is the whole point of Scylla. What even is your claim here? Spend more to get less?
What makes it possible to run cassandra/scylla on nodes with TBs of data density is the TWCS compaction strategy from Jeff Jirsa. He was just a cassandra power user at the time, and I like to think that the invention was possible because of Java.
So, next time you read an ad piece from scylla about replacing 40 mid size boxes running CMS with 4 big boxes, don't forget about TWCS.
Scylla is far more than a compaction strategy. If it was that simple, than Cassandra would already be able to do it.
It's an objectively faster database in every metric. Datastax's enterprise distribution has more functionality but core Cassandra is now entirely outclassed by Scylla in speed and features.
The TWCS is just a good example that raw performance is not everything. Performance also comes from things like compaction strategy, data modeling, and access pattern, while users and stakeholders also care about things like easiness to modify, friendly license, and steady stewardship.
Redpanda uses the Seastar framework which was created by the ScyllaDB project. Scylla is high-performance C++ reimplemation of Cassandra and RedPanda seems to be chasing the same thing as an alternative to Kafka/JVM.
As a 12 year adtech veteran who has built ad networks from scratch 3 times, low-latency and high-throughput are critical to ad serving infrastructure and that's why Scylla is such a better alternative to Cassandra. The only other database that gets close is Aerospike, and possibly Redis Enterprise with Flash persistence. It's entirely valid to want similar improvements for event streams as well, and as long as they keep the same external API then you don't lose any of the ecosystem advantages either.
Seastar is a fundamentally different way of programming from what you mentioned above. Let me give you an example. Seastar takes all the memory up front - never gives it back to the operating system (you can control how much via -m2G, etc). This gives you deterministic allocation latency, is just incrementing a couple of pointers. Memory is split evenly across the number of cores and the way you communicate between cores is message passing - which means you explicitly tell which thread is allowed to read which inbox (similar to actors) - i wrote about it here in 2017 https://www.alexgallego.org/concurrency/smf/2017/12/16/futur...
The point of seastar is to not tune the GC for each application workload So to bring that up means that you missed the whole point of seastar. Instead the programmer explicitly reserves memory units for each subsystem - say 30% for the RPC, 20% for the app specific page-cache (since it's all DMA no kernel page cache), 20% for write-behinds, etc. (obviously in practice most of this is dynamic). It is not one dimension as suggested and not apples to oranges. it is apples to apples. You have a service, you connect your clients - unchanged - and one has better latency. that simple.
It may be your experience that when you download a bin kafka say 2.4.1 you change the GC settings but in a multi-tenant environment that's a moving target. Most enterprises I have talked to, just use the default script to startup kafka w.r.t gc memory settings. (they may change some writers settings, caching, etc)
At the end of the day there is no substitute for testing in your own app with your own firewall settings w/ your own hardware. The result should still give you 10x lower latency.
What Scylla is doing is unlocking new performance potential with a C++ rewrite and an entirely different process-per-core architecture that gets around the fundamental limitations of Cassandra and makes it easier to run. This performance and stability has also led to the team making existing C* features like LWT, secondary indexes, and read-repair even faster and better than the original implementations.
What kinds of security assessments are done to guarantee that Scylla is as secure as Cassandra?
The fact that a tool works out of the box and offers you an extensive array of ways to get more out of it is not in any way worse than a tool that just works out of the box. You can use it in exactly the same way with no more mental effort - stick to the defaults.
It was just released as v1.0 a few weeks ago.
Kakfa?
You can do TLS in Java without using keystores, even if keystores do have some advantages, and megacorps seem to like them for all the wrong reasons. Using them shifts some of the complexity of dealing with certificates from the developer to the system administrator. It's an implementation choice.
The only other complaint of yours I could find ITT was about setting min/max heap size... not only is it extremely convenient to be able to do that, it also hasn't ever technically been required, and the defaults have been Good Enough for most uses since Java 5 came out in 2004. The JVM will figure it out for you, and if it's too conservative or too aggressive for you, you can tune it.
If Kafka is a PITA to deploy, fine, maybe that's fair criticism. Not all things on the JVM are pain to deploy.
I deploy single package jvm files most days having not had to consider the deploy environment for a about 4 years.
I really don't think running a go app vs. a java app is ANY different.
You’ve made a straw man and are now working very hard to defend it.
Are you saying this isn't a major problem with Node? In my experience it's a problem with every language that requires a separate runtime.
The JVM is the worst of two worlds because I have the build complexity of an AOT compiler to make my deployment artifact and the deployment complexity of an interpreted language to get the correct runtime pre-installed.
Just about any of the implentations will work fine. You only have to concern yourself with avoiding use of the Oracle runtime in production without a license. If in doubt, install the latest AdoptOpenJDK[0] Hotspot version that's compatible with your app.
It’s also trivially easy to package as a docker image.
What makes things like Zookeeper and Kafka to deploy is they are complex applications. You can write them in any other language and the deployment will still remain as non trivial.
Consider the fact Ubuntu 18.04's default repos come with a pretty ancient version of Java. And if you want a newer version, you have to use some third-party source. And the fact there are both OpenJDK and Oracle JDKs to choose from. Will that be JDK or JRE? Headless? Depending on the combination of JDK and software, maybe the fonts will go all funny. Got multiple JDKs? Some programs will use your default Java, some will invite you to select the JDK to use, some will bundle their own complete JDK. Oh, a program uses JavaFX? No, of course that isn't installed just because you installed Java.
These problems aren't inherent to Java, of course - it's 98% due to Oracle's attempts to inconvenience people into paying for licenses.
There, you're done. Need fonts? Install fonts. Pain in the butt? Where?
Statically compiled binaries are undoubtedly a plus regarding deployment.
Don't drink the Kool-Aid, golang is just a more opinionated and much less capable cousin of Java/JVM. And for many companies, that could be the right trade-off.
[1] https://groups.google.com/forum/m/#!topic/golang-announce/mV...
[2] https://groups.google.com/forum/m/#!topic/kubernetes-securit...
[3] https://groups.google.com/forum/m/#!topic/golang-announce/65...
A package can easily add a dependency on a specific Java version. And unlike many other tools, if you have multiple versions of Java, all you usually need to do is add the right one to the $PATH. Sometimes you may have to set $JRE_HOME. And that's it. You're done. If that's too much of a pain to deal with, I hope you never have to install a tool that requires a specific PHP, Python or NodeJS version.
I don't know of any serious tools written in Java that don't either explicitly say which version(s) are supported. Maybe some tools are poorly packaged. But again, that's not a problem with Java.
I indeed won't admit that "Java apps" are a pain to deploy, because it's a meaningless and overly generalistic statement. Some apps are a pain to deploy, but it's almost always down to the way they're packaged/distributed, and rarely down to the underlying technology.
I remember when docker first came out and people were losing their minds about how they will be able to deploy their applications with all their dependencies in a single shippable bundle, kind of like an uber jar.
Including the JVM. The paths alone.
From the perspective of a sysadmin who just wants to deploy an app that ships as a jar and doesn't come pre-packaged with any convenience features like you mentioned, it's an awful lot to work out and learn.
The entire HashiCorp stack and Cockroachdb comes to mind. I’m sure there are things that are much more complex in the world, but they do some fairly heavy lifting.
> especially not cross-platform
I find it drastically simpler to download a single binary for each platform I need to deploy to(macOS,Linux,FreeBSD) than it is to get a consistent Java environment setup on those same platforms.
I think the big difference is that Java projects TEND to be more 'up front' about their configuration options, with shipped default files with a lot of the options already set to some default value, while projects like the ones you mentioned require you to look up every property yourself and set it to something if you want to change it from the default.
Are you saying that Java projects tend to have more sane defaults and come bundled with more out of the box?
The last Java project I had to deploy was ElasticSearch/Kibana, and there was a LOT of configuration needed that required consulting a lot of disparate documentation.
It’s a nearly universal experience that different language runtimes have very different experiences with dependencies. Someone who ships a Node, Python or Ruby-based tool for example has a tough time, to the point of needing to write a wrapper installer that vendors all dependencies including the runtime itself, just to be sure.
The JVM with Maven and fat JARs isn’t probably as bad as these cases, as mostly it’s just requiring a compatible JRE. This is why Spring Boot apps for example are so popular for deployment - no more app server insanity!
Jlink makes this more like Golang, but isn’t used enough.
That said Golang dependency management for the developer / builder is a history of horrors. Maven has had its issues but has gotten past most of them.
To go back to the original argument, why is Kafka ported from JVM to Go suddenly much easier to deploy? Is it really due to the JVM, or is it perhaps related to design decisions made while porting the application to Go?
Like any other tool, Java has its pros and its cons, but you’re being hyperbolic.
Otherwise, at least with Ubuntu 19.04 which I have, it's simple `sudo apt install openjdk-11-jdk-headless`
For what I do “how do I get Java on the server” is much less difficult than making deployments quicker and more efficient, which is more a problem for our CI/CD harness and integration tests. Neither of which are Java specific per se.
That said if you want to keep your JDK up to date via Docker or AMI of course that’s fine. a jlink JAR does the same thing but I can see the desire for wanting to decouple runtime upgrades from code upgrades.
sudo add-apt-repository --yes https://adoptopenjdk.jfrog.io/adoptopenjdk/deb/
sudo apt-get install adoptopenjdk-<version>-hotspot
Yes, it’s a few extra steps, but it’s not exactly Sisyphean now is it? Given the complexity of modern CI/CD pipelines, and the trend towards continuous deployment, making sure that a JVM is deployed alongside the uberjar is a pretty small ask. Would a single binary be easier? Sure. Would I change languages just for that? No, other factors are more important IMHO.
There’s definitely some laborious parts of the Java packaging setup (the whole resources, meta-inf / manifest setup reeks of YAGNI problems) but similar to Go’s lack of expressiveness being a feature, the Java ecosystem is fundamentally designed for organizations that separate developers from the systems where the software is run, and this is either helpful or hurtful entirely depending upon the organization’s needs.
I’ve never really had a problem deploying anything based upon the packaging - it’s a minor part compared to various obscure configuration files, ConfigMaps, or environment variable injections that bother me more, and that has nothing to do with a package or even language in itself.
I'm particularly curious about how you deal with TLS in the JVM without keystores? I'd love to hear more about that.
Also, I have to laugh about the heap settings not "technically" being required. Sure, but how many JVM-based services have you actually run in production without setting them? If it's more than zero, then I congratulate you on your luck.
Kafka and ZK are a pain to deploy. Not the JVM's fault. Entirely down to the choices made by their respective community. Last I checked it was still virtually impossible to secure a ZK-ensemble. Maybe they'd welcome patches, but when I tried to submit a minor bugfix years ago, I found them to be an unwelcoming community.
You also seem to have misunderstood what causes some apps to be a pain to deploy: the apps, themselves, suck, and would suck no matter what language they had been written in. Apache tends to become the home for a lot of projects that exhibit that particular issue, and Zookeeper does not seem to be an exception.
Kafka likely also has a case of the Apacheisms going on, but also apparently has a non-trivial amount of Scala in it. Complaining about Kafka fits your requirements far more than Zookeeper alone does.
So, so much of Google depends on Go. It is extremely well understood, has an excellent track record and has zero chance of being "abandoned"
The claim that Go does not have a track record is ridiculous. It's whole premise is that it's made to power essential parts of the largest IT operation in the world.
And there's nothing for Google to cancel, it's open source and they even made the effort of translating it to Go so it's super easy to work on.
And I'm not fanboying over Go, haven't written a Go service in years, but denying Go is an excellent platform to build these kinds of apps on is ridiculous.
Honestly Zookeeper in theory is a great idea: Having a centralized service for maintaining config info saves a lot of heartache when dealing with an open source distributed systems project. But in practice, I've never had a smooth experience getting Zookeeper to run consistently, especially with Kafka.
For sure, part of it is that we treat it like oxygen, in that if it's gone for a few seconds everything just dies. But having dealt with similar systems both proprietary and open source, my opinion is that Zookeeper just hasn't risen to the challenges of its users in the past 5 years. If the next generation of software architects want to use open source streaming or distributed systems, Zookeeper needs to be rewritten or removed.
Also shout out to the confluent.io team: I never paid for your enterprise license, but without your blog posts, docker images, or slack room, I never would have been able to get Kafka working. Thanks again!
in the end broker will do RPC to a service (kafka controller/ etcd) and this service will use raft to replicate the state.
It should be exactly the same. And if anything knowing which node are running the raft algorithm help you be more careful with rolling restart and upgrade.
Also, etcd powers many critical open source projects, so there are many institutional eyes that actively contribute to its improvement. IME if we ever encountered an issue at work with ZK, we found it impossible to trace it down to a bug that we could fix and upstream. Etcd’s been easier in this regard.
With regards to Kafka, it's probably easier and more robust to add their own consensus layer rather than switching to etcd - Kafka is already a distributed system built by a team of distributed systems engineers. It makes sense for them to build their own consensus, deeply integrated with the replication mechanism, rather than relying on an external database.
Would still be great to see ZooKeeper made superfluous, though.
Unfortunately, this is not correct. BookKeeper stores a lot of information in ZookKeeper. By extension, Pulsar (which is based on BookKeeper) also stores a lot of metadata there as well.
For example, from the BK documentation ( https://zookeeper.apache.org/doc/r3.3.6/bookkeeperOverview.h... ):
An application first creates a ledger before writing to bookies through a local BookKeeper client instance. Upon creating a ledger, a BookKeeper client writes metadata about the ledger to ZooKeeper. Each ledger currently has a single writer. This writer has to execute a close ledger operation before any other client can read from it. If the writer of a ledger does not close a ledger properly because, for example, it has crashed before having the opportunity of closing the ledger, then the next client that tries to open a ledger executes a procedure to recover it. As closing a ledger consists essentially of writing the last entry written to a ledger to ZooKeeper, the recovery procedure simply finds the last entry written correctly and writes it to ZooKeeper.
On the other hand I feel like projects should try and use open source “building-blocks”, like etcd and zookeeper, when building their distributed systems. Not only does this help iron out correctness bugs, but it also means that more people are familiar with the quirks, limitations, requirements etc.... of these tools. For example, I think I would be frustrated to hear that K8s were implementing their own raft.
This was actually an early complaint leveled against Kubernetes for things like proxying at the node level or implementing DNS. “Don’t reinvent the wheel!” Sometimes better administrative experiences exist only after a component absorbs some function previous systems expose.
Sometimes it’s better to own the parts of the problem that make your system simpler.
Basically versioning + some form of fault tolerance is sufficient for k8s API server.
I wrote something similar but that provide only (lock/lease) using paxos.
The problem with raft ETCD and zookeeper replicates state machine design is that : - machine leaving and joining the ensemble dynamically is something hard to do correctly - optimal quorum size is no more than 5, you can setup other node as observer but it’s hard to decide which of 1000 node should be in the quorum ...
This way it could become the standard for GO and we can stop wasting effort on multiple raft implementation.
For example, there's a great Raft library for Go [1] that any project can use to implement distributed consensus without a separate running program. I find this to be a better approach with the same collective development and community testing advantages but without requiring more operational overhead.
I'd love to ditch Zookeeper (the weak part that falls over more regularly) in our current Kafka cluster.
https://cwiki.apache.org/confluence/display/KAFKA/KIP-500%3A...
It’s been at least a year and a half since we’ve had severe data loss on our homebrew Kafka cluster, but after the first couple you never look at Kafka+Zk the same way... both iirc were due to leadership election bugs that had been reported many months ago with no progress on a solution.
I have no idea how AWS internally puts up with this. I wouldn’t be surprised if they’d replaced the ZK dependency internally years ago.
It cannot be used for leader election as these events are time sensitive and needs to be consistent across the cluster within a short duration, that is the reason we have raft & paxos.