Install
$ agentstack add skill-kora-projects-kora-skills-kora-kafka-consumer ✓ scanned · ✓ verified, works with Claude Code, Cursor, and more.
Security review
✓ PassedNo issues found. Passed automated security review. · v0.1.0 How review works →
- ✓ Prompt-injection patterns
- ✓ Secret / credential exfiltration
- ✓ Dangerous shell & filesystem operations
- ✓ Untrusted network calls
- ✓ Known-malicious package signatures
What it can access
- ✓ Network access No
- ✓ Filesystem access No
- ✓ Shell / process execution No
- ✓ Environment & secrets No
- ✓ Dynamic code execution No
From automated source analysis of v0.1.0. “Used” means the capability is present in the source — more access means more to trust, not that it’s unsafe.
Verified badge
Passed review? Show it. Paste this badge into your README, it links to the public security report.
Reliability & compatibility
Declared compatibility
Compatibility is declared by the source manifest. End-to-end runtime verification is coming, see below.
We're building live execution health for every listing: tool-call success rate, median latency, uptime, and last-checked timestamps, measured, not self-reported. It isn't live yet, so we don't show numbers we can't stand behind.
How agent discovery & health will work →About
Kora Kafka Consumer Skill
Languages: Java, Kotlin | Build: Gradle
> Level 1 — [Quick Start](#quick-start) | Level 2 — [Signatures](#method-signatures) | Level 3 — [Errors & Offset](#error-handling) | Level 4 — [Batch & Rebalance](#batch-processing)
References: [Consumer config](references/kafka-consumer-reference.md) | [Listener signatures](references/kafka-listener-reference.md) | [Strategies](references/kafka-strategies-reference.md) | [Serialization](references/kafka-serialization-reference.md) | [Errors](references/kafka-error-handling-reference.md) | [Offset](references/kafka-offset-reference.md) | [Batch](references/kafka-batch-reference.md) | [Rebalance](references/kafka-rebalance-reference.md) | [Telemetry](references/kafka-telemetry-reference.md) | [Transactions (producer)](references/kafka-transactions-reference.md) | [Testing](references/kafka-testing-reference.md)
Quick Start
1. Dependencies (Kora artifacts inherit their version from the kora-parent BOM — never pin them individually):
dependencies {
koraBom platform("ru.tinkoff.kora:kora-parent:1.2.17")
annotationProcessor "ru.tinkoff.kora:annotation-processors"
implementation "ru.tinkoff.kora:kafka"
implementation "ru.tinkoff.kora:json-module"
implementation "ru.tinkoff.kora:config-hocon"
implementation "ru.tinkoff.kora:logging-logback"
}
2. Application Module (a @KoraApp interface extends each module — interfaces never use implements):
@KoraApp
public interface Application extends KafkaModule, JsonModule, HoconConfigModule, LogbackModule {
static void main(String[] args) {
KoraApplication.run(ApplicationGraph::graph);
}
}
3. Simple Listener:
@Component
public final class UserEventListener {
@KafkaListener("kafka.consumer.userEvents")
void process(String value) {
log.info("Received: {}", value);
}
}
4. Configuration:
kafka {
consumer {
userEvents {
topics = ["user-events"]
driverProperties {
"bootstrap.servers" = ${KAFKA_BOOTSTRAP:"localhost:9092"}
"group.id" = "user-service"
"auto.offset.reset" = "earliest"
}
}
}
}
Method Signatures
| Signature | Commit | Use Case | |-----------|--------|----------| | void process(String value) | Auto | Simple processing | | void process(String key, String value) | Auto | Key-aware | | void process(ConsumerRecord record) | Auto | Full metadata | | void process(ConsumerRecords records) | Auto/batch | Batch processing | | void process(..., Consumer consumer) | Manual | Exactly-once | | void process(@Nullable T, @Nullable Exception) | Auto | Error handling |
JSON with Error Handling:
@Json record OrderEvent(String orderId, BigDecimal amount) {}
@KafkaListener("kafka.consumer.orders")
void process(@Nullable @Json OrderEvent event, @Nullable Exception error) {
if (error != null) {
log.error("Deserialization failed", error);
return;
}
orderService.process(event);
}
Full signatures: [Listener Reference](references/kafka-listener-reference.md)
Strategies
Subscribe (load balancing):
driverProperties { "group.id" = "order-service"; "bootstrap.servers" = "localhost:9092" }
Assign (broadcast):
driverProperties { "bootstrap.servers" = "localhost:9092" } # No group.id
Details: [Strategies Reference](references/kafka-strategies-reference.md)
Error Handling
Deserialization:
@KafkaListener("kafka.consumer.events")
void process(@Nullable @Json Event event, @Nullable Exception error) {
if (error != null) { log.error("Failed", error); return; }
processEvent(event);
}
Skip Invalid:
@KafkaListener("kafka.consumer.events")
void process(Event event) {
if (event.orderId() == null)
throw new KafkaSkipRecordException(new IllegalArgumentException("Missing orderId"));
processEvent(event);
}
DLQ:
@KafkaListener("kafka.consumer.events")
void process(@Nullable Event event, @Nullable Exception error) {
if (error != null) { dlqPublisher.send("dlq", event, error.getMessage()); return; }
processEvent(event);
}
Details: [Error Handling Reference](references/kafka-error-handling-reference.md)
Offset Management
Auto (Default): Commit after each message/batch.
Manual:
@KafkaListener("kafka.consumer.events")
void process(ConsumerRecord record, Consumer consumer) {
try { process(record.value()); consumer.commitSync(); }
catch (Exception e) { throw e; }
}
Rebalance: provide a ConsumerAwareRebalanceListener as a @Component carrying the consumer's tag. Kora generates a tag per listener (Module.ProcessTag), or you can declare your own via @KafkaListener(value = "...", tag = MyTag.class) and reuse it:
@Tag(MyTag.class) @Component
final class RebalanceListener implements ConsumerAwareRebalanceListener {
public void onPartitionsRevoked(Consumer c, Collection p) {
c.commitSync(); // Commit before rebalance
}
public void onPartitionsAssigned(Consumer c, Collection p) { }
}
Details: [Offset Reference](references/kafka-offset-reference.md)
Batch Processing
@KafkaListener("kafka.consumer.orders")
void process(ConsumerRecords records) {
for (ConsumerRecord record : records)
orderService.process(record.value());
}
Config:
kafka.consumer.batchProcessor {
topics = ["high-volume"]
threads = 4
driverProperties { "max.poll.records" = 500; "fetch.min.bytes" = 1048576 }
}
Details: [Batch Reference](references/kafka-batch-reference.md)
Rebalance Handling
@Tag(MyTag.class) @Component
final class MyRebalanceListener implements ConsumerAwareRebalanceListener {
public void onPartitionsRevoked(Consumer c, Collection p) {
log.info("Revoked: {}", p); c.commitSync(); cache.clear();
}
public void onPartitionsAssigned(Consumer c, Collection p) {
log.info("Assigned: {}", p);
}
public void onPartitionsLost(Consumer c, Collection p) {
log.warn("Lost: {}", p); // Don't commit
}
}
Details: [Rebalance Reference](references/kafka-rebalance-reference.md)
Configuration
Required:
kafka.consumer.myListener {
topics = ["topic1"]
driverProperties { "bootstrap.servers" = "localhost:9092" }
}
Optional:
| Parameter | Default | Description | |-----------|---------|-------------| | offset | latest | earliest, latest, or duration (5m) | | pollTimeout | 5s | Max wait for messages | | backoffTimeout | 15s | Pause after exception | | threads | 1 | Parallel threads | | shutdownWait | 30s | Graceful shutdown |
Telemetry:
telemetry { logging {enabled=true}; metrics {enabled=true}; tracing {enabled=true} }
Common Pitfalls
| Symptom | Cause | Fix | |---------|-------|-----| | @KoraApp does not compile | interface Application implements KafkaModule | An interface extends modules, never implements | | cannot find symbol KafkaSkipRecordException | Wrong import | Import ru.tinkoff.kora.kafka.common.exceptions.KafkaSkipRecordException | | Consumer never starts | threads = 0 in config | Use threads >= 1 (0 disables the consumer entirely) | | Listener restarts in a loop | Handler throws an unhandled exception | Kora restarts the consumer on uncaught exceptions; throw KafkaSkipRecordException to skip, or handle and return | | offset = "5m" ignored | group.id is set | offset/duration applies only in assign mode (no group.id); with a group, committed offsets win | | Deserialization error crashes handler | Plain value signature | Add @Nullable Exception as the last parameter, or catch RecordValueDeserializationException when using ConsumerRecord | | Rebalance listener never invoked | @Tag does not match the listener | Tag the listener (@KafkaListener(tag = MyTag.class)) and the ConsumerAwareRebalanceListener with the same tag | | Nothing generated after refactor | Annotation processor stale | Clean build/generated/, rerun ./gradlew classes |
Do not use field injection — Kora wires components through constructor injection at compile time, and every consumer is a @Component with @KafkaListener methods.
Templates
Assets: ConsumerListener.java.template, JsonMessageListener.java.template, ConsumerListenerTests.java.template, application.conf.template
See [assets/README.md](assets/README.md) for generator script.
Testing
Use @KoraAppTest, inject the listener with @TestComponent, await the consumer's own collected state with Awaitility, and drive Kafka with Testcontainers. The example app uses io.goodforgod:testcontainers-extensions-kafka for a thin KafkaConnection:
@TestcontainersKafka(mode = ContainerMode.PER_RUN, topics = @Topics("my-topic-consumer"))
@KoraAppTest(Application.class)
class MyListenerTests implements KoraAppTestConfigModifier {
@ConnectionKafka
private KafkaConnection connection;
@TestComponent
private MyListener consumer;
@Override
public KoraConfigModification config() {
return KoraConfigModification
.ofSystemProperty("KAFKA_BOOTSTRAP", connection.params().bootstrapServers());
}
@Test
void processed() {
connection.send("my-topic-consumer", Event.ofValueAndRandomKey("hello".getBytes()));
Awaitility.await().atMost(Duration.ofSeconds(15))
.until(() -> consumer.received().size() == 1);
}
}
The matching config keeps ${KAFKA_BOOTSTRAP} as the placeholder used in application.conf. To await the consumer container starting, inject its generated tag as Lifecycle: @Tag(MyListenerModule.MyListenerProcessTag.class) @TestComponent Lifecycle.
Details: [Testing Reference](references/kafka-testing-reference.md)
Source of truth: [Kafka doc](../../.kora-agent/kora-docs/mkdocs/docs/en/documentation/kafka.md) | [Messaging guide](../../.kora-agent/kora-docs/mkdocs/docs/en/guides/messaging-kafka.md) | [Example app](../../.kora-agent/kora-examples/examples/java/kora-java-kafka/)
Source & license
This open-source skill is cataloged on AgentStack and links to its original source — we do not rehost the code.
- Author: kora-projects
- Source: kora-projects/kora-skills
- License: Apache-2.0
- Homepage: http://kora-projects.github.io/kora-docs
Install and usage instructions live in the source repository linked above.
Reviews
No reviews yet, be the first.
Write a review
Versions
- v0.1.0 Imported from the upstream source.