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;
send(Long.toString(i%10), Long.toString(i));
Thread.sleep(500);
}
- semaphore.acquire(MAX_PENDING_MESSAGES);
}
catch (Exception e)
{
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
}
});
- 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
);
}