On efficiently partitioning a topic in Apache Kafka
arxiv.org
arxiv.org
"...Even though Apache Kafka provides some out of the box optimizations, it does not strictly define how each topic shall be efficiently distributed into partitions. The well-formulated fine-tuning that is needed in order to improve an Apache Kafka cluster performance is still an open research problem.
In this paper, we first model the Apache Kafka topic partitioning process for a given topic. Then, given the set of brokers, constraints and application requirements on throughput, OS load, replication latency and unavailability, we formulate the optimization problem of finding how many partitions are needed and show that it is computationally intractable, being an integer program.
Furthermore, we propose two simple, yet efficient heuristics to solve the problem: the first tries to minimize and the second to maximize the number of brokers used in the cluster.
Finally, we evaluate its performance via large-scale simulations, considering as benchmarks some Apache Kafka cluster configuration recommendations provided by Microsoft and Confluent. We demonstrate that, unlike the recommendations, the proposed heuristics respect the hard constraints on replication latency and perform better w.r.t. unavailability time and OS load, using the system resources in a more prudent way..."
My experience so far with Kafka is that error messages are nearly useless. Even an invalid command line parameter raises an exception, rather than printing a usage. In many cases, if I do something wrong, it just stalls and times out, with no real way to debug.
The documentation caps out at a very nominal level. I tried to implement a simple KVS with compact logs. Documentation ran out and I was wading in deep waters.
I'm trying to build a platform that's developer-friendly, and I'm deep enough in to know Kafka shouldn't be part of it.
That's not to mention the growing uncertainty around the JDK/JVM.
I want open-source, archival, and to be able to reply event streams from the beginning.
For my application, I'll take lower performance in return for better usability.
Yeah, Kafka's complex, and it's really a tool that you need to think long and hard about whether or not your problems are complicated enough that introducing Kafka to solve them will ultimately reduce complexity.
And it's definitely something that developers have to learn the semantics of to use well.
> The documentation caps out at a very nominal level. I tried to implement a simple KVS with compact logs. Documentation ran out and I was wading in deep waters
I occasionally contribute to Apache Kafka, any pointers on what documentation was lacking, and where you could have used more? I will see if I can contribute something around that :)
> That's not to mention the growing uncertainty around the JDK/JVM.
What do you mean by this?
> I want open-source, archival, and to be able to reply event streams from the beginning.
Yeah, maybe check out Liftbridge, it's built on NATS. https://liftbridge.io/
Apache Pulsar might be worth a look, but it's actually more complex under the hood than Kafka, but has a lot of features built-in that either aren't in FOSS Kafka yet, like tiered storage, or won't be until Confluent doesn't dominate the PMC (like an integrated schema registry), or just can't be done very nicely, if at all, like decent multi-tenancy.
That said, it's a fast moving target, the code quality last I looked was patchy in places, ditto the documentation for both it and Bookkeeper, and the admin overhead is higher (managing bookies and brokers and Zookeepers vs. just brokers and ZK with Kafka, or when KRaft is production ready, just brokers).
> What do you mean by this?
Agreed - four years ago Oracle changed all the licensing but I thought they did a bit of a 180 on that in the last year or two.
> What do you mean by this?
Oracle. Oracle is trying to milk the cash cow. That's and unpredictable ride.
> I occasionally contribute to Apache Kafka, any pointers on what documentation was lacking, and where you could have used more? I will see if I can contribute something around that :)
Well, if you're contributing, there are three Python interfaces to Kafka: confluent, pykafka, and python-kafka:
- All three have major gaps.
- None use async
Rallying around a single, usable interface would be a huge step. That's where I'd most encourage contribution.
But to see gaps in documentation, take any of the three, and try to fetch the most recent event (without blocking). You'll run into a bunch of holes along the way. I wanted a simple KVS in Kafka -- store with compact, and fetch most recent event. That would allow the system to run without a DB. That was not trivial to figure out.
https://github.com/Parsely/pykafka
pykafka was originally developed and maintained by my team at Parse.ly, but we no longer maintain it. We instead encourage folks to use confluent-kafka-python, which is what we have ourselves switched to in our production systems:
https://github.com/confluentinc/confluent-kafka-python
(pykafka was developed at a time before Confluent invested in their own Python binding. Some of the history of the project is described in this 2016 blog post[1] and our original 2015 announcement[2].)
[1]: https://blog.parse.ly/pykafka-now/
[2]: https://blog.parse.ly/announcing-pykafka-python-support-for-...
What do you mean by "without blocking" though? You can set a poll timeout of 0ms, that should prevent blocking.
It's been awhile since I dug into this but a couple jobs ago I was very concerned with what happens if replaying a topic from 0 for a new consumer: existing, up-to-date consumers are negatively impacted! As I recall this was due to fundamental architecture around partitions, and a notable advantage of Pulsar was not having such issues. Is that correct? Is that still the case?
I may be describing the problem incorrectly, but I know vendors we talked to were aware of this issue and had workarounds; IIRC Aiven had tooling to easily spin up a temporary new "mirror" cluster for the new consumer to catch up.
It sounds like you're asking if you can double the number of readers in a system with no performance impact. If you're at capacity, the answer is obviously no. Yes, every consumer takes some i/o and CPU on the brokers serving the data. I have never used Pulsar but I'm sure that's also the case there.
It definitely prioritises tail consumption over read from 0.
Partition count only limits concurrent consumption within a single consumer group. One consumer group won't impact another unless its consumers are doing sufficiently bad things to bottleneck the network or disk.
The default consumer read sizes are so small you will hit broker CPU and worker thread limits long before network throughput. (Both consumer batch sizes and broker threads can be increased trivially but there's not much documentation around when to do this.)
And fetching small amounts of data repeatedly doesn't impose much overhead unless you're deliberately disconnecting and reconnecting between polls.
And I assure you, you can definitely bottleneck network before CPU and/or network threads, Kafka was literally designed for very large numbers of consumers.
And as for tuning, Kafka The Definitive Guide is pretty much as the name suggests. I've been recommending it very strongly, especially the chapters on monitoring and cluster replication, for years.
You can download a free draft copy of the 2nd edition from Confluent, check it out :)
But max.partition.fetch.bytes is only 1MB.
1 MiB is allegedly the ideal batch size for throughput.