From: Kai Moritz Date: Sun, 13 Sep 2026 13:55:05 +0000 (+0200) Subject: refactor: Implementierung überarbeitet & vereinfacht X-Git-Url: http://juplo.de/gitweb/?a=commitdiff_plain;h=015517f628c6edaecd0fc157218b6ee7a5a4c7e7;p=demos%2Fkafka%2Ftraining refactor: Implementierung überarbeitet & vereinfacht * Die vorliegende Version wurde ausgiebig mit Chat GPT diskutiert. * Das _asynchrone_ Queuing der Nachrichten ist nicht notwendig. * Um das Ziel, die asynchrone Verarbeitung der Ergebnisse asynchron umzusetzen, ohne dabei den IO-Thread der Producer-Instanz zu blockieren, zu erreichen, genügt es, das `CompleteableFuture` für die Verarbeitung des Ergebnisses _synchron_ über `CompletableFuture.completedFuture()` zu erzeugen und erst die blockierende Weiterverarbeitung, über `thenApplyAsync()` an einen anderen Thread zu delegieren! * D.h. auch, dass die verwirrende zusätzliche Log-Meldung "Scheduled queuing of message..." ganz entfallen kann und an deren Stelle wieder unverändert das Logging erfolgt, das festhält, wie lange der Aufruf von `send()` gedauert hat. * Außerdem wird die `Semaphore` nicht mehr benötigt, da alle Aufrufe von `send()` _synchron_ in dem Thread erfolgen, der die `for`-Schleife abarbeitet. D.h., wenn diese Schleife beendet wird, können nicht mehr wie bisher noch nachträglich asynchron weitere Aufrufe von `send()` erfolgen, auf die zuvor mit Hilfe der `Semapohre` gewartet werden musste, damit kein Fehler entsteht, weil `send()` aufgerufen wird, nachdem `close` auf der `Producer`-Instanz aufgerufen wurde. * Insgesamt schrumpfen die Unterschiede zu der Implementierung, die den `Callback` verwendet auch angenehm auf ein Minimum zurück. * *Beachte:* ** In einer produktiven Anwendung, sollte der Aufruf von `send()` in einen `try`/`catch`-Block eingefasst werden, um eine korrekte Fehlerbehandlung der bereits synchron beim Aufruf von `send()` erfolgten Aufrufe sicherzustellen! ** Dies gilt aber genauso für die Implementierung, die den `Callback` verwendet, ist also der einfachen Vorführ-Implementierung geschuldet, die bei jedem synchronen Fehler abbricht. ** Hier kommt lediglich dazu, dass zusätzlich eine `InterruptedException`, die von dem blockierenden Aufruf von `Future.get()` ausgelöst werden kann, das Programm beendent könnte. --- diff --git a/src/main/java/de/juplo/kafka/ExampleProducer.java b/src/main/java/de/juplo/kafka/ExampleProducer.java index d62d12af..78cb778d 100644 --- a/src/main/java/de/juplo/kafka/ExampleProducer.java +++ b/src/main/java/de/juplo/kafka/ExampleProducer.java @@ -8,20 +8,19 @@ import org.apache.kafka.clients.producer.RecordMetadata; import org.apache.kafka.common.serialization.StringSerializer; import java.util.Properties; -import java.util.concurrent.*; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; @Slf4j public class ExampleProducer { - public final static int MAX_PENDING_MESSAGES = 100; - private final String id; private final String topic; private final Producer producer; private final ExecutorService executor = Executors.newSingleThreadExecutor(); - private final Semaphore semaphore = new Semaphore(MAX_PENDING_MESSAGES); private volatile boolean running = true; private volatile boolean done = false; private long produced = 0; @@ -53,7 +52,6 @@ public class ExampleProducer send(Long.toString(i%10), Long.toString(i)); Thread.sleep(500); } - semaphore.acquire(MAX_PENDING_MESSAGES); } catch (Exception e) { @@ -78,22 +76,8 @@ public class ExampleProducer value // Value ); - semaphore.acquire(); CompletableFuture completableFuture = CompletableFuture - .supplyAsync(() -> - { - Future recordMetadataFuture = producer.send(record); - long sendRequestQueued = System.currentTimeMillis(); - semaphore.release(); - log.trace( - "{} - Queued message {}={}, latency={}ms", - id, - key, - value, - sendRequestQueued - sendRequested - ); - return recordMetadataFuture; - }, executor) + .completedFuture(producer.send(record)) .thenApplyAsync(recordMetadataFuture -> { try @@ -137,14 +121,14 @@ public class ExampleProducer } }); - long queuingOfSendRequestScheduled = System.currentTimeMillis(); + long sendRequestQueued = System.currentTimeMillis(); produced++; log.trace( - "{} - Scheduled queuing of message {}={}, latency={}ms", + "{} - Queued message {}={}, latency={}ms", id, key, value, - queuingOfSendRequestScheduled - sendRequested + sendRequestQueued - sendRequested ); }