From 2cecae1dfb6af13c311e48d89dfc38d1f70ba181 Mon Sep 17 00:00:00 2001 From: Amol Gote Date: Sat, 4 Nov 2023 17:18:20 -0400 Subject: [PATCH] Incorporate Review comments --- .../com/baeldung/kafka/message/ordering/Config.java | 2 +- ...ternalSequenceWithTimeWindowIntegrationTest.java} | 12 +++++------- .../ordering/MultiplePartitionIntegrationTest.java | 7 +++---- .../ordering/SinglePartitionIntegrationTest.java | 6 ++---- 4 files changed, 11 insertions(+), 16 deletions(-) rename apache-kafka-2/src/test/java/com/baeldung/kafka/message/ordering/{ExtSeqWithTimeWindowIntegrationTest.java => ExternalSequenceWithTimeWindowIntegrationTest.java} (95%) diff --git a/apache-kafka-2/src/main/java/com/baeldung/kafka/message/ordering/Config.java b/apache-kafka-2/src/main/java/com/baeldung/kafka/message/ordering/Config.java index 9cc6314309..7fae8403b5 100644 --- a/apache-kafka-2/src/main/java/com/baeldung/kafka/message/ordering/Config.java +++ b/apache-kafka-2/src/main/java/com/baeldung/kafka/message/ordering/Config.java @@ -7,5 +7,5 @@ public class Config { public static final int MULTIPLE_PARTITIONS = 5; public static final int SINGLE_PARTITION = 1; - public static short REPLICATION_FACTOR = 1; + public static final short REPLICATION_FACTOR = 1; } diff --git a/apache-kafka-2/src/test/java/com/baeldung/kafka/message/ordering/ExtSeqWithTimeWindowIntegrationTest.java b/apache-kafka-2/src/test/java/com/baeldung/kafka/message/ordering/ExternalSequenceWithTimeWindowIntegrationTest.java similarity index 95% rename from apache-kafka-2/src/test/java/com/baeldung/kafka/message/ordering/ExtSeqWithTimeWindowIntegrationTest.java rename to apache-kafka-2/src/test/java/com/baeldung/kafka/message/ordering/ExternalSequenceWithTimeWindowIntegrationTest.java index 76e4a47d17..0c64f663f3 100644 --- a/apache-kafka-2/src/test/java/com/baeldung/kafka/message/ordering/ExtSeqWithTimeWindowIntegrationTest.java +++ b/apache-kafka-2/src/test/java/com/baeldung/kafka/message/ordering/ExternalSequenceWithTimeWindowIntegrationTest.java @@ -11,7 +11,6 @@ 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.clients.producer.RecordMetadata; -import org.apache.kafka.common.KafkaFuture; import org.apache.kafka.common.serialization.LongDeserializer; import org.apache.kafka.common.serialization.LongSerializer; import org.junit.jupiter.api.AfterAll; @@ -25,8 +24,8 @@ import java.time.Duration; import java.util.*; import java.util.concurrent.ExecutionException; import java.util.concurrent.Future; - -import static org.junit.jupiter.api.Assertions.*; +import com.google.common.collect.ImmutableList; +import static org.assertj.core.api.AssertionsForInterfaceTypes.assertThat; @Testcontainers public class ExternalSequenceWithTimeWindowIntegrationTest { @@ -84,7 +83,6 @@ public class ExternalSequenceWithTimeWindowIntegrationTest { System.out.println("User Event ID: " + userEvent.getUserEventId() + ", Partition : " + metadata.partition()); } - boolean isOrderMaintained = true; consumer.subscribe(Collections.singletonList(Config.MULTI_PARTITION_TOPIC)); List buffer = new ArrayList<>(); long lastProcessedTime = System.nanoTime(); @@ -102,9 +100,9 @@ public class ExternalSequenceWithTimeWindowIntegrationTest { buffer.add(record.value()); }); } - assertThat(receivedUserEventList) - .isEqualTo(sentUserEventList) - .containsExactlyElementsOf(sentUserEventList); + assertThat(receivedUserEventList) + .isEqualTo(sentUserEventList) + .containsExactlyElementsOf(sentUserEventList); } private static void processBuffer(List buffer, List receivedUserEventList) { diff --git a/apache-kafka-2/src/test/java/com/baeldung/kafka/message/ordering/MultiplePartitionIntegrationTest.java b/apache-kafka-2/src/test/java/com/baeldung/kafka/message/ordering/MultiplePartitionIntegrationTest.java index 752514c09a..2fde24114c 100644 --- a/apache-kafka-2/src/test/java/com/baeldung/kafka/message/ordering/MultiplePartitionIntegrationTest.java +++ b/apache-kafka-2/src/test/java/com/baeldung/kafka/message/ordering/MultiplePartitionIntegrationTest.java @@ -11,7 +11,6 @@ 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.clients.producer.RecordMetadata; -import org.apache.kafka.common.KafkaFuture; import org.apache.kafka.common.serialization.LongDeserializer; import org.apache.kafka.common.serialization.LongSerializer; import org.junit.jupiter.api.AfterAll; @@ -25,8 +24,8 @@ import java.time.Duration; import java.util.*; import java.util.concurrent.ExecutionException; import java.util.concurrent.Future; - -import static org.junit.jupiter.api.Assertions.*; +import com.google.common.collect.ImmutableList; +import static org.assertj.core.api.AssertionsForInterfaceTypes.assertThat; @Testcontainers public class MultiplePartitionIntegrationTest { @@ -89,7 +88,7 @@ public class MultiplePartitionIntegrationTest { receivedUserEventList.add(userEvent); System.out.println("User Event ID: " + userEvent.getUserEventId()); }); - assertThat(receivedUserEventList) + assertThat(receivedUserEventList) .isNotEqualTo(sentUserEventList) .containsExactlyInAnyOrderElementsOf(sentUserEventList); } diff --git a/apache-kafka-2/src/test/java/com/baeldung/kafka/message/ordering/SinglePartitionIntegrationTest.java b/apache-kafka-2/src/test/java/com/baeldung/kafka/message/ordering/SinglePartitionIntegrationTest.java index a767133627..0826365f97 100644 --- a/apache-kafka-2/src/test/java/com/baeldung/kafka/message/ordering/SinglePartitionIntegrationTest.java +++ b/apache-kafka-2/src/test/java/com/baeldung/kafka/message/ordering/SinglePartitionIntegrationTest.java @@ -5,7 +5,6 @@ import com.baeldung.kafka.message.ordering.serialization.JacksonDeserializer; import com.baeldung.kafka.message.ordering.serialization.JacksonSerializer; import org.apache.kafka.clients.admin.Admin; import org.apache.kafka.clients.admin.AdminClientConfig; -import org.apache.kafka.clients.admin.CreateTopicsResult; import org.apache.kafka.clients.admin.NewTopic; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.ConsumerRecords; @@ -14,7 +13,6 @@ 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.clients.producer.RecordMetadata; -import org.apache.kafka.common.KafkaFuture; import org.apache.kafka.common.serialization.LongDeserializer; import org.apache.kafka.common.serialization.LongSerializer; import org.junit.jupiter.api.AfterAll; @@ -29,8 +27,8 @@ import java.time.Duration; import java.util.*; import java.util.concurrent.ExecutionException; import java.util.concurrent.Future; - -import static org.junit.jupiter.api.Assertions.assertTrue; +import com.google.common.collect.ImmutableList; +import static org.assertj.core.api.AssertionsForInterfaceTypes.assertThat; @Testcontainers public class SinglePartitionIntegrationTest {