Kafka without ZooKeeper
confluent.io
confluent.io
Kafka implemented in Go without needing Zookeeper.
I just finished writing a book that shows how to build similar distributed services from scratch, it walks though building a simple distributed commit log with built-in consensus and service discovery from nothing to deployment: https://pragprog.com/titles/tjgo/distributed-services-with-g...
I imagine there are no region restrictions, right? I've already contacted support though.
The goal is to make it lightweight and simple to operate, yet very fast.
I took a different path by divorcing from the Kafka protocol and experimenting with what I believe is a simpler to model approach to reliable event processing [1].
[1] https://github.com/dataptive/styx/blob/master/docs/howto/rel...
I will take the opportunity to say that Kafka is kind of painful, with or without ZK. Check out NATS! [0]. It doesn't solve all the same problems, but is so much easier to use (during development especially) and can do a lot of the same things.
NATS doesn't ever store messages persistently; but this might be fine for your application, and then you don't have to worry about setting 5 different config options to make sure Kafka actually frees up disk space like you expect it to ;)
NATS also enables some unique patterns like request/reply via a "reply to" message header.
Anyway, it's been a joy to use!
Not true. Both Nats streaming and the upcoming jetstream (core nats) do.
To the folks at Synadia -- I love NATS, but the naming and organization of these projects could use some work. What's with the `stan.*` repository names? Where did "jetstream" come from? Why is it baked into `nats-server` but `nats-streaming-server` isn't? Is `nats-streaming-server` on the back burner?
[0] https://docs.nats.io/nats-streaming-concepts/intro [1] https://github.com/nats-io/jetstream
Jetstream is GA with the 2.2.0 release. Folks who believe in waiting for "not .0" won't have to wait too much longer.
If you went there for tap water, yeah, maybe there are better options.
That being said, have you checked out NATS Streaming Server? It’s effectively a first party client for NATS that gives it at least once semantics and persistence, and makes it much more applicable to use cases that are currently on Kafka.
Docs here if you’re curious - https://docs.nats.io/nats-streaming-concepts/intro
[0]: https://docs.nats.io/whats_new_22 [1]: https://docs.nats.io/compare-nats
Disclosure: I work for Confluent
> messages from a given single publisher will be delivered to all eligible subscribers in the order in which they were originally published. There are no guarantees of message delivery order amongst multiple publishers.
https://docs.nats.io/faq#does-nats-offer-any-guarantee-of-me...
Messages are ordered within partitions.
> What systems out there require strictly ordered data? It seems like any design that requires something like that is going to be extremely brittle.
TCP/IP ?
Right, but that means you're still "unordered" across those partitions?
> TCP/IP ?
But TCP/IP isn't delivered in order, it rearranges the unordered packages by their ID. I guess ordered delivery would be nice for that, but I just feel like making your protocol not require ordering is far simpler.
Not to mention that both TCP and Kafka have to handle head of line blocking?
I'm not trying to say that ordering is bad or anything, I just feel like it isn't buying me tons.
Right, so related messages have an ordering guarantee but unrelated messages may be processed out of order relative to each other, which is usually what you want. (Of course you do have to set the record key correctly).
> I'm not trying to say that ordering is bad or anything, I just feel like it isn't buying me tons.
It's a lot more lightweight than full ACID, but if you get your dataflow right it achieves everything that a traditional database does. Without ordering you wouldn't be able to do anything that requires any kind of consistency.
To me, it seemed at odds with the parallelism of a partition, but I suppose in this case you'd be partitioning on some sort of semantic key vs, say, a hash.
Thanks for bearing with me on that, this was just an unfamiliar idea for me.
> both TCP and Kafka have to handle head of line blocking
Well which is it?
(If TCP doesn't give you ordered delivery, why would a head block the rest of the line?)
Maybe if you're an e-business, you'll split everything happening on your website by client id, but still want events belonging to a single client to be received in order, for practicality.
Apache Pulsar offers the same distributed log offering with a fundamentally better architecture, but Kafka has closed most of the gaps now and has far more integrations and a bigger ecosystem.
That being said, I don't think this is what differentiates the two systems, the guarantees they do/don't make are likely what will make the decision for your project.
[0]: https://github.com/nats-io/nats-streaming-server/blob/master...
Going forward, you will no longer need to configure and run a separate ZooKeeper service just to run Kafka. For proof-of-concept projects, a single-process Docker image will be available when running in KRaft mode (non-ZK mode).
For bigger projects, you may want to use a managed cloud service. Or if you do choose to manage it yourself, it will be easier running one service than two.
Disclosure: I work for Confluent.
What is the migration strategy here? Is it doc'd up yet? I am having flashbacks to migration for follower partitions recently which required a decent amount of pre planning of partition layout.
Also as it is pulling in the duties of ZK into kafka what sort of CPU/memory changes are you seeing? Is it 'meh' or all the way to 'you may want to add a couple of CPUs and a few more GB'? Also is it working ok with the stretched cluster?
Also if you want to hit an interesting market you may want to look at 'does it run OK on a raspberry PI'.
Is the single process deployment only doable via a container? Or will we actually have OS native process as well?
Kafka cloud offerings like AWS MSK are quite different, as you still have to do much of the Kafka management yourself. It's not a fully managed service. This is also reflected in the pricing model, as you pay per instance-hours (= infra), not by usage (= data). Compare to AWS S3—you don't pay for instance-hours of S3 storage servers here, nor do you have to upgrade or scale in/out your S3 servers (you don't even see 'servers' as an S3 user, just like you don't see Kafka brokers as a Confluent Cloud user).
Secondly, Confluent is available on all three major clouds: AWS, GCP, and Azure. And we also support streaming data across clouds with 'cluster linking'. The other Kafka offerings are "their cloud only".
Thirdly, Confluent includes many additional components of the Kafka ecosystem as (again) fully managed services. This includes e.g. managed connectors, managed schema registry, and managed ksqlDB.
There's a more detailed list at https://www.confluent.io/confluent-cloud/ if you are interested. I am somewhat afraid this comment is coming across as too much marketing already. ;-)
Disclaimer: I work at Confluent.
My preference is MSK but I'm very comfortable with vanilla Kafka in AWS at a good price with auto-updates.
With a Kafka compatibility shim: https://github.com/googleapis/java-pubsublite-kafka
Disclaimer: I work on GCP.
We run a number of Kafka clusters, most are relatively low trafic, and the management overhead is pretty. Earlier version did require a bit more attention, but mostly it’s pretty simple to deal with.
Kafka is awesome, but using it in local envs is a pain in the ass, if this is never becomes PROD ready it is already an immense achievement to be able to run Kafka locally with less complexity and overhead.
With a Kafka compatibility shim: https://github.com/googleapis/java-pubsublite-kafka
Disclaimer: I work for GCP.
https://cloud.google.com/pubsub/lite/docs https://github.com/googleapis/java-pubsublite-kafka
Disclaimer: I work on this product.
In the process of splitting up everything in modules.
( Microservices would be Overkill)
I don't have enough experience with these new vanity licenses to know what the contribution story looks like, either
from a user perspective, existing SASL + SSL + SCRAM will be released next wednesday - so no code changes.
After Raft, it became easier to just implement that layer yourself and so most projects after Raft (or probably more accurately once people started seeing how stable etcd was, ~2014), just used Raft internally where they would have previously used zookeeper.
IMHO, it's still easier to delegate the consensus problem to a third party service like Zookeeper or ETCD.
Also, do any of those projects publish their Raft implementation as a library for other projects to include?
Most NoSQL databases, now, use Raft, which didn't exist at the time when Kafka was created. Other NoSQL databases, at the time, were not as stable as Zookeeper or had silent bugs that ate data (see aphyr's Jepsen series[1], which thourghly tested several NoSQL databases and found many to be failing, except for Zookeeper).
> Yeah! I mean, I find a lot of linearizability errors in various databases, but this was also my very first time doing this kind of test, and it varies from system to system. Could have easily slipped through the cracks.
In summary, aphyr thought Zookeeper is linearizable even though it doesn't provide linearizable ops.
Looks like Zookeeper needs to be tested again.
Wow indeed... I assume that’s a play on “sudo make me a sandwich”?
See the mailing list thread for 2.8.0-RC0 for where to find the bits if you want to test https://lists.apache.org/thread.html/r16894a11aec73abac521ff... and the project site has some "contact" info for mailing lists where these things are announced and advertised (including releases, Kafka Improvement Proposals, and more) https://kafka.apache.org/contact
Setting up a Zookeeper ensemble is not that hard, they're light on resources and and basically zero maintenance.
The ZK development API is also pretty much awful. Apache Curator (which wraps around the ZK API and implements a bunch of common "recipes") makes it less painful, but it really ought to be part of ZK proper.
I’m still glad to see it go away, one less operational dependency the better.
I have a concern, if the Kafka broker provide both coordination service and Kafka service, how to achieve the resource isolation? If some of the topic which on the coordination service with very high throughput, this must cause instability of the coordination service, could this further cause the instability of the entire cluster?
If some of the broker only provide the coordination service, what is this essential difference? Will this cause more problems for expansion and contraction? Will this bring greater risks when users scale down the brokers? I am very afraid that the coordination service will be shut down due to careless operation.
Apart from Confluent wanting you to use Kafka so they can keep leeching money off you by hijacking de facto ownership of an open source project, of course.
More moving parts. Brokers, and Bookies, and ZK, plus proxies etc. Plus an additional ZK for inter-cluster replication.
Immaturity - it's still early days for Pulsar, and there's still a lot of bugs being found - and then rapidly fixed, full credit to them, but yeah, not yet as stable. Documentation is often obsoleted, and I found myself having read the code to figure out what was actually going on.
More complex workflows - there's only really one model for a developer consuming or producing against Kafka. With Pulsar, there's multiple different subscription modes, and choosing the wrong one could produce problems.
Also, the need to explicit ack the messages is something you'd have to always watch for to avoid duplicated reads. Also, if using batch receive, when I was looking at Pulsar, you either had to acknowledge the entire batch, or none of it, so a failure during batch processing would lead to the batch being reprocessed, but I think acking within a batch is in development.
No Pulsar IO S3 sink yet.
That said, there's a lot of cool things it's doing, like the built-in schema registry and far easier multitenancy, and offloading older data into S3 etc. transparently to the consumers, so I'm definitely I'm keeping an eye on it.
Lastly, you're taking aim at Confluent, you realise Pulsar is largely controlled by people employed by StreamNative, yeah?
Not saying that they won’t turn evil at some point, but so far they’re leagues ahead of Confluent in terms of earning developer trust. At a minimum this developer, but also others that I’ve worked with.
Maybe I’m the minority opinion here and that’s fine, but confluent has been far too shady for me to ever consider contracting with them.
We used FOSS Kafka for yonks without hitting any limitations - At one point we were looking at Confluent Replicator, but decided it was just easier to go with Mirror Maker 1 (and you know, no massive licensing fees) - and Mirror Maker 2 largely emulates Replicator in terms of functionality.
I'm aware of a few other things like the MQTT KC Connector, but that was never part of Kafka in the first instance, it's something that Confluent built for paying customers.
And I could argue that the StreamNative "critical functionality" that they lock away behind enterprise agreements is "quick bug fixes for the many bugs we're still finding", if I was feeling mean spirited.
But anyway, it seems your preference for Pulsar is due to Confluent, but they're not the only ones offering managed Kafkas - AWS, IBM, RedHat, etc. etc.
The Kafka community is huge and the velocity of development is very high. It's easy to forget now, but in the beginning, Kafka didn't even have replication. That's a good reminder that things that seem like permanent advantages of system X over Kafka (for various values of X) may very well prove to be temporary. For example, in this very thread, I see people talking about how various system X'es have the advantage over Kafka because they can run without ZK. Those discussions are almost out of date.
Finally, I work at Confluent and I think the company has always been a positive force in the open source community. I respect the Pulsar people as well, but I think they have a difficult challenge to overcome.
No. No it doesn’t. It has at-least-once delivery with client-side deduplication. That’s not new, it’s what TCP does FFS. Why would you lie to people about supporting something long established at best and demonstrably impossible at worst?
> Finally, I work at Confluent....
Oh, that’s why. Never mind then. Continue selling digital snake oil.
Confluent make big bold claims "Exactly once delivery" and have aggressive marketing.
Pulsar on the other hand would say we have "effectivley-once". Reading Pulsar docs vs Kafka, Pulsar are very modest about functionality and have no commercial marketing at all.
These days I have noticed Confluent in blog posts do use effectively once but marketing is as aggressive as ever.
Credit where credit is due. Confluent, the marketing and big bold claims is why almost everyone is using Kafka and not Pulsar and may not of even heard of Pulsar. I do find Pulsar architecture more interesting, since Splunk has brought them though it's remained in the background like it always has with no huge push to sell it.
It's not just deduplication, you can atomically commit a consumer from one topic + produce of records resulting from that. Which is exactly the same exactly-once guarantee that you get from e.g. an SQL database in linearizable mode (a lot of SQL databases will do the same thing internally - optimistically execute transactions and then re-run them in the case of a conflict).
They contribute heavily to the project and offer enterprise support, professional services, and proprietary tech like Cluster linking.
Source: I am an enterprise customer of Confluent.
-------------------------------------------------------
The problem I see with kafka is that it was built before cloud architectures were commonly adopted (with distributed systems everywhere). Confluent has put a lot of effort dragging kafka's architecture to the present, but some major features are missing:
- Auto-scaling: Confluent finally introduced "elastic scaling" a few months ago but it only allows you to scale up and must be triggered by the admin (no threshold-based auto-scaling).
- Multi-tenancy: Planning for a multi-tenant kafka cluster is not for the faint of heart. Achieving isolation tends toward liberal usage of topics of which starts to become unmanageable in the low thousands. This isn't crazy when you've got a few hundred microservices and several tenants to keep isolated.
- Decoupled brokers and storage: Any broker scaling or failure can lead to downtime while event storage is redistributed.
Confluent's Cloud service reduces operational overhead but isn't always feasible due to cost, resource limits (like service accounts or schemas for instance), data controls, etc.-------------------------------------------------------
With the removal of zookeeper and tiered storage (separated from compute), Kafka has caught up on scalability while being simpler to deploy. It also has a far bigger ecosystem with more polished features like ksqldb.
One of the benefits of the Kafka rearchitecture effort is to allow Kafka to "scale down" to run without external dependencies. Using Pulsar would add more dependencies.
(That said, unlike many I consider depending on ZooKeeper to be a positive sign. "We wrote our own consensus protocol" belongs in roughly the same bucket as "we wrote our own crypto." Using ZooKeeper doesn't automatically mean your distributed system will work but at least you'll have a fighting chance.)
BookKeeper is a feature though. Allows to scale the partition beyond the capacity of a storage unit. Effectively unlimited retention for a partition. The problem with Kafka is that the broker is tied to storage.
And BK isn't a feature without cost, it's documentation is... somewhat sparse, and there is significant complexity to maintaining it.
I remember asking for it 5 years ago: https://radek-gruchalski.medium.com/the-case-for-kafka-cold-.... Confluent turned it into a paid feature.
The fact that MM2 happened, and Confluent didn't try to stop it, despite it being awfully similar to Replicator, makes me think that Confluent are acting in good faith.
Incidentally, I quite like how Pulsar solved tiered storage, and it's a definite tick in the Pulsar box - it's transparent from a consumer's POV, although there somewhat of a delay in rehydrating the offloaded block, I don't think anyone's expecting near-realtime performance when loading historical data.
[1]: https://cwiki.apache.org/confluence/display/KAFKA/KIP-405%3A...
> The fact that MM2 happened, and Confluent didn't try to stop it, despite it being awfully similar to Replicator, makes me think that Confluent are acting in good faith.
Let me share an anecdote related to this example. We (Confluent) were actually the ones who contributed the documentation for MirrorMaker v2 to the Apache Kafka docs (https://kafka.apache.org/documentation/#georeplication). The development lead on MM2 was (an engineer at) Cloudera, yet they never spent the time to provide user-facing documentation to the Kafka project. I don't want to speculate about reasons, yet I noticed that MM2 was documented in the Cloudera docs.
If we didn't care for the Kafka community at Confluent, we would not have spent our own resources and time to fill that gap, given that we have a proprietary product similar to MM2 (i.e., Confluent Replicator).
Hardly the most straightforward, and it was rather a gaping hole. Thanks for the background on how that hole developed.
I really appreciate Confluent putting that time into documenting something vital, that could compete with your own product, and IMO that does put a nail in the previous commenter's assertions about Confluent's alleged attempts to wall off necessary features of Kafka.
One layer BookKeeper provide an abstraction similar to HDFS. That is it provide file that are horizontally scalable in size and throughput and reliable append only files.
Pulsar is a service built on top of BookKeeper but could run on top of HDFS or something like Amazon S3 ...
And is only responsible for making sure there is only one writer per BookKeeper file even if multiple process try sending request to Pulsar to write to the same partition. It also try to balance request across all the brokers.
Add to that their insistence in claiming “exactly once delivery semantics” from Kafka despite that being provably impossible and I don’t see any reason to trust them as a company or pay for their software.
I’ve been sticking to pulsar for all new projects and have yet to hit a single drawback. It scales better, has less fiddly knobs needing adjusting, has cluster management already built in, and supports traditional pub/suv as well as worker queue semantics. It even has Kafka compatible adapters so it’s relatively easy to migrate existing systems.
Kafka played an important role in the history of distributed system design but it’s time to move to something better built and better managed IMO.
MM1 and MM2 are free. You might be getting confused with Confluent Replicator.
> Kafka compatible adapters
For very limited subsets of the Kafka APIs.
Exactly once delivery is impossible. Exactly once processing is possible. TBH, semantically there's very little difference between those two from an end user perspective.
1- size for single topic limited to the size of one machine 2- complex stateful client library that need to know which machine is currently the master for each partition. ....
2. This is generally handled by the client library transparently. Have you ever needed to manage this state manually?
That said, curious to hear your experiences :)
Outside of librdkafka and jvm client, it’s gloves off.
IIRC Confluent has started putting resources into it - I would hope so, given how .NET Core is going.
That said, the state of Pulsar clients outside of the official Java ones was far worse, I was looking into .NET Core ones and the "official" one (Pulsar-DotPulsar) lacked some key features, whereas a third party one, pulsar-client-dotnet, had far more features, but was still somewhat behind the Java clients.
Caveat is that I looked into all of this when Pulsar was at version 2.6, it's not at 2.7.1, so my comments may well be out of date.
It's worth noting that Twitter built their own system (EventBus) that Apache Pulsar largely mimics in design (and the people who started Pulsar at Yahoo had worked on EventBus prior), with brokers decoupled from storage, and then eventually just decided to get rid of it and use Kafka.
https://blog.twitter.com/engineering/en_us/topics/insights/2...
> One catch to this is that for extremely bandwidth-heavy workloads (very high fanout-reads), EventBus theoretically might be more efficient since we can scale out the serving layer independently. However, we’ve found in practice that our fanout is not extreme enough to merit separating the serving layer, especially given the bandwidth available on modern hardware.
I'm pretty surprised Twitter didn't see benefit from doing this if they have multiple Kafka clusters with different use cases.
> I'm pretty surprised Twitter didn't see benefit from doing this if they have multiple Kafka clusters with different use cases.
Yeah, I think they were too tbh. I wish I could delve more into what they experienced beyond that single blog post I linked.
It might depend on what you're ingesting and how much. Being able to independently scale ingest and storage is a good alternative to have. It's not only ingest though. It's also consumption. As it stands, having a parallel consumer over a large partition spanning several GBs also requires tons of RAM because a segment must be loaded into memory. In that sense, reprocessing historical data is pretty difficult. There's a lot of complexity hidden in additional Druid, HDFS installations or shoe-horned object storage with their own indexing to support access to historical data semi-fast.
I don't know too much about Kafka's internals, but that's my not experience of reading several terabytes of data from a Kafka topic. Memory didn't blow out, although we did burn through IOPs credits.
Edit: Kafka apparently does not store a complete log segment in memory, only parts but having many consumers may lead to a lot of churn or a lot of memory consumed. Maybe this is getting better.
The benefit is not just independent scaling of resources, but more useful features like archiving and reading from object storage with infinite history, and faster and more reliable data replication across regions.
I have a use case where this would make sense, I’ll dig in directly!
Many of the hiccups with running Hadoop and friends in containers or on cloud VM's boils down to how hostnames are resolved and advertised; not any significant design issue.
Which makes it all the less reasonable to assume it would be in any way cloud native, when "the cloud" was at best a nascent idea at that point.
And how many years did it take AWS to get any serious traction after launch?
1. Locking service
2. Id generator
3. Bloom filter etc.,
https://tomharrisonjr.com/uuid-or-guid-as-primary-keys-be-ca...
https://news.ycombinator.com/item?id=14523523 (the hacker news comments for the above article are really insightful)
https://www.percona.com/blog/2019/11/22/uuids-are-popular-bu...
https://blog.codinghorror.com/primary-keys-ids-versus-guids/