NOTE

2.5 Kafka Producer

Kafka producer send flow, interceptors, serializers, partitioners, idempotence, transactions, and message-loss handling.

Message QueuesCreated Updated 3 min readhistorical

This is a historical learning note and may contain outdated or incomplete understanding.

1. What Is a Producer?

  • One of Kafka’s clients, responsible for producing messages and pushing them to Kafka.

2. Producer Message-Sending Flow

  • Three objects are involved: the main thread, the Sender thread, and a thread-shared variable, RecordAccumulator (which stores data waiting to be sent).
  • Process: the main thread creates a message and sends it through Interceptors -> Serializer -> Partitioner to the RecordAccumulator. The Sender thread continuously pulls messages from the RecordAccumulator and sends them to the Kafka broker.

2.1. Interceptor

  • Performs processing before the producer sends a message to the Broker, such as filtering messages, counting messages, or modifying message content.
  • A custom interceptor needs to implement org.apache.kafka.clients.producer.ProducerInterceptor.

2.2. Serializer

  • Data sent by the producer needs to be converted into a byte array through a serializer.
  • A custom serializer needs to implement org.apache.kafka.common.serialization.Serializer.

2.3. Partitioner

  • Determines which Partition the message sent by the producer goes to.
  • There are four situations:
    • A partition is specified.
      • Directly use the specified value as the partition value.
    • No partition is specified, but there is a key.
      • Use hash(key) % number of partitions in the Topic to decide which partition receives the data.
    • Neither a partition nor a key is specified.
      • The first call randomly generates an integer (each subsequent call increments this integer), then takes this value modulo the total number of available partitions in the topic to obtain the partition value. This is the commonly called round-robin algorithm.
    • A custom partitioner can be used.
      • It needs to implement org.apache.kafka.clients.producer.Partitioner.

3. How a Producer Ensures Messages Are Not Duplicated

3.1. What Does “No Duplicate Messages” Mean?

  • Ensuring messages are not duplicated is essentially the concept of idempotency.

3.2. Why Duplicate Messages Occur

  • Kafka producers use at least once. The reason is:
    • if a producer sends successfully, the message has already been committed to the log file, with the protection of the multi-replica mechanism;
    • if sending fails, the producer can retry until it succeeds.
  • A producer may write a message repeatedly. After idempotence is enabled, writing the same message multiple times is equivalent to writing it once; in other words, there is only one message.

3.3. How to Ensure It

3.3.1. Enable Idempotence

  • Enable it on the producer client with enable.idempotence=true.
  • At the same time, retries must be greater than 0, acks must be -1, and max.in.flight.requests.per.connection cannot be greater than 5.
3.3.1.1. Implementation of Idempotence
  • Producer ID and sequence number
    • Producer: each producer has a producer ID when initialized. When sending a message to a partition, it increments the sequence number by 1.
    • Broker: maintains a sequence number for a producer ID and partition.
      • If the producer’s sn_new = broker's sn_old + 1, accept it.
      • If the producer’s sn_new < broker's sn_old + 1, it is a duplicate and is discarded.
      • If the producer’s sn_new > broker's sn_old + 1, a message was missed and an error is reported.

3.3.2. Enable Transactions

  • The producer enables enable.idempotence=true and also sets transactional.id.
  • Example
    Properties properties= new Properties();
    properties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,StringSerializer.class . getName()) ;
    properties.put(ProducerConfig . VALUE SERIALIZER_CLASS_CONFIG ,StringSerializer.class . getName());
    properties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG , brokerList);
    properties.put(ProducerConfig.TRANSACTIONALIDCONFIG , transactionid);
    
    KafkaProducer<String , String> producer= new KafkaProducer<>(properties );
    producer.initTransactions ();
    producer.beginTransaction();
    
    try {
    	// Process business logic and create ProducerRecords
    	ProducerRecord<String, String> recordl =new ProducerRecord<>(topic,”msgl ”);
    	producer.send(recordl) ;
    	ProducerRecord<String, String> record2 =new ProducerRecord<>(topic,”msg2 ”);
    	producer.send(record2);
    	ProducerRecord<String, String> record3 =new ProducerRecord<>(topic,”msg3 ”);
    	producer.send(record3);
    	// Process some other logic
    	producer.commitTransaction() ;
    } catch (ProducerFencedException e) {
    	producer.abortTransaction() ;
    }
    producer.close();
  • Consumer
    • Set isolation.level.
      • If it is read_commited, the consumer application cannot see (consume) uncommitted transactions.
      • If it is read_uncommiter, the consumer application can see (consume) uncommitted transactions.
3.3.2.1. Why Have Transactions?
  • Idempotence can only apply to one partition, while transactions can span multiple partitions.
  • Transactions can ensure the atomicity of write operations across multiple partitions.
3.3.2.2. Implementation of Transactions
  • Transaction coordinator.

4. How a Producer Prevents Message Loss

4.1. Set the Buffer

  • Use block.on.buffer.full = true. When the asynchronous buffer is full, block there and wait for the buffer to become available; do not clear the buffer.

4.2. Use Asynchronous Callbacks for Message Sending

  • There are three ways to send messages:
    • asynchronous by default: send(xxx);
    • synchronous: send(xxx).get();
    • asynchronous callback: send(xxx, callback).
  • Use KafkaProducer.send(record, callback). After sending a message, the callback function is invoked. If sending succeeds, send the next one; if sending fails, record it in the log and let a scheduled script scan it later (a reported send failure may not mean the actual send failed; it may only mean no feedback was received, so the scheduled script may send it again).

4.3. Set Retries

  • Use retries and retry.backoff.ms to set the number of retries and the interval between retries.

5. References

Discussion

Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub