DUE TO SPAM, SIGN-UP IS DISABLED. Goto Selfserve wiki signup and request an account.
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.
- How it works
- Split the payload into chunks smaller than
message.max.bytes. - Assign the same key to all chunks to ensure they land in the same partition, preserving order.
- Reassemble chunks on the consumer using metadata (e.g., sequence IDs, total chunk count).
- Considerations
- Requires custom logic for splitting/reassembling
- Consumers must handle out-of-order or missing chunks
- Available open-Source implementation:
2. Reference-Based Messaging
Store the large payload externally and send only a reference (e.g URI, databases key) in the Kafka message.
- How it works
- Upload the payload to external storage (e.g., S3, HDFS, or a database).
- Send a reference (e.g.,
s3://bucket/key) via Kafka. - Consumers fetch the payload using the reference
- Considerations
- Reduces Kafka network/storage load.
- Introduces dependency on external systems.
- Available open-source implementation
- There isn't a specific open-source implementation for Reference-Based Messaging. Mostly everyone adopting this has their custom solution
| Pattern | Pros | Cons |
|---|---|---|
| Chunking | No external storage required | Complex client logic to split and reassemble messages |
| Reference-Based | Minimizes Kafka load | External system dependency |
This KIP is proposing
- A serializer in Apache Kafka that implements Reference-Based Messaging as this is the simplest one.
- A notion of a
Composableserializer 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:
- Check if the estimated size of the data (bytes) after applying provided compression (if there is one) it needs to serialize is larger than the configured threshold.
- If it is large than the provided threshold (`
large.message.threshold.bytes`):- Use the provided PayloadStore implementation to publish the large message into payload store, generating an id which is the reference to access this later.
- Encapsulate that ID into a simple Kafka event using a structured format.
- Pass the new Kafka Event down
- Add a
large-message: trueheader
- If it’s not large then provided threshold (
large.message.threshold.bytes):- Do nothing, pass the data as it is.
- If it is large than the provided threshold (`
3. LargeMessageDeserializer
A configurable Kafka deserializer that performs:
- Check if the event is large message by looking for
large-message: trueheader- If it does have `
large-message` header:- Parse the event and retrieve the ID and the needed useful information on the event.
- Use the provided PayloadStore implementation to download the original data from the payload-store.
- Return payload as the final value
- If it doesn't have the
large-messageheader:- Do nothing return the Kafka message as it is
- If it does have `
Proposed Changes
1. Composable Serializer/Deserializer
Configuration
| Key | Description | Default | Required |
| value.serializers | List 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[]>`. | None | Yes |
| value.deserializers | List 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[]>`. | None | Yes |
2. LargeMessageSerializer
Configuration
| Key | Description | Default | Required |
|---|---|---|---|
large.message.payload.store.class | Implementation of org.apache.kafka.common.serialization.largemessage.store.PayloadStore . | None | Yes |
large.message.payload.store.timeout.ms | Timeout 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.PayloadStore | 10000 | No |
large.message.payload.store.retry.count | Number of retries for payload store operations. This is a part of basic configurations for any implementation of org.apache.kafka.common.serialization.largemessage.store.PayloadStore | 5 | No |
large.message.payload.store.retry.max.backoff.ms | Backoff 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.PayloadStore | 10000 | No |
large.message.payload.store.retry.delay.backoff.ms | Delay 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.PayloadStore | 100 | No |
large.message.threshold.bytes | The maximum size of the message is considered large. | 1MB (Default of max.message.bytes) | No |
3. LargeMessageDeserializer
Configuration
| Key | Description | Default | Required |
|---|---|---|---|
large.message.payload.store.class | Implementation of org.apache.kafka.common.serialization.largemessage.store.PayloadStore . | None | Yes |
large.message.payload.store.timeout.ms | Timeout 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.PayloadStore | 10000 | No |
large.message.payload.store.retry.count | Number of retries for payload store operations. This is a part of basic configurations for any implementation of org.apache.kafka.common.serialization.largemessage.store.PayloadStore | 5 | No |
large.message.payload.store.retry.max.backoff.ms | Backoff 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.PayloadStore | 10000 | No |
large.message.payload.store.retry.delay.backoff.ms | Delay 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.PayloadStore | 100 | No |
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. | FALSE | No |
4. 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();
}
}
6. PayloadReferenceValue
/** PayloadReferenceValue represent the final published payload reference into Kafka.
* It includes the full payload path on the store as well as the store class used for the publishing.
* This allow the deserializer to throw a more informative errors if the original payload store class doesn't match the one setup by the consumer client.
*/
{
"apiKey": 0,
"type": "payload-reference",
"name": "PayloadReferenceValue",
// Version 0 KIP-1159.
"validVersions": "0",
"fields": [
{ "name": "fullPayloadPath", "versions": "0+", "type": "string", "about": "The full path to access payload as string."},
{ "name": "payloadStoreClass", "versions": "0+", "type": "string","about": "The used payload store class path to upload payload."}
]
}
7. 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.
*/
public class PayloadResponse {
public final int responseCode;
public final String path;
public final PayloadStoreException payloadStoreException;
public final Boolean retriable;
/**
* Construct payload response with response code and payload id.
*/
public PayloadResponse(int responseCode, String path) {
this(responseCode, path, null, false);
}
/**
* Construct payload response with response code, payload id and exception.
*/
public PayloadResponse(int responseCode, String path, PayloadStoreException payloadStoreException, Boolean retriable) {
this.responseCode = responseCode;
this.fullPayloadPath = path;
this.payloadStoreException = payloadStoreException;
this.retriable = retriable;
}
}
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);
Compatibility, Deprecation, and Migration Plan
- Old clients just need to set the needed configuration to use this feature
Rejected Alternatives
- We have rejected implementing the chunking pattern due to its many potential edge cases and complexity added to the consumer side. A peak of those complexities can be explored more deeply in the LinkedIn presentation.
- Implement this as a separate project outside of Apache Kafka as this seems to be a pattern that needs more use cases and it would be better to have this in Apache Kafka as native implementation instead of a separate project