๐ค Kafka Producers & Advanced Features
๐๏ธ Tuning for Throughput
Anatomy Breakdown: A Producer doesn't send messages one-by-one over the network; that would be incredibly slow. Instead, it waits and groups them into batches.
batch.size: The maximum size (in bytes) of a single batch before it is forced to send. (e.g.,65536for 64KB).linger.ms: The maximum time (in milliseconds) the producer will wait for a batch to fill up before sending it anyway.compression.type: Compresses the batch (e.g.,snappy,lz4) before sending, trading a tiny bit of CPU for a massive reduction in network traffic.
If you want High Throughput, increase batch.size and linger.ms.
If you want Low Latency, decrease linger.ms (or set it to 0).
๐ก๏ธ Reliability (Acks and Retries)
When you send a message, how do you know the Broker actually received it?
acks=0: "Fire and forget." The producer doesn't wait for a reply. Extremely fast, but high risk of data loss.acks=1: The Leader Broker replies as soon as it writes to its own disk. (Default).acks=all: The Leader waits until all replica backup Brokers have successfully written the data before replying. Slowest, but safest.retries: If a send fails, the producer will automatically retry sending the message this many times before throwing an error.enable.idempotence: Set toTrueto ensure that even if the producer retries sending a message due to a network glitch, Kafka will safely deduplicate it. This guarantees exactly-once delivery from the producer.
๐๏ธ Serializers
Follow the Data: Kafka doesn't understand JSON, Strings, or Integers. It only understands raw bytes.
- Your app creates a Python Dictionary
{"id": 101, "name": "Ravi"}. - The Serializer kicks in and converts that dictionary into a byte payload:
b'101,Ravi,IN'. - Kafka stores the bytes.
- The Consumer's Deserializer converts those bytes back into an object.
You can use built-in serializers (StringSerializer) or write a Custom Serializer for complex objects.
๐ Custom Partitioners
By default, Kafka uses a hash of your message key to decide which partition it goes to. If you don't provide a key, it uses a RoundRobin strategy to distribute load evenly.
But what if you want all messages from India (IN) to go to Partition 0, and all from the US (US) to go to Partition 1? You can build a Custom Partitioner to intercept the routing logic before the message leaves your app.
๐ Broker Quotas
What Existed Previously: A single rogue application (e.g., stuck in an infinite loop) could accidentally spam millions of messages per second, crashing the entire Kafka cluster for everyone.
Problems Faced: No single app should be able to monopolize cluster resources.
How Present Technology Solves It:
Kafka Admins can enforce Quotas (produce and consume quotas). If an app exceeds its quota (e.g., 5 MB/s), the broker actively throttles (delays) the responses, forcing the rogue application to slow down without affecting other well-behaved apps.
๐งช Practice Drill
// Try answering these:
Q1. If your app is sending millions of messages per second, should you increase or decrease `linger.ms`?
Q2. Which `acks` setting guarantees that data is written to all backup nodes before confirming success?
Q3. Does Kafka natively understand Python objects or JSON strings?
Q4. If you want to force specific user IDs into specific partitions based on geographic location, what must you implement?
Q5. How does a broker prevent a single spamming application from taking down the cluster?
๐ก Click for Solutions
A1. Increase it. Waiting slightly longer (e.g., 20ms) allows massive batches to form, significantly improving network throughput.
A2. acks=all.
A3. No. Kafka only stores raw byte arrays. The Producer must use a Serializer.
A4. A Custom Partitioner.
A5. By using Quotas and Throttling.
โ โ๏ธ Kafka Cluster Configuration | Next Topic โ ๐ค Spark Streaming & Kafka Integration