Status

Current state: "Under Discussion"

Discussion thread: here

JIRA: here

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

Note: This KIP is a result of working with members of Apache Cassandra community who are working on  CEP-44: Kafka integration for Cassandra CDC using Sidecar and one of the limitations they have faced was large message sizes 

Motivation

Kafka has a limit for message size which limits some use cases where they might have messages that are larger than message.max.bytes even after enabling compression on the producer side or after applying serialization formats like Apache Avro or Protocol Buffers to reduce payload size.  Increasing message size indefinitely is not a viable solution as it can lead to performance degradation, memory issues, and instability of the broker. And by looking at some of the enterprise/cloud offerings of Kafka you can see that on average they can offer 8MB to 10MB as max message size.

Popular Patterns to solve this

At the moment of writing this KIP, there are two famous patterns to handle without increasing message.max.bytes

1. Message Chunking (Splitting and Reassembly message)

Break large messages into smaller chunks, send them sequentially, and reassemble on the consumer side.

2. Reference-Based Messaging

Store the large payload externally and send only a reference (e.g URI, databases key) in the Kafka message.

PatternProsCons
ChunkingNo external storage required Complex client logic to split and reassemble messages
Reference-BasedMinimizes Kafka loadExternal system dependency

This KIP is proposing

  1. A serializer in Apache Kafka that implements Reference-Based Messaging as this is the simplest one.
  2. A notion of a Composable  serializer where the client can apply a list of serializers before applying a large message serializer


Note: 

This KIP will benefit CEP-44: Kafka integration for Cassandra CDC using Sidecar 

Public Interfaces

1. Composable Serializer/Deserializer

In some use-cases we might need a way to be able to chain of Kafka serializer/deserializer before applying large message serializer for example need to apply schema or special format like avro or protobuf. We could just create a single LargeMessageSerializer that implemented the Kafka Serializer interface, but we would need to create the different versions (AvroSerializer, ProtobufSerializer) with their support for the schema.
Instead the KIP proposing a composable serializer, that still implements the Kafka Serializer interface, but allows concatenating several serializers to perform what we are looking for here. 

2. LargeMessageSerializer

A configurable Kafka serializer that:

3. LargeMessageDeserializer

A configurable Kafka deserializer that performs:

Proposed Changes

1. Composable Serializer/Deserializer

org.apache.kafka.common.serilalization.ComposableSerializer and org.apache.kafka.common.serilalization.ComposableDeserializer

Configuration

KeyDescriptionTypeDefaultRequired
value.serializersList of serializers class that implements the `org.apache.kafka.common.serialization.Serializer` interface in order. The first serializer is `Serializer<T>` while the remaining serializers must be `Serializer<byte[]>`.StringNoneNo, If not exist the code will fail back to value.serializer 
value.deserializersList of deserializer classes that implements the `org.apache.kafka.common.serialization.Deserializer` interface in order. The last deserializer must be `Deserializer<T>` while the rest must be `Deserializer<byte[]>`.StringsNoneNo, If not exist the code will fail back to value.deserializer 

2. org.apache.kafka.common.serilalization.large.message.LargeMessageSerializer

Configuration

KeyDescriptionTypeDefaultRequired
large.message.payload.store.classImplementation of org.apache.kafka.common.serialization.largemessage.store.PayloadStore .StringNoneYes
large.message.payload.store.timeout.msTimeout for payload store operations. This can't exceed max.block.ms for Kafka producers or max.poll.interval.ms in Kafka consumer. This is a part of basic configurations for any implementation of org.apache.kafka.common.serialization.largemessage.store.PayloadStoreLong10000No
large.message.payload.store.retry.countNumber of retries for payload store operations. This is a part of basic configurations for any implementation of org.apache.kafka.common.serialization.largemessage.store.PayloadStoreInt5No
large.message.payload.store.retry.max.backoff.msBackoff time between retries for payload store operations. This is a part of basic configurations for any implementation of org.apache.kafka.common.serialization.largemessage.store.PayloadStoreLong10000No
large.message.payload.store.retry.delay.backoff.msDelay time between retries back off for payload store operations. This is a part of basic configurations for any implementation of org.apache.kafka.common.serialization.largemessage.store.PayloadStoreLong100No
large.message.threshold.bytes The maximum size of the message is considered large.Long1MB (Default of max.message.bytes)No





3.org.apache.kafka.common.serilalization.large.message.LargeMessageDeserializer

Configuration

KeyDescriptionTypeDefaultRequired
large.message.payload.store.classImplementation of org.apache.kafka.common.serialization.largemessage.store.PayloadStore .StringNoneYes
large.message.payload.store.timeout.msTimeout for payload store operations. This can't exceed max.block.ms for Kafka producers or max.poll.interval.ms in Kafka consumer. This is a part of basic configurations for any implementation of org.apache.kafka.common.serialization.largemessage.store.PayloadStoreLong10000No
large.message.payload.store.retry.countNumber of retries for payload store operations. This is a part of basic configurations for any implementation of org.apache.kafka.common.serialization.largemessage.store.PayloadStoreInt5No
large.message.payload.store.retry.max.backoff.msBackoff time between retries for payload store operations. This is a part of basic configurations for any implementation of org.apache.kafka.common.serialization.largemessage.store.PayloadStoreLong10000No
large.message.payload.store.retry.delay.backoff.msDelay time between retries back off for payload store operations. This is a part of basic configurations for any implementation of org.apache.kafka.common.serialization.largemessage.store.PayloadStoreLong100No
large.message.skip.not.found.error Skip not found error when the payload is not found in the store. This allows the deserializer to skip not-found messages and return empty bytes instead. This is important when external store ttl is smaller than kafka retention.BooleanFALSENo





