Skip to content

Commit 021f55d

Browse files
committed
Make Kafka producer/consumer testable
Refactor KafkaMessageProducer and KafkaMessageConsumer to depend on the Producer/Consumer interfaces and add constructors that accept mockable instances. Extract default producer/consumer creation into factory methods so tests can inject MockProducer/MockConsumer. Update tests across the microservices-messaging module to use MockProducer/MockConsumer, add more meaningful assertions, error/interrupt handling tests, and simplify AppTest. These changes improve unit testability and remove the need for a running Kafka instance while preserving runtime behavior.
1 parent 4c08cdc commit 021f55d

9 files changed

Lines changed: 194 additions & 98 deletions

File tree

microservices-messaging/src/main/java/com/iluwatar/messaging/KafkaMessageConsumer.java

Lines changed: 24 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -30,7 +30,7 @@
3030
import java.util.Collections;
3131
import java.util.Properties;
3232
import java.util.concurrent.atomic.AtomicBoolean;
33-
import java.util.function.Consumer;
33+
import org.apache.kafka.clients.consumer.Consumer;
3434
import org.apache.kafka.clients.consumer.ConsumerConfig;
3535
import org.apache.kafka.clients.consumer.ConsumerRecords;
3636
import org.apache.kafka.clients.consumer.KafkaConsumer;
@@ -41,10 +41,10 @@
4141
/** Kafka message consumer that subscribes to topics and processes messages. */
4242
public class KafkaMessageConsumer implements AutoCloseable, Runnable {
4343
private static final Logger LOGGER = LoggerFactory.getLogger(KafkaMessageConsumer.class);
44-
private final KafkaConsumer<String, String> consumer;
44+
private final Consumer<String, String> consumer;
4545
private final ObjectMapper objectMapper;
4646
private final String topic;
47-
private final Consumer<Message> messageHandler;
47+
private final java.util.function.Consumer<Message> messageHandler;
4848
private final AtomicBoolean running = new AtomicBoolean(true);
4949

5050
/**
@@ -56,20 +56,34 @@ public class KafkaMessageConsumer implements AutoCloseable, Runnable {
5656
* @param messageHandler handler for received messages
5757
*/
5858
public KafkaMessageConsumer(
59-
String bootstrapServers, String groupId, String topic, Consumer<Message> messageHandler) {
59+
String bootstrapServers,
60+
String groupId,
61+
String topic,
62+
java.util.function.Consumer<Message> messageHandler) {
63+
this(createDefaultConsumer(bootstrapServers, groupId), topic, messageHandler);
64+
}
65+
66+
KafkaMessageConsumer(
67+
Consumer<String, String> consumer,
68+
String topic,
69+
java.util.function.Consumer<Message> messageHandler) {
70+
this.consumer = consumer;
71+
this.objectMapper = new ObjectMapper();
72+
this.objectMapper.registerModule(new JavaTimeModule());
73+
this.topic = topic;
74+
this.messageHandler = messageHandler;
75+
}
76+
77+
private static Consumer<String, String> createDefaultConsumer(
78+
String bootstrapServers, String groupId) {
6079
Properties props = new Properties();
6180
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
6281
props.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
6382
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
6483
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
6584
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
6685
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "true");
67-
68-
this.consumer = new KafkaConsumer<>(props);
69-
this.objectMapper = new ObjectMapper();
70-
this.objectMapper.registerModule(new JavaTimeModule());
71-
this.topic = topic;
72-
this.messageHandler = messageHandler;
86+
return new KafkaConsumer<>(props);
7387
}
7488

7589
@Override

microservices-messaging/src/main/java/com/iluwatar/messaging/KafkaMessageProducer.java

