...
1. Composable Serializer/Deserializer
org.apache.kafka.common.serilalization.ComposableSerializer and org.apache.kafka.common.serilalization.ComposableDeserializer
Configuration
| Key | Description | Type | 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[]>`. | String | None | No, If not exist the code will fail back to value.serializer |
| 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[]>`. | Strings | None | No, If not exist the code will fail back to value.deserializer |
2. org.apache.kafka.common.serilalization.large.message.LargeMessageSerializer
Configuration
| Key | Description | Type | Default | Required |
|---|
large.message.payload.store.class | Implementation of org.apache.kafka.common.serialization.largemessage.store.PayloadStore . | String | None | Yes |
large.message.payload.store.timeout.ms | Timeout for payload store operations. This can't exceedmax.block.msfor Kafka producers ormax.poll.interval.msin Kafka consumer. This is a part of basic configurations for any implementation of org.apache.kafka.common.serialization.largemessage.store.PayloadStore | Long | 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 | Int | 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 | Long | 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 | Long | 100 | No |
large.message.threshold.bytes | The maximum size of the message is considered large. | Long | 1MB (Default of max.message.bytes) | No |
|
|
|
|
|
3.org.apache.kafka.common.serilalization.large.message.LargeMessageDeserializer
Configuration
| Key | Description | Type | Default | Required |
|---|
large.message.payload.store.class | Implementation of org.apache.kafka.common.serialization.largemessage.store.PayloadStore . | String | None | Yes |
large.message.payload.store.timeout.ms | Timeout for payload store operations. This can't exceedmax.block.msfor Kafka producers ormax.poll.interval.msin Kafka consumer. This is a part of basic configurations for any implementation of org.apache.kafka.common.serialization.largemessage.store.PayloadStore | Long | 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 | Int | 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 | Long | 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 | Long | 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. | Boolean | FALSE | No |
|
|
|
|
|
4.org.apache.kafka.common.serilalization.large.message.PayloadStore
| Code Block |
|---|
|
/**
* 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
| Code Block |
|---|
|
/**
* 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
| Code Block |
|---|
/**
* 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;
}
}
|
...