Posts

Showing posts with the label Multi-thread

Improving performance of Kafka Producer

Image
  As we know, Kafka uses an asynchronous publish/subscribe model. While our producer calls the send() command, the result returned is a future. That future offers methods to check the status of the information in the process. Moreover, as the batch is ready, the producer sends it to the broker. Basically, the broker waits for an event, then, receives the result, and further responds that the transaction is complete. For latency and throughput, two parameters are particularly important for Kafka performance Tuning: Batch Size Instead of the number of messages, batch.size measures batch size in total bytes. That means it controls how many bytes of data to collect, before sending messages to the Kafka broker. So, without exceeding available memory, set this as high as possible. Make sure the default value is 16384. However, it might never get full, if we increase the size of our buffer. On the basis of other triggers, such as linger time in milliseconds, the Producer sends the informa...

Multi-threaded Apache Kafka Consumer

Image
Why do we need multi-thread consumer model?  Suppose we implement a notification module which allow users to subscribe for notifications from other users, other applications. Our module reads messages which will be written by other users, applications to a Kafka clusters. In this case, we can get all notifications of the others written to a Kafka topic and our module will create a consumer to subscribe to that topic.   Everything seems to be fine at the beginning. However, what will happen if the number of notifications produced by other applications, users is increased fast and exceed the rate that can be processed by our module?   All the messages/notifications that haven’t been processed by our module, are still in the Kafka topic. However, things get more danger when the number of messages is too much. Some of them will be lost when the retention policy is met (Note that Kafka retention policy can be time-based, partition size-based, key-based). And more im...