From 6c736c0c88083bb936169dd5b42ba1a04335336b Mon Sep 17 00:00:00 2001 From: "emanuel.trandafir" Date: Wed, 27 Dec 2023 18:29:16 +0100 Subject: [PATCH 1/3] BAEL-6552: kafka consumer's min/max fetch size --- .../VariableFetchSizeKafkaListener.java | 41 ++++++ ...ariableFetchSizeKafkaListenerLiveTest.java | 117 ++++++++++++++++++ 2 files changed, 158 insertions(+) create mode 100644 apache-kafka-2/src/main/java/com/baeldung/kafka/consumer/VariableFetchSizeKafkaListener.java create mode 100644 apache-kafka-2/src/test/java/com/baeldung/kafka/consumer/VariableFetchSizeKafkaListenerLiveTest.java diff --git a/apache-kafka-2/src/main/java/com/baeldung/kafka/consumer/VariableFetchSizeKafkaListener.java b/apache-kafka-2/src/main/java/com/baeldung/kafka/consumer/VariableFetchSizeKafkaListener.java new file mode 100644 index 0000000000..199a9e7b5c --- /dev/null +++ b/apache-kafka-2/src/main/java/com/baeldung/kafka/consumer/VariableFetchSizeKafkaListener.java @@ -0,0 +1,41 @@ +package com.baeldung.kafka.consumer; + +import java.time.Duration; +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; +import java.util.Properties; + +import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.apache.kafka.clients.consumer.KafkaConsumer; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +public class VariableFetchSizeKafkaListener implements Runnable { + private static Logger log = LoggerFactory.getLogger(VariableFetchSizeKafkaListener.class); + private final String topic; + private final KafkaConsumer consumer; + + public VariableFetchSizeKafkaListener(String topic, Properties consumerProperties) { + this.topic = topic; + this.consumer = new KafkaConsumer<>(consumerProperties); + } + + @Override + public void run() { + consumer.subscribe(Collections.singletonList(topic)); + int pollCount = 1; + while (true) { + List> records = new ArrayList<>(); + for (ConsumerRecord record : consumer.poll(Duration.ofMillis(1_000))) { + records.add(record); + } + if(!records.isEmpty()) { + String batchOffsets = String.format("%s -> %s", records.get(0).offset(), records.get(records.size() - 1).offset()); + String groupId = consumer.groupMetadata().groupId(); + log.info("groupId: {}, poll: #{}, fetched: #{} records, offsets: {}", groupId, pollCount++, records.size(), batchOffsets); + } + } + } + +} diff --git a/apache-kafka-2/src/test/java/com/baeldung/kafka/consumer/VariableFetchSizeKafkaListenerLiveTest.java b/apache-kafka-2/src/test/java/com/baeldung/kafka/consumer/VariableFetchSizeKafkaListenerLiveTest.java new file mode 100644 index 0000000000..2465139d34 --- /dev/null +++ b/apache-kafka-2/src/test/java/com/baeldung/kafka/consumer/VariableFetchSizeKafkaListenerLiveTest.java @@ -0,0 +1,117 @@ +package com.baeldung.kafka.consumer; + +import java.util.List; +import java.util.Properties; +import java.util.concurrent.CompletableFuture; +import java.util.stream.Collectors; +import java.util.stream.IntStream; + +import org.apache.kafka.clients.consumer.ConsumerConfig; +import org.apache.kafka.clients.producer.KafkaProducer; +import org.apache.kafka.clients.producer.ProducerConfig; +import org.apache.kafka.clients.producer.ProducerRecord; +import org.apache.kafka.common.serialization.StringDeserializer; +import org.apache.kafka.common.serialization.StringSerializer; +import org.junit.jupiter.api.Test; +import org.testcontainers.containers.KafkaContainer; +import org.testcontainers.junit.jupiter.Container; +import org.testcontainers.junit.jupiter.Testcontainers; +import org.testcontainers.utility.DockerImageName; + +@Testcontainers +class VariableFetchSizeKafkaListenerLiveTest { + + @Container + private static final KafkaContainer KAFKA_CONTAINER = new KafkaContainer(DockerImageName.parse("confluentinc/cp-kafka:latest")); + + @Test + void whenChangingMaxPartitionFetchBytesProperty_thenAdjustBatchSizesWhilePolling() throws Exception { + publishSensorData(300, "engine.sensors.temperature"); + + // max.partition.fetch.bytes = 500 Bytes + Properties fetchSize_500B = kafkaConsumerProperties(); + fetchSize_500B.setProperty(ConsumerConfig.GROUP_ID_CONFIG, "max_fetch_size_500B"); + fetchSize_500B.setProperty(ConsumerConfig.MAX_PARTITION_FETCH_BYTES_CONFIG, 500 + ""); + CompletableFuture.runAsync( + new VariableFetchSizeKafkaListener("engine.sensors.temperature", fetchSize_500B) + ); + + // max.partition.fetch.bytes = 5.000 Bytes + Properties fetchSize_5KB = kafkaConsumerProperties(); + fetchSize_5KB.setProperty(ConsumerConfig.GROUP_ID_CONFIG, "max_fetch_size_5KB"); + fetchSize_5KB.setProperty(ConsumerConfig.MAX_PARTITION_FETCH_BYTES_CONFIG, 5_000 + ""); + CompletableFuture.runAsync( + new VariableFetchSizeKafkaListener("engine.sensors.temperature", fetchSize_5KB) + ); + + Thread.sleep(10_000L); + } + + @Test + void whenChangingMinFetchBytesProperty_thenAdjustWaitTimeWhilePolling() throws Exception { + publishSensorData(300, "engine.sensors.temperature", 100L); + + // fetch.min.bytes = 1 byte (default) + Properties minFetchSize_1B = kafkaConsumerProperties(); + minFetchSize_1B.setProperty(ConsumerConfig.GROUP_ID_CONFIG, "min_fetch_size_1B"); + minFetchSize_1B.setProperty(ConsumerConfig.FETCH_MIN_BYTES_CONFIG, 1 + ""); + CompletableFuture.runAsync( + new VariableFetchSizeKafkaListener("engine.sensors.temperature", minFetchSize_1B) + ); + + // fetch.min.bytes = 500 bytes + Properties minFetchSize_500B = kafkaConsumerProperties(); + minFetchSize_500B.setProperty(ConsumerConfig.GROUP_ID_CONFIG, "mim_fetch_size_500B"); + minFetchSize_500B.setProperty(ConsumerConfig.FETCH_MIN_BYTES_CONFIG, 500 + ""); + CompletableFuture.runAsync( + new VariableFetchSizeKafkaListener("engine.sensors.temperature", minFetchSize_500B) + ); + + Thread.sleep(10_000L); + } + + + private static Properties kafkaConsumerProperties() { + Properties props = new Properties(); + props.setProperty(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, KAFKA_CONTAINER.getBootstrapServers()); + props.setProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); + props.setProperty(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); + props.setProperty(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); + return props; + } + + private void publishSensorData(int measurementsCount, String topic) { + publishSensorData(measurementsCount, topic, 0L); + } + + private void publishSensorData(int measurementsCount, String topic, long delayInMillis) { + List> records = IntStream.range(0, measurementsCount) + .mapToObj(__ -> new ProducerRecord<>(topic, "key1", "temperature=255F")) + .collect(Collectors.toList()); + + CompletableFuture.runAsync(() -> { + try (KafkaProducer producer = testKafkaProducer()) { + for (ProducerRecord rec : records) { + producer.send(rec); + sleep(delayInMillis); + } + } + }); + } + + private static KafkaProducer testKafkaProducer() { + Properties props = new Properties(); + props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, KAFKA_CONTAINER.getBootstrapServers()); + props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); + props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); + return new KafkaProducer<>(props); + } + + private static void sleep(long delayInMillis) { + try { + Thread.sleep(delayInMillis); + } catch (InterruptedException e) { + throw new RuntimeException(e); + } + } +} \ No newline at end of file From 1ac666749b96342a97df1ebc54f9020f9ac79570 Mon Sep 17 00:00:00 2001 From: emanueltrandafir1993 Date: Wed, 27 Dec 2023 22:33:40 +0100 Subject: [PATCH 2/3] BAEL-6552: added test using the default configuration --- .../VariableFetchSizeKafkaListener.java | 7 ++-- ...ariableFetchSizeKafkaListenerLiveTest.java | 39 +++++++++++++------ 2 files changed, 32 insertions(+), 14 deletions(-) diff --git a/apache-kafka-2/src/main/java/com/baeldung/kafka/consumer/VariableFetchSizeKafkaListener.java b/apache-kafka-2/src/main/java/com/baeldung/kafka/consumer/VariableFetchSizeKafkaListener.java index 199a9e7b5c..5de2c0e8cb 100644 --- a/apache-kafka-2/src/main/java/com/baeldung/kafka/consumer/VariableFetchSizeKafkaListener.java +++ b/apache-kafka-2/src/main/java/com/baeldung/kafka/consumer/VariableFetchSizeKafkaListener.java @@ -16,18 +16,19 @@ public class VariableFetchSizeKafkaListener implements Runnable { private final String topic; private final KafkaConsumer consumer; - public VariableFetchSizeKafkaListener(String topic, Properties consumerProperties) { + public VariableFetchSizeKafkaListener(String topic, KafkaConsumer consumer) { this.topic = topic; - this.consumer = new KafkaConsumer<>(consumerProperties); + this.consumer = consumer; } @Override public void run() { consumer.subscribe(Collections.singletonList(topic)); int pollCount = 1; + while (true) { List> records = new ArrayList<>(); - for (ConsumerRecord record : consumer.poll(Duration.ofMillis(1_000))) { + for (ConsumerRecord record : consumer.poll(Duration.ofMillis(500))) { records.add(record); } if(!records.isEmpty()) { diff --git a/apache-kafka-2/src/test/java/com/baeldung/kafka/consumer/VariableFetchSizeKafkaListenerLiveTest.java b/apache-kafka-2/src/test/java/com/baeldung/kafka/consumer/VariableFetchSizeKafkaListenerLiveTest.java index 2465139d34..1a794af48e 100644 --- a/apache-kafka-2/src/test/java/com/baeldung/kafka/consumer/VariableFetchSizeKafkaListenerLiveTest.java +++ b/apache-kafka-2/src/test/java/com/baeldung/kafka/consumer/VariableFetchSizeKafkaListenerLiveTest.java @@ -7,6 +7,7 @@ import java.util.stream.Collectors; import java.util.stream.IntStream; import org.apache.kafka.clients.consumer.ConsumerConfig; +import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.clients.producer.ProducerRecord; @@ -24,24 +25,40 @@ class VariableFetchSizeKafkaListenerLiveTest { @Container private static final KafkaContainer KAFKA_CONTAINER = new KafkaContainer(DockerImageName.parse("confluentinc/cp-kafka:latest")); + @Test + void whenUsingDefaultConfiguration_thenProcessInBatchesOf() throws Exception { + String topic = "engine.sensors.temperature"; + publishSensorData(300, topic); + + Properties props = commonConsumerProperties(); + props.setProperty(ConsumerConfig.GROUP_ID_CONFIG, "default_config"); + KafkaConsumer kafkaConsumer = new KafkaConsumer<>(props); + + CompletableFuture.runAsync( + new VariableFetchSizeKafkaListener(topic, kafkaConsumer) + ); + + Thread.sleep(10_000L); + } + @Test void whenChangingMaxPartitionFetchBytesProperty_thenAdjustBatchSizesWhilePolling() throws Exception { publishSensorData(300, "engine.sensors.temperature"); // max.partition.fetch.bytes = 500 Bytes - Properties fetchSize_500B = kafkaConsumerProperties(); + Properties fetchSize_500B = commonConsumerProperties(); fetchSize_500B.setProperty(ConsumerConfig.GROUP_ID_CONFIG, "max_fetch_size_500B"); fetchSize_500B.setProperty(ConsumerConfig.MAX_PARTITION_FETCH_BYTES_CONFIG, 500 + ""); CompletableFuture.runAsync( - new VariableFetchSizeKafkaListener("engine.sensors.temperature", fetchSize_500B) + new VariableFetchSizeKafkaListener("engine.sensors.temperature", new KafkaConsumer<>(fetchSize_500B)) ); // max.partition.fetch.bytes = 5.000 Bytes - Properties fetchSize_5KB = kafkaConsumerProperties(); + Properties fetchSize_5KB = commonConsumerProperties(); fetchSize_5KB.setProperty(ConsumerConfig.GROUP_ID_CONFIG, "max_fetch_size_5KB"); fetchSize_5KB.setProperty(ConsumerConfig.MAX_PARTITION_FETCH_BYTES_CONFIG, 5_000 + ""); CompletableFuture.runAsync( - new VariableFetchSizeKafkaListener("engine.sensors.temperature", fetchSize_5KB) + new VariableFetchSizeKafkaListener("engine.sensors.temperature", new KafkaConsumer<>(fetchSize_5KB)) ); Thread.sleep(10_000L); @@ -52,26 +69,26 @@ class VariableFetchSizeKafkaListenerLiveTest { publishSensorData(300, "engine.sensors.temperature", 100L); // fetch.min.bytes = 1 byte (default) - Properties minFetchSize_1B = kafkaConsumerProperties(); + Properties minFetchSize_1B = commonConsumerProperties(); minFetchSize_1B.setProperty(ConsumerConfig.GROUP_ID_CONFIG, "min_fetch_size_1B"); minFetchSize_1B.setProperty(ConsumerConfig.FETCH_MIN_BYTES_CONFIG, 1 + ""); CompletableFuture.runAsync( - new VariableFetchSizeKafkaListener("engine.sensors.temperature", minFetchSize_1B) + new VariableFetchSizeKafkaListener("engine.sensors.temperature", new KafkaConsumer<>(minFetchSize_1B)) ); // fetch.min.bytes = 500 bytes - Properties minFetchSize_500B = kafkaConsumerProperties(); + Properties minFetchSize_500B = commonConsumerProperties(); minFetchSize_500B.setProperty(ConsumerConfig.GROUP_ID_CONFIG, "mim_fetch_size_500B"); minFetchSize_500B.setProperty(ConsumerConfig.FETCH_MIN_BYTES_CONFIG, 500 + ""); CompletableFuture.runAsync( - new VariableFetchSizeKafkaListener("engine.sensors.temperature", minFetchSize_500B) + new VariableFetchSizeKafkaListener("engine.sensors.temperature", new KafkaConsumer<>(minFetchSize_500B)) ); Thread.sleep(10_000L); } - private static Properties kafkaConsumerProperties() { + private static Properties commonConsumerProperties() { Properties props = new Properties(); props.setProperty(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, KAFKA_CONTAINER.getBootstrapServers()); props.setProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); @@ -86,8 +103,8 @@ class VariableFetchSizeKafkaListenerLiveTest { private void publishSensorData(int measurementsCount, String topic, long delayInMillis) { List> records = IntStream.range(0, measurementsCount) - .mapToObj(__ -> new ProducerRecord<>(topic, "key1", "temperature=255F")) - .collect(Collectors.toList()); + .mapToObj(__ -> new ProducerRecord<>(topic, "key1", "temperature=255F")) + .collect(Collectors.toList()); CompletableFuture.runAsync(() -> { try (KafkaProducer producer = testKafkaProducer()) { From c5892eec6edf535f78fdda443d4853974cbe641f Mon Sep 17 00:00:00 2001 From: "emanuel.trandafir" Date: Thu, 28 Dec 2023 15:10:10 +0100 Subject: [PATCH 3/3] BAEL-6552: small changes --- ...ariableFetchSizeKafkaListenerLiveTest.java | 45 ++++++++++++------- 1 file changed, 28 insertions(+), 17 deletions(-) diff --git a/apache-kafka-2/src/test/java/com/baeldung/kafka/consumer/VariableFetchSizeKafkaListenerLiveTest.java b/apache-kafka-2/src/test/java/com/baeldung/kafka/consumer/VariableFetchSizeKafkaListenerLiveTest.java index 1a794af48e..1dcf4784f9 100644 --- a/apache-kafka-2/src/test/java/com/baeldung/kafka/consumer/VariableFetchSizeKafkaListenerLiveTest.java +++ b/apache-kafka-2/src/test/java/com/baeldung/kafka/consumer/VariableFetchSizeKafkaListenerLiveTest.java @@ -3,6 +3,7 @@ package com.baeldung.kafka.consumer; import java.util.List; import java.util.Properties; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ExecutionException; import java.util.stream.Collectors; import java.util.stream.IntStream; @@ -28,9 +29,13 @@ class VariableFetchSizeKafkaListenerLiveTest { @Test void whenUsingDefaultConfiguration_thenProcessInBatchesOf() throws Exception { String topic = "engine.sensors.temperature"; - publishSensorData(300, topic); + publishTestData(300, topic); - Properties props = commonConsumerProperties(); + Properties props = new Properties(); + props.setProperty(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, KAFKA_CONTAINER.getBootstrapServers()); + props.setProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); + props.setProperty(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); + props.setProperty(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); props.setProperty(ConsumerConfig.GROUP_ID_CONFIG, "default_config"); KafkaConsumer kafkaConsumer = new KafkaConsumer<>(props); @@ -38,27 +43,29 @@ class VariableFetchSizeKafkaListenerLiveTest { new VariableFetchSizeKafkaListener(topic, kafkaConsumer) ); - Thread.sleep(10_000L); + Thread.sleep(5_000L); } @Test void whenChangingMaxPartitionFetchBytesProperty_thenAdjustBatchSizesWhilePolling() throws Exception { - publishSensorData(300, "engine.sensors.temperature"); + String topic = "engine.sensors.temperature"; + publishTestData(300, topic); + Thread.sleep(1_000L); // max.partition.fetch.bytes = 500 Bytes Properties fetchSize_500B = commonConsumerProperties(); fetchSize_500B.setProperty(ConsumerConfig.GROUP_ID_CONFIG, "max_fetch_size_500B"); - fetchSize_500B.setProperty(ConsumerConfig.MAX_PARTITION_FETCH_BYTES_CONFIG, 500 + ""); + fetchSize_500B.setProperty(ConsumerConfig.MAX_PARTITION_FETCH_BYTES_CONFIG, "500"); CompletableFuture.runAsync( - new VariableFetchSizeKafkaListener("engine.sensors.temperature", new KafkaConsumer<>(fetchSize_500B)) + new VariableFetchSizeKafkaListener(topic, new KafkaConsumer<>(fetchSize_500B)) ); // max.partition.fetch.bytes = 5.000 Bytes Properties fetchSize_5KB = commonConsumerProperties(); fetchSize_5KB.setProperty(ConsumerConfig.GROUP_ID_CONFIG, "max_fetch_size_5KB"); - fetchSize_5KB.setProperty(ConsumerConfig.MAX_PARTITION_FETCH_BYTES_CONFIG, 5_000 + ""); + fetchSize_5KB.setProperty(ConsumerConfig.MAX_PARTITION_FETCH_BYTES_CONFIG, "5000"); CompletableFuture.runAsync( - new VariableFetchSizeKafkaListener("engine.sensors.temperature", new KafkaConsumer<>(fetchSize_5KB)) + new VariableFetchSizeKafkaListener(topic, new KafkaConsumer<>(fetchSize_5KB)) ); Thread.sleep(10_000L); @@ -66,22 +73,22 @@ class VariableFetchSizeKafkaListenerLiveTest { @Test void whenChangingMinFetchBytesProperty_thenAdjustWaitTimeWhilePolling() throws Exception { - publishSensorData(300, "engine.sensors.temperature", 100L); + String topic = "engine.sensors.temperature"; + publishTestData(300, topic, 100L); // fetch.min.bytes = 1 byte (default) Properties minFetchSize_1B = commonConsumerProperties(); minFetchSize_1B.setProperty(ConsumerConfig.GROUP_ID_CONFIG, "min_fetch_size_1B"); - minFetchSize_1B.setProperty(ConsumerConfig.FETCH_MIN_BYTES_CONFIG, 1 + ""); CompletableFuture.runAsync( - new VariableFetchSizeKafkaListener("engine.sensors.temperature", new KafkaConsumer<>(minFetchSize_1B)) + new VariableFetchSizeKafkaListener(topic, new KafkaConsumer<>(minFetchSize_1B)) ); // fetch.min.bytes = 500 bytes Properties minFetchSize_500B = commonConsumerProperties(); minFetchSize_500B.setProperty(ConsumerConfig.GROUP_ID_CONFIG, "mim_fetch_size_500B"); - minFetchSize_500B.setProperty(ConsumerConfig.FETCH_MIN_BYTES_CONFIG, 500 + ""); + minFetchSize_500B.setProperty(ConsumerConfig.FETCH_MIN_BYTES_CONFIG, "500"); CompletableFuture.runAsync( - new VariableFetchSizeKafkaListener("engine.sensors.temperature", new KafkaConsumer<>(minFetchSize_500B)) + new VariableFetchSizeKafkaListener(topic, new KafkaConsumer<>(minFetchSize_500B)) ); Thread.sleep(10_000L); @@ -97,11 +104,11 @@ class VariableFetchSizeKafkaListenerLiveTest { return props; } - private void publishSensorData(int measurementsCount, String topic) { - publishSensorData(measurementsCount, topic, 0L); + private void publishTestData(int measurementsCount, String topic) { + publishTestData(measurementsCount, topic, 0L); } - private void publishSensorData(int measurementsCount, String topic, long delayInMillis) { + private void publishTestData(int measurementsCount, String topic, long delayInMillis) { List> records = IntStream.range(0, measurementsCount) .mapToObj(__ -> new ProducerRecord<>(topic, "key1", "temperature=255F")) .collect(Collectors.toList()); @@ -109,13 +116,17 @@ class VariableFetchSizeKafkaListenerLiveTest { CompletableFuture.runAsync(() -> { try (KafkaProducer producer = testKafkaProducer()) { for (ProducerRecord rec : records) { - producer.send(rec); + producer.send(rec).get(); sleep(delayInMillis); } + } catch (ExecutionException | InterruptedException e) { + throw new RuntimeException(e); } }); } + + private static KafkaProducer testKafkaProducer() { Properties props = new Properties(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, KAFKA_CONTAINER.getBootstrapServers());