4.org.apache.kafka.common.serilalization.large.message.PayloadStore

/**
* The contract for any PayloadStore implementation. 
* This parent abstract class will validate the initial configurations that any payload store must have, like large.message.payload.store.timeout.ms, large.message.payload.store.retry.count, large.message.payload.store.retry.max.backoff.ms and large.message.payload.store.retry.delay.backoff.ms. 
* And extract them from the provided configs. 
*/
public abstract class PayloadStore implements Configurable, Closeable {

    Integer maxRetries;
    Integer timeoutMs;
    Long maxBackoffMs;
    Long delayBackoffMs;
    protected Metrics metrics;     
    
    @Override
    public void configure(Map<String, ?> configs) {
        // configure
    }

     /**
     * Publish data into the store.
     *
     * @param data data that will be published to the store.
     * @return {@link PayloadResponse}.
     */
    public abstract PayloadResponse publish(String topic, byte[] data);

    /**
     * Download full data from the store.
     *
     * @param path id of the data's reference in the store
     * @return content of the object as bytes.
     */
    public abstract byte[] download(String path);

    /** 
    * Generate an id for the data's reference in the store. 
    * By default the id is a random UUID however some stores might need more smarter way to calculate its reference id. 
    * In such a case please override this method. 
    * @param data data that will be published to the store.
    public String id(byte[] data) {
       return UUID.randomUUID().toString();
    }
}

5. org.apache.kafka.common.serilalization.large.message.PayloadResponse

/**
* Response from publish / download from PayloadStore back to the serialization layer
* It contains the final path, response code and the encountered exception if there was any. 
* If PlayloadResponse contains PayloadStoreException with isRetryable flag then it will serialiser will
* retry. 
* If the responseCode 404 the deserialiser will skip the error if large.message.skip.not.found.error set to true.
 */
public class PayloadResponse {
    public final int responseCode;
    public final String fullPayloadPath;
    public final PayloadStoreException payloadStoreException;
    /**
     * Construct payload response with response code and payload id.
     */
    public PayloadResponse(int responseCode, String fullPayloadPath) {
        this(responseCode, fullPayloadPath, null);
    }

    /**
     * Construct payload response with response code, payload id and exception.
     */
    public PayloadResponse(int responseCode, String fullPayloadPath, PayloadSto
reException payloadStoreException) {
        this.responseCode = responseCode;
        this.fullPayloadPath = fullPayloadPath;
        this.payloadStoreException = payloadStoreException;
     }
}

6. org.apache.kafka.common.serilalization.large.message.PayloadStoreException

/**
* Exception class that can either be reliable or not
* this helps the serializer/desrializer to decided either to retry or to crash. 
* One subclass will be added is PayloadNotFoundException which is used to indicated if the payload not found 
* This is used by deserializer to skip or not.
**/
public class PayloadStoreException extends RuntimeException {
    protected boolean isRetryable = false;

    /**
     * Constructor PayloadStoreException with message and throwable.
     */
    public PayloadStoreException(String message, Throwable t) {
        super(message, t);
    }

    /**
     * Constructor PayloadStoreException with message.
     */
    public PayloadStoreException(String message) {
        super(message);
    }

    /**
     * Constructor PayloadStoreException with throwable.
     */
    public PayloadStoreException(Throwable t) {
        super(t);
    }

    /**
     * Constructor PayloadStoreException with message, throwable and if it is retryable or not.
     */
    public PayloadStoreException(String message, Throwable t, boolean retryable) {
        this(message, t);
        isRetryable = retryable;
    }

    /**
     * return whether the exception is retryable or not.
     */
    public boolean isRetryable() {
        return isRetryable;
    }
}


Example

Map<String, Object> producerConfig = new HashMap<>();
producerConfig.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
producerConfig.put("value.serializers",
        "org.apache.kafka.common.serialization.DoubleSerializer,org.apache.kafka.common.serialization.LargeMessageSerializer");  producerConfig.put("large.message.payload.store.class", "CustomS3Store")
 producerConfig.put("s3.bucket", "my-bucket")
producerConfig.put("bootstrap.servers", "localhost:9092");

KafkaProducer<String, Double> producer = new KafkaProducer<>(producerConfig);

Consideration: 

Recommendation: Set TTL duration to exceed your Kafka topic retention period by a reasonable buffer (e.g., 10-20%) to ensure payload availability throughout the message lifecycle while preventing indefinite storage growth.

Compatibility, Deprecation, and Migration Plan

Rejected Alternatives