import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
+import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.springframework.http.HttpStatus;
private final String id;
private final String topic;
private final Integer partition;
- private final KafkaProducer<String, String> producer;
+ private final Producer<String, String> producer;
private long produced = 0;
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();