]> juplo.de Git - demos/kafka/training/commitdiff
refactor: Implementierung überarbeitet & vereinfacht grundlagen/simple-producer--completablefuture
authorKai Moritz <kai@juplo.de>
Sun, 13 Sep 2026 13:55:05 +0000 (15:55 +0200)
committerKai Moritz <kai@juplo.de>
Sun, 13 Sep 2026 14:27:56 +0000 (16:27 +0200)
* 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.

src/main/java/de/juplo/kafka/ExampleProducer.java

index d62d12afeccb6a49ab6481fc27c7efe2863f6a66..78cb778df5094b5165a1735b96c88bfb3df239d2 100644 (file)
@@ -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<String, String> 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<RecordMetadata> completableFuture = CompletableFuture
-      .supplyAsync(() ->
-      {
-        Future<RecordMetadata> 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
     );
   }