Versions Compared

Key

  • This line was added.
  • This line was removed.
  • Formatting was changed.

...

This page is meant as a template for writing a KIP. To create a KIP choose Tools->Copy on this page and modify with your content and replace the heading with the next KIP number and a description of your issue. Replace anything in italics with your own description.

Status

Current state: "Under Discussion"

Discussion thread: here [Change the link from the KIP proposal email archive to your own email thread]

JIRA: here [Change the link from KAFKA-1 to your own ticket]

Please keep the discussion on the mailing list rather than commenting on the wiki (wiki discussions get unwieldy fast).

Motivation

KafkaProducer is designed to be thread-safe and we encourage users to share a single producer instance across multiple threads. This approach is effective for scenarios where numerous threads append records infrequently and it helps to avoid creating many producers on resource-limited systems. These threads [link]. These threads often need to send records to different topics, which is supported by the current producer implementation. However, another important configuration—configurations acks —is and compression are currently set at the producer level only.

As a result, users needing different settings for each topic must create additional producer instances instead of reusing the existing one and this approach requires 3x producer instances at least.

This is harmful for edge device(resource-constrained device) since users need a producer pool to handle different acks and compression. Additionally, this approach creates many idle producers if a sensor with a specific setting has no data for a while.

In this KIP, we propose adding topic-level acks and compression, enabling users reuse a single KafkaProducer even when different topics require different settings. This change will reduce the need for multiple producers and minimize resource usage on edge devicesThis limitation prevents users from reusing KafkaProducer instance when different topics require distinct asks settings, forcing them to create separate producers for each configuration. Furthermore, this capability has always been available at the RPC level, and with this update, we're simply bringing it back into users' hands.

Public Interfaces

We only add a new field into ProducerRecord constructor.

Code Block
languagejava
titleProducerRecord
linenumberstrue
public ProducerRecord(String topic, Integer partition, Long timestamp, K key, V value, Iterable<Header> headers, Short acks) {
        if (topic == null)
            throw new IllegalArgumentException("Topic cannot be null.");
        if (timestamp != null && timestamp < 0)
            throw new IllegalArgumentException(
                    String.format("Invalid timestamp: %d. Timestamp should always be non-negative or null.", timestamp));
        if (partition != null && partition < 0)
            throw new IllegalArgumentException(
                    String.format("Invalid partition: %d. Partition number should always be non-negative or null.", partition));
        this.topic = topic;
        this.partition = partition;
        this.key = key;
        this.value = value;
        this.timestamp = timestamp;
        this.headers = new RecordHeaders(headers);
		this.acks = acks; // new field
    }

Proposed Changes

two new config TOPIC_ACKS_CONFIG  and TOPIC_ACKS_DOC into ProducerConfig and define the format.

  • Format  : acks.topic => acks.<Topic-A>=<acks-1>:acks.<Topic-B>=<acks-2>
    • example: config.put(acks.topic, acks.money=-1:acks.temperature=1)
  • Format  : compression.type.topic => compression.type.<Topic-A>=<compression-1>:compression.type.<Topic-B>=<compression-2>
    • example: config.put(compression.type.topic, compression.type.money=lz4:compression.type.temperature=gzip)

Name

Type

Importance

Default

Description

acks.topicStringLOW
null

This configuration item will set acks for specific topics.

compression.type.topicStringHIGHnone

This configuration item will set compression for specific topics.

Proposed Changes

  • Adding a new configuration in  ProducerConfig : acks.topic and compression.type.topic,  By setting this configuration item, users can customize the acks and compression for specific topics.
  • Behavior change:
    1. Attach topic-level acks and compression settings to RecordAccumulator#TopicInfo.
    2. Return Map<Acks, List<ProducerBatch>> when RecordAccumulator#drainBatchesForOneNode is called.
    3. Finally, we can get the acks information and group same acks into List<ProducerBatch>> for a node in sender#sendProduceRequests and then send request
    The proposed change is straightforward, as the ProduceRequest already contains an acks field. All we need is to add a new field to the ProduceRecord constructor. This new field will be null by default, and the producer-level acks setting will be used unless users specified otherwise
    1. .

Compatibility, Deprecation, and Migration Plan

  • The old public interface will still be supported is new one so there are no compatibility issues.
  • The ProduceRequest already contains an acks field so we don't need add any additional change at PRC-level [link].

Performance issue

  • There is no performance issue since we will group the same acks batches to broker.
  • producer request is asynchronous operation so we do not need to wait the response.

Test Plan

  • Relevant unit tests and integration tests will be added to demonstrate the functionality.

Rejected Alternatives

  • N/A