Implementierung zum Anfügen der Header wiederhergestellt
[demos/kafka/training] / src / main / java / de / juplo / kafka / RestProducer.java
index 73bec5b..cecb980 100644 (file)
@@ -42,6 +42,12 @@ public class RestProducer
         value   // Value
     );
 
+    record.headers().add("source", id.getBytes());
+    if (correlationId != null)
+    {
+      record.headers().add("id", BigInteger.valueOf(correlationId).toByteArray());
+    }
+
     producer.send(record, (metadata, e) ->
     {
       long now = System.currentTimeMillis();