Kafka
The Kafka module provides declarative integration with Apache Kafka: reading messages through
@KafkaListener, sending messages through @KafkaPublisher, serialization, deserialization, transactions, processing errors,
and telemetry.
Apache Kafka is a distributed event streaming platform. Applications write events to a topic, while other applications read
them through a consumer group or directly assigned partitions. Kora creates the required Consumer and Producer at compile time,
binds them to the dependency graph, and lets most of the contract be described through method signatures.
For a step-by-step walkthrough before the reference details, see Kafka Messaging.
Dependency¶
Dependency build.gradle:
Module:
Dependency build.gradle.kts:
Module:
Consumer¶
Consumer reads records from a topic and passes them to an application method. Kora creates the consumer container,
calls poll(), applies deserialization, invokes the handler, and commits the offset unless the method signature requires
manual Consumer control.
Creating a Consumer requires using the @KafkaListener annotation over a method:
The @KafkaListener annotation parameter points to the Consumer configuration path.
In case you need different behavior for different topics, it is possible to create several such containers, each with its own individual configuration. It looks like this:
The value in the annotation indicates which part of the configuration file should be used.
Conceptually, it is similar to @ConfigSource: the annotation value selects the configuration branch for a specific container.
Configuration¶
Configuration describes the settings of a particular @KafkaListener and an example for the configuration at path kafka.someConsumer is given below.
Basic configuration parameters:
kafka {
someConsumer {
topics = ["topic1", "topic2"] //(1)!
offset = "latest" //(2)!
pollTimeout = "5s" //(3)!
threads = 1 //(4)!
driverProperties { //(5)!
"bootstrap.servers": "localhost:9093"
"group.id": "my-group-id"
}
}
}
- List of
topics to subscribe to (requiredto specify eithertopicsortopicsPattern) - Initial read position (default:
latest). Allowed values:earliest,latest, or time offset (e.g.5m) - Maximum time to wait for messages (default:
5s) - Number of threads for the consumer (default:
1) - Official
Kafka ConsumerProperties(required, no default)
kafka:
someConsumer:
topics:
- "topic1"
- "topic2" #(1)!
offset: "latest" #(2)!
pollTimeout: "5s" #(3)!
threads: 1 #(4)!
driverProperties: #(5)!
"bootstrap.servers": "localhost:9093"
"group.id": "my-group-id"
- List of
topics to subscribe to (requiredto specify eithertopicsortopicsPattern) - Initial read position (default:
latest). Allowed values:earliest,latest, or time offset (e.g.5m) - Maximum time to wait for messages (default:
5s) - Number of threads for the consumer (default:
1) - Official
Kafka ConsumerProperties(required, no default)
Full Configuration
Example of the complete configuration described in the KafkaListenerConfig class (default or example values are specified):
In a real configuration, either topics or topicsPattern is usually specified.
kafka {
someConsumer {
topics = ["topic1", "topic2"] //(1)!
topicsPattern = "topic*" //(2)!
partitions = ["0", "1"] //(3)!
allowEmptyRecords = false //(4)!
offset = "latest" //(5)!
pollTimeout = "5s" //(6)!
backoffTimeout = "15s" //(7)!
partitionRefreshInterval = "1m" //(8)!
threads = 1 //(9)!
shutdownWait = "30s" //(10)!
driverProperties { //(11)!
"bootstrap.servers": "localhost:9093"
"group.id": "my-group-id"
}
telemetry {
logging {
enabled = false //(12)!
}
metrics {
enabled = true //(13)!
slo = [1, 10, 50, 100, 200, 500, 1000, 2000, 5000, 10000, 20000, 30000, 60000, 90000] //(14)!
tags = { // (15)!
"key1" = "value1"
"key2" = "value2"
}
}
tracing {
enabled = true //(16)!
attributes = { // (17)!
"key1" = "value1"
"key2" = "value2"
}
}
}
}
}
- List of
topicvalues theConsumersubscribes to (not set by default, optional; eithertopicsortopicsPatternmust be specified) topicpattern theConsumersubscribes to (not set by default, optional; eithertopicsortopicsPatternmust be specified)- List of partitions used only for consumer name construction when
group.id,topics, andtopicsPatternare not specified; partition assignment is controlled by theassigncontainer (not set by default, optional) If specified, consumer will read only from specified partitions. Example:["0", "1", "2"]. - Whether to process empty batches when the signature accepts
ConsumerRecords(default:false) IffalseandConsumerRecordsis empty (no messages), consumer method will not be called. Iftrue, method will be called with emptyConsumerRecords(useful for periodic checks). - Initial read position for the
assignstrategy whengroup.idis not specified (default:latest). Valid values:earliest- earliest availableoffsetlatest- latest availableoffset- string in
Durationformat, for example5m, - shift back by the specified duration Format: number + unit (ms, s, m, h, d). Examples:5m= 5 minutes ago,1h= 1 hour ago.
- Maximum time to wait for messages from a
topicwithin onepoll()call (default:5s) - Initial delay between unexpected processing errors; with repeated errors the delay increases up to
60s(default:15s) If consumer throws unexpected exception (notKafkaSkipRecordException), Kora will restart consumer withbackoffTimeoutdelay to prevent cyclic errors. - Partition list refresh period for the
assignstrategy (default:1m) - Number of threads the consumer starts on; if set to
0, the consumer is not started (default:1) - Time to wait for processing before stopping the consumer during graceful shutdown (default:
30s) - Official
Kafka ConsumerProperties; see Apache Kafka Consumer Configs (required, not set by default) - Enables module logging (default:
false) - Enables module metrics (default:
true) - Configures SLO for metrics (default:
ru.tinkoff.kora.telemetry.common.TelemetryConfig.MetricsConfig#DEFAULT_SLO) - Configures metric tags (default:
{}) - Enables module tracing (default:
true) - Configures tracing attributes (default:
{})
kafka:
someConsumer:
topics: #(1)!
- "topic1"
- "topic2"
topicsPattern: "topic*" #(2)!
partitions: #(3)!
- "0"
- "1"
allowEmptyRecords: false #(4)!
offset: "latest" #(5)!
pollTimeout: "5s" #(6)!
backoffTimeout: "15s" #(7)!
partitionRefreshInterval: "1m" #(8)!
threads: 1 #(9)!
shutdownWait: "30s" #(10)!
driverProperties: #(11)!
bootstrap.servers: "localhost:9093"
group.id: "my-group-id"
telemetry:
logging:
enabled: false #(12)!
metrics:
enabled: true #(13)!
slo: [ 1, 10, 50, 100, 200, 500, 1000, 2000, 5000, 10000, 20000, 30000, 60000, 90000 ] #(14)!
tags: #(15)!
key1: value1
key2: value2
tracing:
enabled: true #(16)!
attributes: #(17)!
key1: value1
key2: value2
- List of
topicvalues theConsumersubscribes to (not set by default, optional; eithertopicsortopicsPatternmust be specified) topicpattern theConsumersubscribes to (not set by default, optional; eithertopicsortopicsPatternmust be specified)- List of partitions used only for consumer name construction when
group.id,topics, andtopicsPatternare not specified; partition assignment is controlled by theassigncontainer (not set by default, optional) If specified, consumer will read only from specified partitions. Example:["0", "1", "2"]. - Whether to process empty batches when the signature accepts
ConsumerRecords(default:false) IffalseandConsumerRecordsis empty (no messages), consumer method will not be called. Iftrue, method will be called with emptyConsumerRecords(useful for periodic checks). - Initial read position for the
assignstrategy whengroup.idis not specified (default:latest). Valid values:earliest- earliest availableoffsetlatest- latest availableoffset- string in
Durationformat, for example5m, - shift back by the specified duration Format: number + unit (ms, s, m, h, d). Examples:5m= 5 minutes ago,1h= 1 hour ago.
- Maximum time to wait for messages from a
topicwithin onepoll()call (default:5s) - Initial delay between unexpected processing errors; with repeated errors the delay increases up to
60s(default:15s) If consumer throws unexpected exception (notKafkaSkipRecordException), Kora will restart consumer withbackoffTimeoutdelay to prevent cyclic errors. - Partition list refresh period for the
assignstrategy (default:1m) - Number of threads the consumer starts on; if set to
0, the consumer is not started (default:1) - Time to wait for processing before stopping the consumer during graceful shutdown (default:
30s) - Official
Kafka ConsumerProperties; see Apache Kafka Consumer Configs (required, not set by default) - Enables module logging (default:
false) - Enables module metrics (default:
true) - Configures SLO for metrics (default:
ru.tinkoff.kora.telemetry.common.TelemetryConfig.MetricsConfig#DEFAULT_SLO) - Configures metric tags (default:
{}) - Enables module tracing (default:
true) - Configures tracing attributes (default:
{})
Module metrics are described in the Metrics Reference section.
Consume strategy¶
subscribe strategy involves the use of group.id,
to group the executors so that they do not duplicate the reading of records from their queue across multiple application instances.
Example of subscribe strategy configuration:
assign connection strategy implies that each instance of the application reads messages from the topic simultaneously with others,
i.e., messages are duplicated between all instances of the application within the topic.
This strategy is useful, for example, when all application replicas must receive the same message at once: to reset a local cache,
update local reference data, or handle a service event.
To use this strategy, simply do not specify group.id in the consumer configuration.
However, only one topic can be specified at a time in this strategy.
Example of assign strategy configuration:
Signatures¶
Available signatures for out-of-the-box Kafka Consumer methods, where K refers to the key type and V to the message value type.
The generator supports three signature families: separate key/value arguments, a single ConsumerRecord<K, V>, or a whole ConsumerRecords<K, V> batch.
These families cannot be mixed in the same method.
Key and value¶
A signature with separate arguments accepts value, optional key, optional Headers, optional Consumer<K, V>, and optional deserialization errors.
One user argument is treated as value; two user arguments are treated as key and value in that exact order.
If key is not declared, the key deserialization type is considered to be byte[].
To handle deserialization errors, add Exception, RecordKeyDeserializationException, or RecordValueDeserializationException.
When such an argument is present, Kora passes the deserialization error to it, and the corresponding key or value is passed as null.
Without such an argument, the deserialization error is thrown from the handler, and the record is read again without committing the current offset.
Whole record¶
A signature with ConsumerRecord<K, V> accepts one whole record, optional Consumer<K, V>, and optional deserialization errors:
Exception, RecordKeyDeserializationException, or RecordValueDeserializationException.
Headers, separate key/value arguments, and the telemetry context are not supported in this signature.
If error arguments are not declared, the deserialization error can be thrown when calling record.key() or record.value().
If error arguments are declared, Kora calls key() and/or value() beforehand, catches the deserialization error, and passes it to the method.
@KafkaListener("kafka.someConsumer")
void process(ConsumerRecord<K, V> record) {
try {
var key = record.key();
var value = record.value();
// some value handling work
} catch (RecordKeyDeserializationException e) {
// do deserialization handling work
} catch (RecordValueDeserializationException e) {
// do deserialization handling work
}
}
@KafkaListener("kafka.someConsumer")
fun process(record: ConsumerRecord<K, V>) {
try {
val key = record.key()
val value = record.value()
// some value handling work
} catch (e: RecordKeyDeserializationException) {
// do deserialization handling work
} catch (e: RecordValueDeserializationException) {
// do deserialization handling work
}
}
@KafkaListener("kafka.someConsumer")
void process(ConsumerRecord<K, V> record,
@Nullable RecordKeyDeserializationException keyException,
@Nullable RecordValueDeserializationException valueException) {
if (keyException != null || valueException != null) {
// do deserialization handling work
return;
}
var key = record.key();
var value = record.value();
// some value handling work
}
@KafkaListener("kafka.someConsumer")
fun process(
record: ConsumerRecord<K, V>,
keyException: RecordKeyDeserializationException?,
valueException: RecordValueDeserializationException?,
) {
if (keyException != null || valueException != null) {
// do deserialization handling work
return
}
val key = record.key()
val value = record.value()
// some value handling work
}
Batch of records¶
A signature with ConsumerRecords<K, V> accepts the whole batch of records from one poll().
Together with it, only Consumer<K, V> and KafkaConsumerRecordsTelemetryContext<K, V> can be declared.
Separate key/value arguments, Headers, and deserialization error arguments are not supported in this signature; deserialization errors should be handled while iterating over records.
@KafkaListener("kafka.someConsumer")
void process(ConsumerRecords<K, V> records) {
for (var record : records) {
try {
var key = record.key();
var value = record.value();
// some value handling work
} catch (RecordKeyDeserializationException e) {
// do deserialization handling work
} catch (RecordValueDeserializationException e) {
// do deserialization handling work
}
}
}
@KafkaListener("kafka.someConsumer")
fun process(records: ConsumerRecords<K, V>) {
for (record in records) {
try {
val key = record.key()
val value = record.value()
// some value handling work
} catch (e: RecordKeyDeserializationException) {
// do deserialization handling work
} catch (e: RecordValueDeserializationException) {
// do deserialization handling work
}
}
}
Offset commit¶
If the signature does not declare a Consumer<K, V> argument, Kora commits the offset automatically: after each record for key/value and ConsumerRecord<K, V> signatures, or after the whole batch for ConsumerRecords<K, V>.
It does this by calling commitSync().
If the signature declares a Consumer<K, V> argument, automatic offset commit is disabled, and the handler is fully responsible for calling commitSync() or commitAsync().
This mode is useful when the offset should be committed only after an external operation, several records should be committed together, or the read position should be controlled manually.
In subscribe mode, a manual commit commits the offset inside the consumer group.
In assign mode, partitions are not coordinated through a consumer group, so it is usually more important to manage the position manually with seek(), pause(), and resume() instead of relying on a group offset commit.
If the handler fails before the manual commit, the record or batch will be read again according to the current consumer position.
@KafkaListener("kafka.someConsumer")
void process(ConsumerRecord<K, V> record, Consumer<K, V> consumer) {
try {
var key = record.key();
var value = record.value();
// some value handling work
} catch (RecordKeyDeserializationException e) {
// do deserialization handling work
} catch (RecordValueDeserializationException e) {
// do deserialization handling work
} finally {
consumer.commitSync();
}
}
@KafkaListener("kafka.someConsumer")
fun process(record: ConsumerRecord<K, V>, consumer: Consumer<K, V>) {
try {
val key = record.key()
val value = record.value()
// some value handling work
} catch (e: RecordKeyDeserializationException) {
// do deserialization handling work
} catch (e: RecordValueDeserializationException) {
// do deserialization handling work
} finally {
consumer.commitSync()
}
}
Deserialization¶
Deserializer is used to deserialize ConsumerRecord keys and values.
Kora provides Deserializer components for basic types: String, UUID, byte[], Bytes, ByteBuffer, Double, Float,
Integer, Long, Short, and Void.
Tags are supported to better customize the Deserializer.
Tags can be set on parameter-key, parameter-value, as well as on parameters of type ConsumerRecord and ConsumerRecords.
These tags will be set on container dependencies.
@Component
final class ConsumerService {
@KafkaListener("kafka.someConsumer1")
void process1(@Tag(Sometag1.class) String key, @Tag(Sometag2.class) String value) {
// some handler code
}
@KafkaListener("kafka.someConsumer2")
void process2(ConsumerRecord<@Tag(Sometag1.class) String, @Tag(Sometag2.class) String> record) {
// some handler code
}
}
@Component
class ConsumerService {
@KafkaListener("kafka.someConsumer1")
fun process1(@Tag(Sometag1::class) key: String, @Tag(Sometag2::class) value: String) {
// some handler code
}
@KafkaListener("kafka.someConsumer2")
fun process2(record: ConsumerRecord<@Tag(Sometag1::class) String, @Tag(Sometag2::class) String>) {
// some handler code
}
}
If deserialization from JSON is required, use the @Json tag.
In this case, Kora uses JsonReader<T> and JsonKafkaDeserializer<T> from the JSON module:
@Component
final class ConsumerService {
@Json
public record JsonEvent(String name, Integer code) {}
@KafkaListener("kafka.someConsumer1")
void process1(String key, @Json JsonEvent value) {
// some handler code
}
@KafkaListener("kafka.someConsumer2")
void process2(ConsumerRecord<String, @Json JsonEvent> record) {
// some handler code
}
}
@Component
class ConsumerService {
@Json
data class JsonEvent(val name: String, val code: Int)
@KafkaListener("kafka.someConsumer1")
fun process1(key: String, @Json value: JsonEvent) {
// some handler code
}
@KafkaListener("kafka.someConsumer2")
fun process2(record: ConsumerRecord<String, @Json JsonEvent>) {
// some handler code
}
}
For non-key handlers, the default is Deserializer<byte[]> since it simply returns unhandled bytes.
Custom Deserializer¶
If custom deserialization is required, you can implement your own Deserializer.
Option 1: Default deserializer for type
If you provide Deserializer<T> as a component without a tag, it will be used for all consumers of this type:
@Component
public static class MyEventDeserializer implements Deserializer<MyEvent> {
private final JsonReader<MyEvent> reader;
public MyEventDeserializer(JsonReader<MyEvent> reader) {
this.reader = reader;
}
@Override
public MyEvent deserialize(String topic, byte[] data) {
try {
return reader.read(data);
} catch (IOException e) {
throw new IllegalArgumentException(e);
}
}
}
@Component
final class SomeConsumer {
@KafkaListener("kafka.someConsumer")
void process(MyEvent value) { // Uses MyEventDeserializer
// event handling
}
}
@Component
class MyEventDeserializer(
private val reader: JsonReader<MyEvent>
) : Deserializer<MyEvent> {
override fun deserialize(topic: String, data: ByteArray): MyEvent {
return try {
reader.read(data)
} catch (e: IOException) {
throw IllegalArgumentException(e)
}
}
}
@Component
class SomeConsumer {
@KafkaListener("kafka.someConsumer")
fun process(value: MyEvent) { // Uses MyEventDeserializer
// event handling
}
}
Option 2: Point deserializer for specific consumer
If you need to use different deserialization for different consumers of the same type, you can use tags:
@Component
final class SomeConsumer {
@Json
public record MyEvent(String username, int code) {}
@Tag(MyEvent.class)
@Component
public static class MyDeserializer implements Deserializer<MyEvent> {
private final JsonReader<MyEvent> reader;
public MyDeserializer(JsonReader<MyEvent> reader) {
this.reader = reader;
}
@Override
public MyEvent deserialize(String topic, byte[] data) {
try {
return reader.read(data);
} catch (IOException e) {
throw new IllegalArgumentException(e);
}
}
}
@KafkaListener("kafka.someConsumer")
void process(@Tag(MyEvent.class) MyEvent value) {
// event handling
}
}
@Component
class SomeConsumer {
@Json
data class MyEvent(val username: String, val code: Int)
@Tag(MyEvent::class)
@Component
class MyDeserializer(
private val reader: JsonReader<MyEvent>
) : Deserializer<MyEvent> {
override fun deserialize(topic: String, data: ByteArray): MyEvent {
return try {
reader.read(data)
} catch (e: IOException) {
throw IllegalArgumentException(e)
}
}
}
@KafkaListener("kafka.someConsumer")
fun process(@Tag(MyEvent::class) value: MyEvent) {
// event handling
}
}
Exception handling¶
If the method labeled @KafkaListener throws an exception, Consumer will be restarted,
because there is no general solution on how to handle this and the developer must decide how to handle it.
Exception skipping¶
If you need to skip processing a specific event (ConsumerRecord) during processing for business logic reasons,
you can throw a KafkaSkipRecordException by passing the actual exception to the constructor.
In this case, all metrics will be correctly accounted for and recorded, processing of the corresponding event will be skipped, and the next event will begin to be processed.
If you want to implement your own skippable exceptions,
you can use the SkippableRecordException interface, which should be implemented in your exceptions.
Deserialization errors¶
If you use a signature with ConsumerRecord or ConsumerRecords, you will get a value deserialization exception at the moment of calling the key or value methods.
At that point, it is worth handling it in the way you want.
The following exceptions are thrown:
ru.tinkoff.kora.kafka.common.exceptions.RecordKeyDeserializationException.ru.tinkoff.kora.kafka.common.exceptions.RecordValueDeserializationException.
From these exceptions, you can get a raw ConsumerRecord<byte[], byte[]> using getRecord() method:
@Component
final class ConsumerService {
@KafkaListener("kafka.someConsumer")
public void process(String key, String value) {
try {
// Access key/value which may throw deserialization exception
} catch (RecordKeyDeserializationException e) {
ConsumerRecord<byte[], byte[]> rawRecord = e.getRecord();
// Handle raw record (log, send to DLQ, etc.)
} catch (RecordValueDeserializationException e) {
ConsumerRecord<byte[], byte[]> rawRecord = e.getRecord();
// Handle raw record (log, send to DLQ, etc.)
}
}
}
@Component
class ConsumerService {
@KafkaListener("kafka.someConsumer")
fun process(key: String, value: String) {
try {
// Access key/value which may throw deserialization exception
} catch (e: RecordKeyDeserializationException) {
val rawRecord = e.getRecord()
// Handle raw record (log, send to DLQ, etc.)
} catch (e: RecordValueDeserializationException) {
val rawRecord = e.getRecord()
// Handle raw record (log, send to DLQ, etc.)
}
}
}
If you use a signature with unpacked key/value/headers, you can add Exception, Throwable, RecordKeyDeserializationException, RecordKeyDeserializationException with the last argument
or RecordValueDeserializationException.
Note that all arguments become optional, meaning we expect to either have a key and value or an exception.
Custom tag¶
Automatic tag is created for the consumer by default, it can be viewed in the generated module at compile time.
If for some reason you need to override the consumer tag, you can set it as an argument to the @KafkaListener annotation:
Rebalance events¶
You can listen and react to rebalance events with your implementation of the ConsumerAwareRebalanceListener interface,
it should be provided as a component by the consumer tag:
@Tag(SomeListenerProcessTag.class)
@Component
public final class SomeListener implements ConsumerAwareRebalanceListener {
@Override
public void onPartitionsRevoked(Consumer<?, ?> consumer, Collection<TopicPartition> partitions) {
// Called before partitions are revoked from this consumer.
// Use this to commit offsets or cleanup state.
}
@Override
public void onPartitionsAssigned(Consumer<?, ?> consumer, Collection<TopicPartition> partitions) {
// Called when partitions are assigned to this consumer.
// Use this to initialize state for assigned partitions.
}
@Override
public void onPartitionsLost(Consumer<?, ?> consumer, Collection<TopicPartition> partitions) {
// Called when partitions are lost (e.g., consumer failure, group rebalance).
// Unlike onPartitionsRevoked, this is called when the consumer is no longer
// part of the group and cannot commit offsets.
// Use this to cleanup local state for lost partitions.
}
}
@Tag(SomeListenerProcessTag::class)
@Component
class SomeListener : ConsumerAwareRebalanceListener {
override fun onPartitionsRevoked(consumer: Consumer<*, *>, partitions: Collection<TopicPartition>) {
// Called before partitions are revoked from this consumer.
// Use this to commit offsets or cleanup state.
}
override fun onPartitionsAssigned(consumer: Consumer<*, *>, partitions: Collection<TopicPartition>) {
// Called when partitions are assigned to this consumer.
// Use this to initialize state for assigned partitions.
}
override fun onPartitionsLost(consumer: Consumer<*, *>, partitions: Collection<TopicPartition>) {
// Called when partitions are lost (e.g., consumer failure, group rebalance).
// Unlike onPartitionsRevoked, this is called when the consumer is no longer
// part of the group and cannot commit offsets.
// Use this to cleanup local state for lost partitions.
}
}
Manual override¶
Kora provides a small wrapper over KafkaConsumer that allows you to easily trigger the handling of incoming events.
The container constructor is as follows:
public KafkaSubscribeConsumerContainer(KafkaListenerConfig config,
Deserializer<K> keyDeserializer,
Deserializer<V> valueDeserializer,
BaseKafkaRecordsHandler<K, V> handler)
BaseKafkaRecordsHandler<K,V> is the basic functional interface of the handler:
@FunctionalInterface
public interface BaseKafkaRecordsHandler<K, V> {
void handle(ConsumerRecords<K, V> records, KafkaConsumer<K, V> consumer);
}
Telemetry¶
Kafka uses a telemetry contract for logging, metrics, and tracing of messages.
Telemetry configuration (the telemetry { logging / metrics / tracing } section) is described in the Configuration section.
For each batch & message, KafkaListener creates a separate telemetry context, which is closed upon completion of processing.
The message is described via telemetry handler parameters, including topic, partition, offset, and processing duration.
The default factory, DefaultKafkaListenerTelemetryFactory, combines three factories:
- KafkaListenerLoggerFactory builds a KafkaListenerLogger to log the start and end of message processing;
- KafkaListenerMetricsFactory builds a KafkaListenerMetrics to record message metrics;
- KafkaListenerTracerFactory builds a KafkaListenerTracer for distributed tracing.
Metrics and tracing are described in the Metrics Reference section.
Producer¶
Producer sends records to a topic. Kora creates an implementation of the interface annotated with @KafkaPublisher,
selects a Serializer for the key and value, calls KafkaProducer#send, and connects sending with telemetry.
To create a Producer, use the @KafkaPublisher annotation on an interface.
To send messages to an arbitrary topic, declare a method with a ProducerRecord parameter:
The annotation parameter indicates the path to the producer configuration.
Topic¶
If typed methods are required for specific topic values, use the @KafkaPublisher.Topic annotation:
The annotation parameter indicates the path for the topic configuration.
Configuration¶
Configuration describes the settings of a particular @KafkaPublisher; below is an example for the kafka.someProducer configuration path.
Basic configuration parameters:
Full Configuration
Example of the complete configuration described in the KafkaPublisherConfig class (default or example values are specified):
kafka {
someProducer {
driverProperties { //(1)!
"bootstrap.servers": "localhost:9093"
}
telemetry {
logging {
enabled = false //(2)!
}
metrics {
enabled = true //(3)!
slo = [ 1, 10, 50, 100, 200, 500, 1000, 2000, 5000, 10000, 20000, 30000, 60000, 90000 ] //(4)!
tags = { //(5)!
"key1" = "value1"
"key2" = "value2"
}
}
tracing {
enabled = true //(6)!
attributes = { //(7)!
"key1" = "value1"
"key2" = "value2"
}
}
}
}
}
- Official
Kafka ProducerProperties; see Apache Kafka Producer Configs (required, not set by default) - Enables module logging (default:
false) - Enables module metrics (default:
true) - Configures SLO for metrics (default:
ru.tinkoff.kora.telemetry.common.TelemetryConfig.MetricsConfig#DEFAULT_SLO) - Configures metric tags (default:
{}) - Enables module tracing (default:
true) - Configures tracing attributes (default:
{})
kafka:
someProducer:
driverProperties: #(1)!
bootstrap.servers: "localhost:9093"
telemetry:
logging:
enabled: true #(2)!
metrics:
enabled: true #(3)!
slo: [ 1, 10, 50, 100, 200, 500, 1000, 2000, 5000, 10000, 20000, 30000, 60000, 90000 ] #(4)!
tags: #(5)!
key1: value1
key2: value2
tracing:
enabled: true #(6)!
attributes: #(7)!
key1: value1
key2: value2
- Official
Kafka ProducerProperties; see Apache Kafka Producer Configs (required, not set by default) - Enables module logging (default:
false) - Enables module metrics (default:
true) - Configures SLO for metrics (default:
ru.tinkoff.kora.telemetry.common.TelemetryConfig.MetricsConfig#DEFAULT_SLO) - Configures metric tags (default:
{}) - Enables module tracing (default:
true) - Configures tracing attributes (default:
{})
topic configuration describes the settings of a particular @KafkaPublisher.Topic; below is an example for the kafka.someProducer.someTopic configuration path.
Example of the complete configuration described in the KafkaPublisherConfig.TopicConfig class (default or example values are specified):
topicwhere the method sends data (required, not set by default)topicpartition where the method sends data (not set by default, optional) If specified, all messages will be sent to the specified partition. If not specified, standard Kafka partitioning is used (by key or random).
topicwhere the method sends data (required, not set by default)topicpartition where the method sends data (not set by default, optional) If specified, all messages will be sent to the specified partition. If not specified, standard Kafka partitioning is used (by key or random).
Signatures¶
Allows value and key (optional) and headers (optional) to be sent from ProducerRecord:
Can be received as the result of a RecordMetadata operation:
Can be also obtained as the result of a Future<RecordMetadata> or CompletionStage<RecordMetadata> operation:
@KafkaPublisher("kafka.someProducer")
public interface MyPublisher {
@KafkaPublisher.Topic("kafka.someProducer.someTopic")
RecordMetadata send(String value);
@KafkaPublisher.Topic("kafka.someProducer.someTopic")
Future<RecordMetadata> send(String value);
@KafkaPublisher.Topic("kafka.someProducer.someTopic")
CompletionStage<RecordMetadata> sendStage(V value);
}
Can be also obtained as the result of a suspend or Future<RecordMetadata> or CompletionStage<RecordMetadata> or Deferred<RecordMetadata> operation:
@KafkaPublisher("kafka.someProducer")
interface MyPublisher {
@KafkaPublisher.Topic("kafka.someProducer.someTopic")
fun send(value: String): RecordMetadata
@KafkaPublisher.Topic("kafka.someProducer.someTopic")
suspend fun sendSuspend(value: V): RecordMetadata
@KafkaPublisher.Topic("kafka.someProducer.someTopic")
fun send(value: String): Future<RecordMetadata>
@KafkaPublisher.Topic("kafka.someProducer.someTopic")
fun send(value: String): CompletionStage<RecordMetadata>
@KafkaPublisher.Topic("kafka.someProducer.someTopic")
fun send(value: String): Deferred<RecordMetadata>
}
It is possible to send ProducerRecord with or without Callback and combine the response with the examples above:
Serialization¶
Serializer is used to serialize ProducerRecord keys and values.
Kora provides Serializer components for basic types: String, UUID, byte[], Bytes, ByteBuffer, Double, Float,
Integer, Long, Short, and Void.
To specify which Serializer to take from the container, tags can be used.
Tags should be set on ProducerRecord or key/value parameters of methods:
If serialization to JSON is required, use the @Json tag.
In this case, Kora uses JsonWriter<T> and JsonKafkaSerializer<T> from the JSON module:
Custom Serializer¶
If custom serialization is required, you can implement your own Serializer.
Option 1: Default serializer for type
If you provide Serializer<T> as a component without a tag, it will be used for all producers of this type:
@Component
public static class MyEventSerializer implements Serializer<MyEvent> {
private final JsonWriter<MyEvent> writer;
public MyEventSerializer(JsonWriter<MyEvent> writer) {
this.writer = writer;
}
@Override
public byte[] serialize(String topic, MyEvent data) {
try {
return writer.toByteArray(data);
} catch (IOException e) {
throw new IllegalArgumentException(e);
}
}
}
@KafkaPublisher("kafka.someProducer")
public interface MyPublisher {
@KafkaPublisher.Topic("kafka.someProducer.topic")
void send(MyEvent value); // Uses MyEventSerializer
}
@Component
class MyEventSerializer(
private val writer: JsonWriter<MyEvent>
) : Serializer<MyEvent> {
override fun serialize(topic: String, data: MyEvent): ByteArray {
return try {
writer.toByteArray(data)
} catch (e: IOException) {
throw IllegalArgumentException(e)
}
}
}
@KafkaPublisher("kafka.someProducer")
interface MyPublisher {
@KafkaPublisher.Topic("kafka.someProducer.topic")
fun send(value: MyEvent) // Uses MyEventSerializer
}
Option 2: Point serializer for specific producer
If you need to use different serialization for different producers of the same type, you can use tags:
@KafkaPublisher("kafka.someProducer")
public interface MyKafkaProducer {
@Json
record MyEvent(String username, int code) {}
@Tag(MyEvent.class)
@Component
class MySerializer implements Serializer<MyEvent> {
private final JsonWriter<MyEvent> writer;
public MySerializer(JsonWriter<MyEvent> writer) {
this.writer = writer;
}
@Override
public byte[] serialize(String topic, MyEvent data) {
try {
return writer.toByteArray(data);
} catch (IOException e) {
throw new IllegalArgumentException(e);
}
}
}
void send(ProducerRecord<String, @Tag(MyEvent.class) MyEvent> record);
}
@KafkaPublisher("kafka.someProducer")
interface MyKafkaProducer {
@Json
data class MyEvent(val username: String, val code: Int)
@Tag(MyEvent::class)
@Component
class MySerializer(
private val writer: JsonWriter<MyEvent>
) : Serializer<MyEvent> {
override fun serialize(topic: String, data: MyEvent): ByteArray {
return try {
writer.toByteArray(data)
} catch (e: IOException) {
throw IllegalArgumentException(e)
}
}
}
fun send(record: ProducerRecord<String, @Tag(MyEvent::class) MyEvent>)
}
Exception handling¶
If a send error happens in a method annotated with @KafkaPublisher.Topic that does not return Future<RecordMetadata>,
ru.tinkoff.kora.kafka.common.exceptions.KafkaPublishException is thrown.
The original error from KafkaProducer is available in cause.
To get the failed record, use the getRecord() method:
@Component
class SomeService {
private final MyPublisher publisher;
public SomeService(MyPublisher publisher) {
this.publisher = publisher;
}
void sendMessage() {
try {
publisher.send("key", "value");
} catch (KafkaPublishException e) {
ProducerRecord<byte[], byte[]> failedRecord = e.getRecord();
// Handle the failed record (log, retry, etc.)
}
}
}
Serialization errors¶
If a key or value serialization error happens in a method annotated with @KafkaPublisher.Topic,
org.apache.kafka.common.errors.SerializationException is thrown, just like with a direct org.apache.kafka.clients.producer.Producer#send call.
Transactions¶
Messages can be sent to Kafka within a transaction.
For this, use the @KafkaPublisher annotation and extend TransactionalPublisher.
First, describe a regular KafkaProducer, and then use its type to create a transactional Producer:
@KafkaPublisher("kafka.someProducer")
public interface MyPublisher {
@KafkaPublisher.Topic("kafka.someProducer.someTopic")
void send(String key, String value);
}
@KafkaPublisher("kafka.someTransactionalProducer")
public interface MyTransactionalPublisher extends TransactionalPublisher<MyPublisher> {
}
Use inTx methods to send messages in a transaction: all messages inside the lambda are committed on successful execution
and aborted on error.
It is also possible to manage the transaction manually through begin():
Configuration¶
KafkaPublisherConfig.TransactionConfig is used to configure @KafkaPublisher with the TransactionalPublisher interface:
kafka {
someTransactionalProducer {
idPrefix = "kora-app-" //(1)!
maxPoolSize = 10 //(2)!
maxWaitTime = "10s" //(3)!
}
}
- Transaction identifier prefix; a random
UUIDwill be appended to it (default:kora-app-) Format:{idPrefix}-{uuid}. Example:kafka-app-550e8400-e29b-41d4-a716-446655440000. - Maximum size of the transactional
Producerpool (default:10) - Maximum time to wait for a free
Producerfrom the pool (default:10s)
kafka:
someTransactionalProducer:
idPrefix: "kora-app-" #(1)!
maxPoolSize: 10 #(2)!
maxWaitTime: "10s" #(3)!
- Transaction identifier prefix; a random
UUIDwill be appended to it (default:kora-app-) Format:{idPrefix}-{uuid}. Example:kafka-app-550e8400-e29b-41d4-a716-446655440000. - Maximum size of the transactional
Producerpool (default:10) - Maximum time to wait for a free
Producerfrom the pool (default:10s)
Advanced Transaction Usage¶
Transaction Interface¶
The begin() method returns a Transaction<P> object that provides advanced transaction management capabilities:
try (var tx = transactionalPublisher.begin()) {
// Sending messages
tx.publisher().send("key1", "value1");
tx.publisher().send("key2", "value2");
// Commit consumer offsets within the transaction (exactly-once semantics)
Map<TopicPartition, OffsetAndMetadata> offsets = ...;
ConsumerGroupMetadata groupMetadata = ...;
tx.sendOffsetsToTransaction(offsets, groupMetadata);
// Explicit flush to guarantee sending before commit
tx.flush();
// commit() is called automatically on try-with-resources close
}
transactionalPublisher.begin().use { tx ->
// Sending messages
tx.publisher().send("key1", "value1")
tx.publisher().send("key2", "value2")
// Commit consumer offsets within the transaction (exactly-once semantics)
val offsets: Map<TopicPartition, OffsetAndMetadata> = ...
val groupMetadata: ConsumerGroupMetadata = ...
tx.sendOffsetsToTransaction(offsets, groupMetadata)
// Explicit flush to guarantee sending before commit
tx.flush()
// commit() is called automatically on use close
}
Transaction<P> methods:
| Method | Description |
|---|---|
publisher() |
Returns typed publisher for sending messages |
producer() |
Returns raw Producer<byte[], byte[]> for low-level operations |
sendOffsetsToTransaction(offsets, groupMetadata) |
Commits consumer offsets within the same transaction |
flush() |
Guarantees all messages are sent before commit |
abort() |
Aborts the transaction |
abort(cause) |
Aborts the transaction with specified cause |
close() |
Closes the transaction (commits if no abort) |
Transaction methods¶
TransactionalPublisher provides 4 methods for working with transactions:
| Method | Passes to callback | Returns value |
|---|---|---|
inTx(TransactionalConsumer) |
P publisher |
void |
inTx(TransactionalFunction) |
P publisher |
R |
withTx(TransactionConsumer) |
Transaction<P> tx |
void |
withTx(TransactionFunction) |
Transaction<P> tx |
R |
Example with return value:
// inTx with return value
Long messageId = transactionalPublisher.inTx(producer -> {
producer.send("key", "value");
return System.currentTimeMillis();
});
// withTx with Transaction access
transactionalPublisher.withTx(tx -> {
tx.publisher().send("key", "value");
tx.sendOffsetsToTransaction(offsets, groupMetadata);
tx.flush(); // Explicit flush
});
// inTx with return value
val messageId = transactionalPublisher.inTx { producer ->
producer.send("key", "value")
System.currentTimeMillis()
}
// withTx with Transaction access
transactionalPublisher.withTx { tx ->
tx.publisher().send("key", "value")
tx.sendOffsetsToTransaction(offsets, groupMetadata)
tx.flush() // Explicit flush
}
Default Serializers and Deserializers¶
KafkaModule automatically provides serializers and deserializers for base types via KafkaSerializersModule and KafkaDeserializersModule.
These serializers/deserializers are provided as components without tags and are used by default for all consumers/producers of corresponding types.
Supported types out of the box:
| Type | Serializer | Deserializer |
|---|---|---|
String |
StringSerializer |
StringDeserializer |
byte[] |
ByteArraySerializer |
ByteArrayDeserializer |
ByteBuffer |
ByteBufferSerializer |
ByteBufferDeserializer |
Bytes |
BytesSerializer |
BytesDeserializer |
UUID |
UUIDSerializer |
UUIDDeserializer |
Integer |
IntegerSerializer |
IntegerDeserializer |
Long |
LongSerializer |
LongDeserializer |
Short |
ShortSerializer |
ShortDeserializer |
Double |
DoubleSerializer |
DoubleDeserializer |
Float |
FloatSerializer |
FloatDeserializer |
Void |
VoidSerializer |
VoidDeserializer |
To use, simply specify the type in the publisher/consumer method:
Signatures¶
Available signatures for out-of-the-box Kafka Producer methods, where K refers to the key type and V to the message value type.
The generator supports two signature families: sending a ready ProducerRecord<K, V> and sending through a method annotated with @KafkaPublisher.Topic.
These families cannot be mixed in the same method.
Prepared event¶
A method with ProducerRecord<K, V> is used when the topic, partition, timestamp, or Headers should be set by the calling code.
Such a method cannot be annotated with @KafkaPublisher.Topic, because all send details are already contained in the ProducerRecord.
One Callback can be passed additionally.
Methods per topic¶
A method with key, value, and Headers must be annotated with @KafkaPublisher.Topic.
One user argument is treated as value; two user arguments are treated as key and value in that exact order.
Headers and Callback can be declared additionally, but only one argument of each type is allowed.
If Headers is not passed, Kora creates empty headers.
@KafkaPublisher("kafka.someProducer")
public interface MyPublisher {
@KafkaPublisher.Topic("kafka.someProducer.someTopic")
void send(V value);
@KafkaPublisher.Topic("kafka.someProducer.someTopic")
void send(K key, V value);
@KafkaPublisher.Topic("kafka.someProducer.someTopic")
void send(K key, V value, Headers headers);
@KafkaPublisher.Topic("kafka.someProducer.someTopic")
void send(K key, V value, Headers headers, Callback callback);
}
@KafkaPublisher("kafka.someProducer")
interface MyPublisher {
@KafkaPublisher.Topic("kafka.someProducer.someTopic")
fun send(value: V)
@KafkaPublisher.Topic("kafka.someProducer.someTopic")
fun send(key: K, value: V)
@KafkaPublisher.Topic("kafka.someProducer.someTopic")
fun send(key: K, value: V, headers: Headers)
@KafkaPublisher.Topic("kafka.someProducer.someTopic")
fun send(key: K, value: V, headers: Headers, callback: Callback)
}
Send result¶
For a synchronous method, the return type can be void/Unit or RecordMetadata.
In this case, Kora calls KafkaProducer#send, waits for send completion through Future#get(), and only then returns control to the caller.
For asynchronous sending, the return type can be Future<RecordMetadata>, CompletionStage<RecordMetadata>, or CompletableFuture<RecordMetadata>.
In Kotlin, suspend methods and Deferred<RecordMetadata> are also supported.
If the signature contains a Callback, Kora first completes its own send telemetry and then calls the user Callback.
@KafkaPublisher("kafka.someProducer")
public interface MyPublisher {
@KafkaPublisher.Topic("kafka.someProducer.someTopic")
RecordMetadata send(V value);
@KafkaPublisher.Topic("kafka.someProducer.someTopic")
Future<RecordMetadata> sendFuture(V value);
@KafkaPublisher.Topic("kafka.someProducer.someTopic")
CompletionStage<RecordMetadata> sendStage(V value);
@KafkaPublisher.Topic("kafka.someProducer.someTopic")
CompletableFuture<RecordMetadata> sendCompletableFuture(V value);
}
@KafkaPublisher("kafka.someProducer")
interface MyPublisher {
@KafkaPublisher.Topic("kafka.someProducer.someTopic")
fun send(value: V): RecordMetadata
@KafkaPublisher.Topic("kafka.someProducer.someTopic")
suspend fun sendSuspend(value: V): RecordMetadata
@KafkaPublisher.Topic("kafka.someProducer.someTopic")
fun send(value: String): Future<RecordMetadata>
@KafkaPublisher.Topic("kafka.someProducer.someTopic")
fun send(value: String): CompletionStage<RecordMetadata>
@KafkaPublisher.Topic("kafka.someProducer.someTopic")
fun send(value: String): CompletableFuture<RecordMetadata>
@KafkaPublisher.Topic("kafka.someProducer.someTopic")
fun send(value: String): Deferred<RecordMetadata>
}
Invalid combinations are: ProducerRecord<K, V> together with @KafkaPublisher.Topic, ProducerRecord<K, V> together with separate key/value/Headers, more than one Headers, more than one Callback, and a method with separate key/value without @KafkaPublisher.Topic.
Telemetry¶
Kafka uses a telemetry contract for logging, metrics, and tracing of messages.
Telemetry configuration (the telemetry { logging / metrics / tracing } section) is described in the Configuration section.
For each message, KafkaPublisher creates a telemetry context, which is closed upon completion of processing.
The message is described via telemetry handler parameters, including topic, partition, offset, and processing duration.
The default factory, DefaultKafkaPublisherTelemetryFactory, combines three factories:
- KafkaPublisherLoggerFactory builds a KafkaPublisherLogger to log the start and end of message processing;
- KafkaPublisherMetricsFactory builds a KafkaPublisherMetrics to record message metrics;
- KafkaPublisherTracerFactory builds a KafkaPublisherTracer for distributed tracing.
Metrics and tracing are described in the Metrics Reference section.