Lines changed: 13 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,7 @@
2828
import com.fasterxml.jackson.datatype.jsr310.JavaTimeModule;
2929
import java.util.Properties;
3030
import org.apache.kafka.clients.producer.KafkaProducer;
31+
import org.apache.kafka.clients.producer.Producer;
3132
import org.apache.kafka.clients.producer.ProducerConfig;
3233
import org.apache.kafka.clients.producer.ProducerRecord;
3334
import org.apache.kafka.common.serialization.StringSerializer;
@@ -37,7 +38,7 @@
3738
/** Kafka message producer that publishes messages to Kafka topics. */
3839
public class KafkaMessageProducer implements AutoCloseable {
3940
private static final Logger LOGGER = LoggerFactory.getLogger(KafkaMessageProducer.class);
40-
private final KafkaProducer<String, String> producer;
41+
private final Producer<String, String> producer;
4142
private final ObjectMapper objectMapper;
4243

4344
/**
@@ -46,16 +47,23 @@ public class KafkaMessageProducer implements AutoCloseable {
4647
* @param bootstrapServers Kafka bootstrap servers
4748
*/
4849
public KafkaMessageProducer(String bootstrapServers) {
50+
this(createDefaultProducer(bootstrapServers));
51+
}
52+
53+
KafkaMessageProducer(Producer<String, String> producer) {
54+
this.producer = producer;
55+
this.objectMapper = new ObjectMapper();
56+
this.objectMapper.registerModule(new JavaTimeModule());
57+
}
58+
59+
private static Producer<String, String> createDefaultProducer(String bootstrapServers) {
4960
Properties props = new Properties();
5061
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
5162
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
5263
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
5364
props.put(ProducerConfig.ACKS_CONFIG, "all");
5465
props.put(ProducerConfig.RETRIES_CONFIG, 3);
55-
56-
this.producer = new KafkaProducer<>(props);
57-
this.objectMapper = new ObjectMapper();
58-
this.objectMapper.registerModule(new JavaTimeModule());
66+
return new KafkaProducer<>(props);
5967
}
6068

6169
/**

microservices-messaging/src/test/java/com/iluwatar/messaging/AppTest.java

Lines changed: 2 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -32,17 +32,7 @@
3232
class AppTest {
3333

3434
@Test
35-
void testMainMethodDoesNotThrowException() {
36-
// Note: This test requires a running Kafka instance
37-
// In a real scenario, we would use embedded Kafka for testing
38-
// For now, we just verify the method can be called without compilation errors
39-
40-
// Act & Assert
41-
assertDoesNotThrow(
42-
() -> {
43-
// Main method requires Kafka to be running, so we don't actually call it in unit tests
44-
// This is a placeholder to ensure the class structure is correct
45-
},
46-
"App should be instantiable");
35+
void testAppConstructor() {
36+
assertDoesNotThrow(() -> new App());
4737
}
4838
}

microservices-messaging/src/test/java/com/iluwatar/messaging/InventoryServiceTest.java

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -104,4 +104,15 @@ void testHandleMultipleMessages() {
104104
},
105105
"Should handle multiple messages without error");
106106
}
107+
108+
@Test
109+
void testHandleMessagesWhenInterrupted() {
110+
Thread.currentThread().interrupt();
111+
inventoryService.handleMessage(new Message("Order Created: ORDER-001"));
112+
org.junit.jupiter.api.Assertions.assertTrue(Thread.interrupted());
113+
114+
Thread.currentThread().interrupt();
115+
inventoryService.handleMessage(new Message("Order Cancelled: ORDER-001"));
116+
org.junit.jupiter.api.Assertions.assertTrue(Thread.interrupted());
117+
}
107118
}

microservices-messaging/src/test/java/com/iluwatar/messaging/KafkaMessageConsumerTest.java

Lines changed: 67 additions & 32 deletions
Original file line numberDiff line numberDiff line change
@@ -26,57 +26,92 @@
2626

2727
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
2828
import static org.junit.jupiter.api.Assertions.assertNotNull;
29+
import static org.junit.jupiter.api.Assertions.assertTrue;
2930

31+
import com.fasterxml.jackson.databind.ObjectMapper;
32+
import com.fasterxml.jackson.datatype.jsr310.JavaTimeModule;
33+
import java.util.Collections;
34+
import java.util.HashMap;
35+
import java.util.concurrent.atomic.AtomicBoolean;
36+
import org.apache.kafka.clients.consumer.ConsumerRecord;
37+
import org.apache.kafka.clients.consumer.MockConsumer;
38+
import org.apache.kafka.clients.consumer.OffsetResetStrategy;
39+
import org.apache.kafka.common.TopicPartition;
40+
import org.junit.jupiter.api.BeforeEach;
3041
import org.junit.jupiter.api.Test;
3142

3243
/**
33-
* Unit tests for {@link KafkaMessageConsumer}. Note: These tests verify basic functionality without
34-
* requiring a Kafka instance. For integration tests with Kafka, use embedded Kafka or
35-
* testcontainers.
44+
* Unit tests for {@link KafkaMessageConsumer}.
3645
*/
3746
class KafkaMessageConsumerTest {
3847

48+
private MockConsumer<String, String> mockConsumer;
49+
private KafkaMessageConsumer kafkaMessageConsumer;
50+
private AtomicBoolean handlerCalled;
51+
52+
@BeforeEach
53+
void setUp() {
54+
mockConsumer = new MockConsumer<>(OffsetResetStrategy.EARLIEST);
55+
handlerCalled = new AtomicBoolean(false);
56+
kafkaMessageConsumer = new KafkaMessageConsumer(
57+
mockConsumer,
58+
"test-topic",
59+
msg -> handlerCalled.set(true));
60+
}
61+
3962
@Test
4063
void testConsumerCanBeInstantiated() {
41-
// Arrange & Act & Assert
42-
// Note: Don't actually create consumer in unit test as it requires Kafka
43-
assertNotNull(KafkaMessageConsumer.class, "KafkaMessageConsumer class should exist");
64+
assertNotNull(kafkaMessageConsumer, "KafkaMessageConsumer should be instantiated");
4465
}
4566

4667
@Test
47-
void testConsumerImplementsRunnable() {
48-
// Arrange & Act & Assert
49-
var interfaces = KafkaMessageConsumer.class.getInterfaces();
50-
for (var i : interfaces) {
51-
if (i.equals(Runnable.class)) {
52-
break;
68+
void testRunProcessesValidMessageAndStops() throws Exception {
69+
ObjectMapper mapper = new ObjectMapper().registerModule(new JavaTimeModule());
70+
Message msg = new Message("Order Created: 123");
71+
String jsonStr = mapper.writeValueAsString(msg);
72+
73+
TopicPartition tp = new TopicPartition("test-topic", 0);
74+
mockConsumer.updateBeginningOffsets(new HashMap<>() {
75+
{
76+
put(tp, 0L);
5377
}
54-
}
55-
assertNotNull(interfaces, "Should have interfaces");
56-
// Note: Runnable is implemented for threading
78+
});
79+
80+
mockConsumer.schedulePollTask(() -> {
81+
mockConsumer.rebalance(Collections.singletonList(tp));
82+
mockConsumer.addRecord(new ConsumerRecord<>("test-topic", 0, 0L, "key", jsonStr));
83+
});
84+
85+
mockConsumer.schedulePollTask(() -> kafkaMessageConsumer.stop());
86+
87+
kafkaMessageConsumer.run();
88+
89+
assertTrue(handlerCalled.get(), "Handler should have been invoked");
90+
assertTrue(mockConsumer.closed(), "Consumer should be closed");
5791
}
5892

5993
@Test
60-
void testConsumerImplementsAutoCloseable() {
61-
// Arrange & Act & Assert
62-
var interfaces = KafkaMessageConsumer.class.getInterfaces();
63-
for (var i : interfaces) {
64-
if (i.equals(AutoCloseable.class)) {
65-
break;
94+
void testRunHandlesInvalidJsonMessage() {
95+
TopicPartition tp = new TopicPartition("test-topic", 0);
96+
mockConsumer.updateBeginningOffsets(new HashMap<>() {
97+
{
98+
put(tp, 0L);
6699
}
67-
}
68-
assertNotNull(interfaces, "Should have interfaces");
69-
// Note: AutoCloseable is implemented
100+
});
101+
102+
mockConsumer.schedulePollTask(() -> {
103+
mockConsumer.rebalance(Collections.singletonList(tp));
104+
mockConsumer.addRecord(new ConsumerRecord<>("test-topic", 0, 0L, "key", "{invalid json"));
105+
});
106+
107+
mockConsumer.schedulePollTask(() -> kafkaMessageConsumer.stop());
108+
109+
assertDoesNotThrow(() -> kafkaMessageConsumer.run());
110+
assertTrue(mockConsumer.closed(), "Consumer should be closed");
70111
}
71112

72113
@Test
73-
void testConsumerClassHasStopMethod() {
74-
// Arrange & Act & Assert
75-
assertDoesNotThrow(
76-
() -> {
77-
var method = KafkaMessageConsumer.class.getDeclaredMethod("stop");
78-
assertNotNull(method, "stop method should exist");
79-
},
80-
"KafkaMessageConsumer should have stop method");
114+
void testCloseStopsConsumer() {
115+
assertDoesNotThrow(() -> kafkaMessageConsumer.close());
81116
}
82117
}

microservices-messaging/src/test/java/com/iluwatar/messaging/KafkaMessageProducerTest.java

Lines changed: 38 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -24,48 +24,61 @@
2424
*/
2525
package com.iluwatar.messaging;
2626

27+
import static org.junit.jupiter.api.Assertions.assertEquals;
2728
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
2829
import static org.junit.jupiter.api.Assertions.assertNotNull;
30+
import static org.junit.jupiter.api.Assertions.assertTrue;
2931

32+
import org.apache.kafka.clients.producer.MockProducer;
33+
import org.apache.kafka.common.serialization.StringSerializer;
34+
import org.junit.jupiter.api.BeforeEach;
3035
import org.junit.jupiter.api.Test;
3136

3237
/**
33-
* Unit tests for {@link KafkaMessageProducer}. Note: These tests verify basic functionality without
34-
* requiring a Kafka instance. For integration tests with Kafka, use embedded Kafka or
35-
* testcontainers.
38+
* Unit tests for {@link KafkaMessageProducer}.
3639
*/
3740
class KafkaMessageProducerTest {
3841

42+
private MockProducer<String, String> mockProducer;
43+
private KafkaMessageProducer kafkaMessageProducer;
44+
45+
@BeforeEach
46+
void setUp() {
47+
mockProducer = new MockProducer<>(true, new StringSerializer(), new StringSerializer());
48+
kafkaMessageProducer = new KafkaMessageProducer(mockProducer);
49+
}
50+
3951
@Test
4052
void testProducerCanBeInstantiated() {
41-
// Arrange & Act & Assert
42-
// Note: Don't actually create producer in unit test as it requires Kafka
43-
// This test just verifies the class structure is correct
44-
assertNotNull(KafkaMessageProducer.class, "KafkaMessageProducer class should exist");
53+
assertNotNull(kafkaMessageProducer, "KafkaMessageProducer should be instantiated");
4554
}
4655

4756
@Test
48-
void testProducerClassHasPublishMethod() {
49-
// Arrange & Act & Assert
50-
assertDoesNotThrow(
51-
() -> {
52-
var method =
53-
KafkaMessageProducer.class.getDeclaredMethod("publish", String.class, Message.class);
54-
assertNotNull(method, "publish method should exist");
55-
},
56-
"KafkaMessageProducer should have publish method");
57+
void testPublishMessageSuccess() {
58+
Message message = new Message("Test Order");
59+
60+
assertDoesNotThrow(() -> kafkaMessageProducer.publish("test-topic", message));
61+
assertEquals(1, mockProducer.history().size());
62+
assertEquals("test-topic", mockProducer.history().get(0).topic());
63+
assertEquals(message.getId(), mockProducer.history().get(0).key());
5764
}
5865

5966
@Test
60-
void testProducerImplementsAutoCloseable() {
61-
// Arrange & Act & Assert
62-
var interfaces = KafkaMessageProducer.class.getInterfaces();
63-
for (var i : interfaces) {
64-
if (i.equals(AutoCloseable.class)) {
65-
break;
66-
}
67+
void testPublishMessageErrorCallback() {
68+
MockProducer<String, String> failingProducer =
69+
new MockProducer<>(false, new StringSerializer(), new StringSerializer());
70+
try (KafkaMessageProducer producerWithError = new KafkaMessageProducer(failingProducer)) {
71+
Message message = new Message("Test Order");
72+
73+
assertDoesNotThrow(() -> producerWithError.publish("test-topic", message));
74+
assertEquals(1, failingProducer.history().size());
75+
failingProducer.errorNext(new RuntimeException("Kafka publish error"));
6776
}
68-
assertNotNull(interfaces, "Should have interfaces");
69-
// Note: AutoCloseable is implemented
77+
}
78+
79+
@Test
80+
void testClose() {
81+
assertDoesNotThrow(() -> kafkaMessageProducer.close());
82+
assertTrue(mockProducer.closed());
7083
}
7184
}

microservices-messaging/src/test/java/com/iluwatar/messaging/NotificationServiceTest.java

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -116,4 +116,19 @@ void testHandleMultipleOrdersSequentially() {
116116
},
117117
"Should handle multiple orders sequentially without error");
118118
}
119+
120+
@Test
121+
void testHandleMessagesWhenInterrupted() {
122+
Thread.currentThread().interrupt();
123+
notificationService.handleMessage(new Message("Order Created: ORDER-001"));
124+
org.junit.jupiter.api.Assertions.assertTrue(Thread.interrupted());
125+
126+
Thread.currentThread().interrupt();
127+
notificationService.handleMessage(new Message("Order Updated: ORDER-001"));
128+
org.junit.jupiter.api.Assertions.assertTrue(Thread.interrupted());
129+
130+
Thread.currentThread().interrupt();
131+
notificationService.handleMessage(new Message("Order Cancelled: ORDER-001"));
132+
org.junit.jupiter.api.Assertions.assertTrue(Thread.interrupted());
133+
}
119134
}

0 commit comments

Comments
 (0)