org.apache.kafka:kafka-clients, the reference implementation, ships with every Kafka 129 release (4.3.1 here), gets protocol features first (such as Share Groups's KafkaShareConsumer) and underlies Streams, Connect and Spring. The producer is thread-safe; a KafkaConsumer is not, so give each thread its own. The analytics consumer once more, configured with group.protocol=consumer, manual commits and max.poll.records=5000, and Jackson for the JSON:
try (var consumer = new KafkaConsumer<String, String>(config)) {
consumer.subscribe(List.of("booknest.order-events"));
for (int idle = 0; idle < 3; ) { // stop after 3 empty polls
var batch = consumer.poll(Duration.ofSeconds(1));
if (batch.isEmpty()) { idle += events == 0 ? 0 : 1; continue; }
for (var r : batch) {
var header = r.headers().lastHeader("event-type");
if (header == null) continue; // test records
var kind = new String(header.value());
events++;
if (kind.equals("order_placed"))
totals.put(r.key(), json.readTree(r.value()).get("total").asDouble());
else if (kind.equals("order_cancelled")) cancelled.add(r.key());
}
consumer.commitSync(); // at least once, per batch
last = System.nanoTime();
}
}java: 390737 events in 3.210 s orders 100000, cancelled 5890, kept 3270268.02
clients/java_run.sh ran the single file with java -cp 'libs/*' OrderStats.java, using jars copied from the broker image; a project declares kafka-clients in Maven 129 or Gradle 19,597 . Without an SLF4J binding such as Logback the client warns that it logs nothing. Java was the fastest of the three, taking records in batches of 5,000.