Merge pull request #52 from bellwethr/master
* bellwethr-master: INTEXT-81 Add unit tests INTEXT-81 Fixing bug in decoding message streams * adding a null check on the key decoder
This commit is contained in:
@@ -182,10 +182,11 @@ public class ConsumerConfiguration {
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
public Map<String, List<KafkaStream<byte[], byte[]>>> getConsumerMapWithMessageStreams() {
|
||||
if (consumerMetadata.getValueDecoder() != null) {
|
||||
if (consumerMetadata.getValueDecoder() != null &&
|
||||
consumerMetadata.getKeyDecoder() != null) {
|
||||
return getConsumerConnector().createMessageStreams(
|
||||
consumerMetadata.getTopicStreamMap(),
|
||||
consumerMetadata.getValueDecoder(),
|
||||
consumerMetadata.getKeyDecoder(),
|
||||
consumerMetadata.getValueDecoder());
|
||||
}
|
||||
|
||||
|
||||
@@ -15,15 +15,13 @@
|
||||
*/
|
||||
package org.springframework.integration.kafka.support;
|
||||
|
||||
import org.junit.Assert;
|
||||
import kafka.consumer.ConsumerIterator;
|
||||
import kafka.consumer.KafkaStream;
|
||||
import kafka.javaapi.consumer.ConsumerConnector;
|
||||
import kafka.message.MessageAndMetadata;
|
||||
import org.junit.Test;
|
||||
import org.mockito.Mockito;
|
||||
import org.mockito.invocation.InvocationOnMock;
|
||||
import org.mockito.stubbing.Answer;
|
||||
import static org.junit.Assert.assertNull;
|
||||
import static org.mockito.Mockito.atLeast;
|
||||
import static org.mockito.Mockito.atMost;
|
||||
import static org.mockito.Mockito.mock;
|
||||
import static org.mockito.Mockito.times;
|
||||
import static org.mockito.Mockito.verify;
|
||||
import static org.mockito.Mockito.when;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.Collection;
|
||||
@@ -31,71 +29,83 @@ import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
|
||||
import kafka.consumer.ConsumerIterator;
|
||||
import kafka.consumer.KafkaStream;
|
||||
import kafka.javaapi.consumer.ConsumerConnector;
|
||||
import kafka.message.MessageAndMetadata;
|
||||
import kafka.serializer.Decoder;
|
||||
|
||||
import org.junit.Assert;
|
||||
import org.junit.Test;
|
||||
import org.mockito.invocation.InvocationOnMock;
|
||||
import org.mockito.stubbing.Answer;
|
||||
|
||||
/**
|
||||
* @author Soby Chacko
|
||||
* @author Gunnar Hillert
|
||||
* @since 0.5
|
||||
*/
|
||||
public class ConsumerConfigurationTests {
|
||||
@Test
|
||||
@SuppressWarnings("unchecked")
|
||||
public void testReceiveMessageForSingleTopicFromSingleStream() {
|
||||
final ConsumerMetadata consumerMetadata = Mockito.mock(ConsumerMetadata.class);
|
||||
final ConsumerMetadata consumerMetadata = mock(ConsumerMetadata.class);
|
||||
final ConsumerConnectionProvider consumerConnectionProvider =
|
||||
Mockito.mock(ConsumerConnectionProvider.class);
|
||||
final MessageLeftOverTracker messageLeftOverTracker = Mockito.mock(MessageLeftOverTracker.class);
|
||||
final ConsumerConnector consumerConnector = Mockito.mock(ConsumerConnector.class);
|
||||
mock(ConsumerConnectionProvider.class);
|
||||
final MessageLeftOverTracker messageLeftOverTracker = mock(MessageLeftOverTracker.class);
|
||||
final ConsumerConnector consumerConnector = mock(ConsumerConnector.class);
|
||||
|
||||
Mockito.when(consumerConnectionProvider.getConsumerConnector()).thenReturn(consumerConnector);
|
||||
when(consumerConnectionProvider.getConsumerConnector()).thenReturn(consumerConnector);
|
||||
|
||||
final ConsumerConfiguration consumerConfiguration = new ConsumerConfiguration(consumerMetadata,
|
||||
consumerConnectionProvider, messageLeftOverTracker);
|
||||
consumerConfiguration.setMaxMessages(1);
|
||||
|
||||
final KafkaStream stream = Mockito.mock(KafkaStream.class);
|
||||
final KafkaStream stream = mock(KafkaStream.class);
|
||||
final List<KafkaStream<byte[], byte[]>> streams = new ArrayList<KafkaStream<byte[], byte[]>>();
|
||||
streams.add(stream);
|
||||
final Map<String, List<KafkaStream<byte[], byte[]>>> messageStreams = new HashMap<String, List<KafkaStream<byte[], byte[]>>>();
|
||||
messageStreams.put("topic", streams);
|
||||
|
||||
Mockito.when(consumerConfiguration.getConsumerMapWithMessageStreams()).thenReturn(messageStreams);
|
||||
final ConsumerIterator iterator = Mockito.mock(ConsumerIterator.class);
|
||||
Mockito.when(stream.iterator()).thenReturn(iterator);
|
||||
final MessageAndMetadata messageAndMetadata = Mockito.mock(MessageAndMetadata.class);
|
||||
Mockito.when(iterator.next()).thenReturn(messageAndMetadata);
|
||||
Mockito.when(messageAndMetadata.message()).thenReturn("got message");
|
||||
Mockito.when(messageAndMetadata.topic()).thenReturn("topic");
|
||||
Mockito.when(messageAndMetadata.partition()).thenReturn(1);
|
||||
when(consumerConfiguration.getConsumerMapWithMessageStreams()).thenReturn(messageStreams);
|
||||
final ConsumerIterator iterator = mock(ConsumerIterator.class);
|
||||
when(stream.iterator()).thenReturn(iterator);
|
||||
final MessageAndMetadata messageAndMetadata = mock(MessageAndMetadata.class);
|
||||
when(iterator.next()).thenReturn(messageAndMetadata);
|
||||
when(messageAndMetadata.message()).thenReturn("got message");
|
||||
when(messageAndMetadata.topic()).thenReturn("topic");
|
||||
when(messageAndMetadata.partition()).thenReturn(1);
|
||||
|
||||
final Map<String, Map<Integer, List<Object>>> messages = consumerConfiguration.receive();
|
||||
Assert.assertEquals(messages.size(), 1);
|
||||
Assert.assertEquals(messages.get("topic").size(), 1);
|
||||
Assert.assertEquals(messages.get("topic").get(1).get(0), "got message");
|
||||
|
||||
Mockito.verify(stream, Mockito.times(1)).iterator();
|
||||
Mockito.verify(iterator, Mockito.times(1)).next();
|
||||
Mockito.verify(messageAndMetadata, Mockito.times(1)).message();
|
||||
Mockito.verify(messageAndMetadata, Mockito.times(1)).topic();
|
||||
verify(stream, times(1)).iterator();
|
||||
verify(iterator, times(1)).next();
|
||||
verify(messageAndMetadata, times(1)).message();
|
||||
verify(messageAndMetadata, times(1)).topic();
|
||||
}
|
||||
|
||||
@Test
|
||||
@SuppressWarnings("unchecked")
|
||||
public void testReceiveMessageForSingleTopicFromMultipleStreams() {
|
||||
final ConsumerMetadata consumerMetadata = Mockito.mock(ConsumerMetadata.class);
|
||||
final ConsumerMetadata consumerMetadata = mock(ConsumerMetadata.class);
|
||||
final ConsumerConnectionProvider consumerConnectionProvider =
|
||||
Mockito.mock(ConsumerConnectionProvider.class);
|
||||
final MessageLeftOverTracker messageLeftOverTracker = Mockito.mock(MessageLeftOverTracker.class);
|
||||
mock(ConsumerConnectionProvider.class);
|
||||
final MessageLeftOverTracker messageLeftOverTracker = mock(MessageLeftOverTracker.class);
|
||||
|
||||
final ConsumerConnector consumerConnector = Mockito.mock(ConsumerConnector.class);
|
||||
final ConsumerConnector consumerConnector = mock(ConsumerConnector.class);
|
||||
|
||||
Mockito.when(consumerConnectionProvider.getConsumerConnector()).thenReturn(consumerConnector);
|
||||
when(consumerConnectionProvider.getConsumerConnector()).thenReturn(consumerConnector);
|
||||
|
||||
final ConsumerConfiguration consumerConfiguration = new ConsumerConfiguration(consumerMetadata,
|
||||
consumerConnectionProvider, messageLeftOverTracker);
|
||||
consumerConfiguration.setMaxMessages(3);
|
||||
|
||||
final KafkaStream stream1 = Mockito.mock(KafkaStream.class);
|
||||
final KafkaStream stream2 = Mockito.mock(KafkaStream.class);
|
||||
final KafkaStream stream3 = Mockito.mock(KafkaStream.class);
|
||||
final KafkaStream stream1 = mock(KafkaStream.class);
|
||||
final KafkaStream stream2 = mock(KafkaStream.class);
|
||||
final KafkaStream stream3 = mock(KafkaStream.class);
|
||||
final List<KafkaStream<byte[], byte[]>> streams = new ArrayList<KafkaStream<byte[], byte[]>>();
|
||||
streams.add(stream1);
|
||||
streams.add(stream2);
|
||||
@@ -103,33 +113,33 @@ public class ConsumerConfigurationTests {
|
||||
final Map<String, List<KafkaStream<byte[], byte[]>>> messageStreams = new HashMap<String, List<KafkaStream<byte[], byte[]>>>();
|
||||
messageStreams.put("topic", streams);
|
||||
|
||||
Mockito.when(consumerConfiguration.getConsumerMapWithMessageStreams()).thenReturn(messageStreams);
|
||||
final ConsumerIterator iterator1 = Mockito.mock(ConsumerIterator.class);
|
||||
final ConsumerIterator iterator2 = Mockito.mock(ConsumerIterator.class);
|
||||
final ConsumerIterator iterator3 = Mockito.mock(ConsumerIterator.class);
|
||||
when(consumerConfiguration.getConsumerMapWithMessageStreams()).thenReturn(messageStreams);
|
||||
final ConsumerIterator iterator1 = mock(ConsumerIterator.class);
|
||||
final ConsumerIterator iterator2 = mock(ConsumerIterator.class);
|
||||
final ConsumerIterator iterator3 = mock(ConsumerIterator.class);
|
||||
|
||||
Mockito.when(stream1.iterator()).thenReturn(iterator1);
|
||||
Mockito.when(stream2.iterator()).thenReturn(iterator2);
|
||||
Mockito.when(stream3.iterator()).thenReturn(iterator3);
|
||||
final MessageAndMetadata messageAndMetadata1 = Mockito.mock(MessageAndMetadata.class);
|
||||
final MessageAndMetadata messageAndMetadata2 = Mockito.mock(MessageAndMetadata.class);
|
||||
final MessageAndMetadata messageAndMetadata3 = Mockito.mock(MessageAndMetadata.class);
|
||||
when(stream1.iterator()).thenReturn(iterator1);
|
||||
when(stream2.iterator()).thenReturn(iterator2);
|
||||
when(stream3.iterator()).thenReturn(iterator3);
|
||||
final MessageAndMetadata messageAndMetadata1 = mock(MessageAndMetadata.class);
|
||||
final MessageAndMetadata messageAndMetadata2 = mock(MessageAndMetadata.class);
|
||||
final MessageAndMetadata messageAndMetadata3 = mock(MessageAndMetadata.class);
|
||||
|
||||
Mockito.when(iterator1.next()).thenReturn(messageAndMetadata1);
|
||||
Mockito.when(iterator2.next()).thenReturn(messageAndMetadata2);
|
||||
Mockito.when(iterator3.next()).thenReturn(messageAndMetadata3);
|
||||
when(iterator1.next()).thenReturn(messageAndMetadata1);
|
||||
when(iterator2.next()).thenReturn(messageAndMetadata2);
|
||||
when(iterator3.next()).thenReturn(messageAndMetadata3);
|
||||
|
||||
Mockito.when(messageAndMetadata1.message()).thenReturn("got message");
|
||||
Mockito.when(messageAndMetadata1.topic()).thenReturn("topic");
|
||||
Mockito.when(messageAndMetadata1.partition()).thenReturn(1);
|
||||
when(messageAndMetadata1.message()).thenReturn("got message");
|
||||
when(messageAndMetadata1.topic()).thenReturn("topic");
|
||||
when(messageAndMetadata1.partition()).thenReturn(1);
|
||||
|
||||
Mockito.when(messageAndMetadata2.message()).thenReturn("got message");
|
||||
Mockito.when(messageAndMetadata2.topic()).thenReturn("topic");
|
||||
Mockito.when(messageAndMetadata2.partition()).thenReturn(2);
|
||||
when(messageAndMetadata2.message()).thenReturn("got message");
|
||||
when(messageAndMetadata2.topic()).thenReturn("topic");
|
||||
when(messageAndMetadata2.partition()).thenReturn(2);
|
||||
|
||||
Mockito.when(messageAndMetadata3.message()).thenReturn("got message");
|
||||
Mockito.when(messageAndMetadata3.topic()).thenReturn("topic");
|
||||
Mockito.when(messageAndMetadata3.partition()).thenReturn(3);
|
||||
when(messageAndMetadata3.message()).thenReturn("got message");
|
||||
when(messageAndMetadata3.topic()).thenReturn("topic");
|
||||
when(messageAndMetadata3.partition()).thenReturn(3);
|
||||
|
||||
final Map<String, Map<Integer, List<Object>>> messages = consumerConfiguration.receive();
|
||||
Assert.assertEquals(messages.size(), 1);
|
||||
@@ -147,22 +157,22 @@ public class ConsumerConfigurationTests {
|
||||
@Test
|
||||
@SuppressWarnings("unchecked")
|
||||
public void testReceiveMessageForMultipleTopicsFromMultipleStreams() {
|
||||
final ConsumerMetadata consumerMetadata = Mockito.mock(ConsumerMetadata.class);
|
||||
final ConsumerMetadata consumerMetadata = mock(ConsumerMetadata.class);
|
||||
final ConsumerConnectionProvider consumerConnectionProvider =
|
||||
Mockito.mock(ConsumerConnectionProvider.class);
|
||||
final MessageLeftOverTracker messageLeftOverTracker = Mockito.mock(MessageLeftOverTracker.class);
|
||||
mock(ConsumerConnectionProvider.class);
|
||||
final MessageLeftOverTracker messageLeftOverTracker = mock(MessageLeftOverTracker.class);
|
||||
|
||||
final ConsumerConnector consumerConnector = Mockito.mock(ConsumerConnector.class);
|
||||
final ConsumerConnector consumerConnector = mock(ConsumerConnector.class);
|
||||
|
||||
Mockito.when(consumerConnectionProvider.getConsumerConnector()).thenReturn(consumerConnector);
|
||||
when(consumerConnectionProvider.getConsumerConnector()).thenReturn(consumerConnector);
|
||||
|
||||
final ConsumerConfiguration consumerConfiguration = new ConsumerConfiguration(consumerMetadata,
|
||||
consumerConnectionProvider, messageLeftOverTracker);
|
||||
consumerConfiguration.setMaxMessages(9);
|
||||
|
||||
final KafkaStream stream1 = Mockito.mock(KafkaStream.class);
|
||||
final KafkaStream stream2 = Mockito.mock(KafkaStream.class);
|
||||
final KafkaStream stream3 = Mockito.mock(KafkaStream.class);
|
||||
final KafkaStream stream1 = mock(KafkaStream.class);
|
||||
final KafkaStream stream2 = mock(KafkaStream.class);
|
||||
final KafkaStream stream3 = mock(KafkaStream.class);
|
||||
final List<KafkaStream<byte[], byte[]>> streams = new ArrayList<KafkaStream<byte[], byte[]>>();
|
||||
streams.add(stream1);
|
||||
streams.add(stream2);
|
||||
@@ -172,33 +182,33 @@ public class ConsumerConfigurationTests {
|
||||
messageStreams.put("topic2", streams);
|
||||
messageStreams.put("topic3", streams);
|
||||
|
||||
Mockito.when(consumerConfiguration.getConsumerMapWithMessageStreams()).thenReturn(messageStreams);
|
||||
final ConsumerIterator iterator1 = Mockito.mock(ConsumerIterator.class);
|
||||
final ConsumerIterator iterator2 = Mockito.mock(ConsumerIterator.class);
|
||||
final ConsumerIterator iterator3 = Mockito.mock(ConsumerIterator.class);
|
||||
when(consumerConfiguration.getConsumerMapWithMessageStreams()).thenReturn(messageStreams);
|
||||
final ConsumerIterator iterator1 = mock(ConsumerIterator.class);
|
||||
final ConsumerIterator iterator2 = mock(ConsumerIterator.class);
|
||||
final ConsumerIterator iterator3 = mock(ConsumerIterator.class);
|
||||
|
||||
Mockito.when(stream1.iterator()).thenReturn(iterator1);
|
||||
Mockito.when(stream2.iterator()).thenReturn(iterator2);
|
||||
Mockito.when(stream3.iterator()).thenReturn(iterator3);
|
||||
final MessageAndMetadata messageAndMetadata1 = Mockito.mock(MessageAndMetadata.class);
|
||||
final MessageAndMetadata messageAndMetadata2 = Mockito.mock(MessageAndMetadata.class);
|
||||
final MessageAndMetadata messageAndMetadata3 = Mockito.mock(MessageAndMetadata.class);
|
||||
when(stream1.iterator()).thenReturn(iterator1);
|
||||
when(stream2.iterator()).thenReturn(iterator2);
|
||||
when(stream3.iterator()).thenReturn(iterator3);
|
||||
final MessageAndMetadata messageAndMetadata1 = mock(MessageAndMetadata.class);
|
||||
final MessageAndMetadata messageAndMetadata2 = mock(MessageAndMetadata.class);
|
||||
final MessageAndMetadata messageAndMetadata3 = mock(MessageAndMetadata.class);
|
||||
|
||||
Mockito.when(iterator1.next()).thenReturn(messageAndMetadata1);
|
||||
Mockito.when(iterator2.next()).thenReturn(messageAndMetadata2);
|
||||
Mockito.when(iterator3.next()).thenReturn(messageAndMetadata3);
|
||||
when(iterator1.next()).thenReturn(messageAndMetadata1);
|
||||
when(iterator2.next()).thenReturn(messageAndMetadata2);
|
||||
when(iterator3.next()).thenReturn(messageAndMetadata3);
|
||||
|
||||
Mockito.when(messageAndMetadata1.message()).thenReturn("got message1");
|
||||
Mockito.when(messageAndMetadata1.topic()).thenReturn("topic1");
|
||||
Mockito.when(messageAndMetadata1.partition()).thenAnswer(getAnswer());
|
||||
when(messageAndMetadata1.message()).thenReturn("got message1");
|
||||
when(messageAndMetadata1.topic()).thenReturn("topic1");
|
||||
when(messageAndMetadata1.partition()).thenAnswer(getAnswer());
|
||||
|
||||
Mockito.when(messageAndMetadata2.message()).thenReturn("got message2");
|
||||
Mockito.when(messageAndMetadata2.topic()).thenReturn("topic2");
|
||||
Mockito.when(messageAndMetadata1.partition()).thenAnswer(getAnswer());
|
||||
when(messageAndMetadata2.message()).thenReturn("got message2");
|
||||
when(messageAndMetadata2.topic()).thenReturn("topic2");
|
||||
when(messageAndMetadata1.partition()).thenAnswer(getAnswer());
|
||||
|
||||
Mockito.when(messageAndMetadata3.message()).thenReturn("got message3");
|
||||
Mockito.when(messageAndMetadata3.topic()).thenReturn("topic3");
|
||||
Mockito.when(messageAndMetadata1.partition()).thenAnswer(getAnswer());
|
||||
when(messageAndMetadata3.message()).thenReturn("got message3");
|
||||
when(messageAndMetadata3.topic()).thenReturn("topic3");
|
||||
when(messageAndMetadata1.partition()).thenAnswer(getAnswer());
|
||||
|
||||
final Map<String, Map<Integer, List<Object>>> messages = consumerConfiguration.receive();
|
||||
int sum = 0;
|
||||
@@ -234,13 +244,13 @@ public class ConsumerConfigurationTests {
|
||||
@Test
|
||||
@SuppressWarnings("unchecked")
|
||||
public void testReceiveMessageAndVerifyMessageLeftoverFromPreviousPollAreTakenFirst() {
|
||||
final ConsumerMetadata consumerMetadata = Mockito.mock(ConsumerMetadata.class);
|
||||
final ConsumerMetadata consumerMetadata = mock(ConsumerMetadata.class);
|
||||
final ConsumerConnectionProvider consumerConnectionProvider =
|
||||
Mockito.mock(ConsumerConnectionProvider.class);
|
||||
final MessageLeftOverTracker messageLeftOverTracker = Mockito.mock(MessageLeftOverTracker.class);
|
||||
final ConsumerConnector consumerConnector = Mockito.mock(ConsumerConnector.class);
|
||||
mock(ConsumerConnectionProvider.class);
|
||||
final MessageLeftOverTracker messageLeftOverTracker = mock(MessageLeftOverTracker.class);
|
||||
final ConsumerConnector consumerConnector = mock(ConsumerConnector.class);
|
||||
|
||||
Mockito.when(messageLeftOverTracker.getCurrentCount()).thenReturn(3);
|
||||
when(messageLeftOverTracker.getCurrentCount()).thenReturn(3);
|
||||
final MessageAndMetadata m1 = new MessageAndMetadata("key1", "value1", "topic1", 1, 1L);
|
||||
final MessageAndMetadata m2 = new MessageAndMetadata("key2", "value2", "topic2", 1, 1L);
|
||||
final MessageAndMetadata m3 = new MessageAndMetadata("key1", "value3", "topic3", 1, 1L);
|
||||
@@ -250,28 +260,28 @@ public class ConsumerConfigurationTests {
|
||||
mList.add(m2);
|
||||
mList.add(m3);
|
||||
|
||||
Mockito.when(messageLeftOverTracker.getMessageLeftOverFromPreviousPoll()).thenReturn(mList);
|
||||
when(messageLeftOverTracker.getMessageLeftOverFromPreviousPoll()).thenReturn(mList);
|
||||
|
||||
Mockito.when(consumerConnectionProvider.getConsumerConnector()).thenReturn(consumerConnector);
|
||||
when(consumerConnectionProvider.getConsumerConnector()).thenReturn(consumerConnector);
|
||||
|
||||
final ConsumerConfiguration consumerConfiguration = new ConsumerConfiguration(consumerMetadata,
|
||||
consumerConnectionProvider, messageLeftOverTracker);
|
||||
consumerConfiguration.setMaxMessages(5);
|
||||
|
||||
final KafkaStream stream = Mockito.mock(KafkaStream.class);
|
||||
final KafkaStream stream = mock(KafkaStream.class);
|
||||
final List<KafkaStream<byte[], byte[]>> streams = new ArrayList<KafkaStream<byte[], byte[]>>();
|
||||
streams.add(stream);
|
||||
final Map<String, List<KafkaStream<byte[], byte[]>>> messageStreams = new HashMap<String, List<KafkaStream<byte[], byte[]>>>();
|
||||
messageStreams.put("topic1", streams);
|
||||
|
||||
Mockito.when(consumerConfiguration.getConsumerMapWithMessageStreams()).thenReturn(messageStreams);
|
||||
final ConsumerIterator iterator = Mockito.mock(ConsumerIterator.class);
|
||||
Mockito.when(stream.iterator()).thenReturn(iterator);
|
||||
final MessageAndMetadata messageAndMetadata = Mockito.mock(MessageAndMetadata.class);
|
||||
Mockito.when(iterator.next()).thenReturn(messageAndMetadata);
|
||||
Mockito.when(messageAndMetadata.message()).thenReturn("got message");
|
||||
Mockito.when(messageAndMetadata.topic()).thenReturn("topic1");
|
||||
Mockito.when(messageAndMetadata.partition()).thenReturn(1);
|
||||
when(consumerConfiguration.getConsumerMapWithMessageStreams()).thenReturn(messageStreams);
|
||||
final ConsumerIterator<String, String> iterator = mock(ConsumerIterator.class);
|
||||
when(stream.iterator()).thenReturn(iterator);
|
||||
final MessageAndMetadata<String, String> messageAndMetadata = mock(MessageAndMetadata.class);
|
||||
when(iterator.next()).thenReturn(messageAndMetadata);
|
||||
when(messageAndMetadata.message()).thenReturn("got message");
|
||||
when(messageAndMetadata.topic()).thenReturn("topic1");
|
||||
when(messageAndMetadata.partition()).thenReturn(1);
|
||||
|
||||
final Map<String, Map<Integer, List<Object>>> messages = consumerConfiguration.receive();
|
||||
int sum = 0;
|
||||
@@ -295,6 +305,74 @@ public class ConsumerConfigurationTests {
|
||||
Assert.assertTrue(valueFound(messages.get("topic3").get(1), "value3"));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testGetConsumerMapWithMessageStreamsWithNullDecoders() {
|
||||
|
||||
final ConsumerMetadata<?,?> mockedConsumerMetadata = mock(ConsumerMetadata.class);
|
||||
|
||||
assertNull(mockedConsumerMetadata.getKeyDecoder());
|
||||
assertNull(mockedConsumerMetadata.getValueDecoder());
|
||||
|
||||
final Map<String, Integer> topicsStreamMap = new HashMap<String, Integer>();
|
||||
when(mockedConsumerMetadata.getTopicStreamMap()).thenReturn(topicsStreamMap);
|
||||
|
||||
final ConsumerConnectionProvider mockedConsumerConnectionProvider = mock(ConsumerConnectionProvider.class);
|
||||
final MessageLeftOverTracker mockedMessageLeftOverTracker = mock(MessageLeftOverTracker.class);
|
||||
final ConsumerConnector mockedConsumerConnector = mock(ConsumerConnector.class);
|
||||
|
||||
when(mockedConsumerConnectionProvider.getConsumerConnector()).thenReturn(mockedConsumerConnector);
|
||||
|
||||
final Map<String, List<KafkaStream<byte[], byte[]>>> messageStreams = new HashMap<String, List<KafkaStream<byte[],byte[]>>>();
|
||||
when(mockedConsumerConnector.createMessageStreams(topicsStreamMap)).thenReturn(messageStreams);
|
||||
|
||||
final ConsumerConfiguration consumerConfiguration = new ConsumerConfiguration(mockedConsumerMetadata,
|
||||
mockedConsumerConnectionProvider, mockedMessageLeftOverTracker);
|
||||
|
||||
consumerConfiguration.getConsumerMapWithMessageStreams();
|
||||
|
||||
verify(mockedConsumerMetadata, atLeast(1)).getTopicStreamMap();
|
||||
verify(mockedConsumerConnector, atLeast(1)).createMessageStreams(topicsStreamMap);
|
||||
verify(mockedConsumerConnector, atMost(0)).createMessageStreams(topicsStreamMap, null, null);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testGetConsumerMapWithMessageStreamsWithDecoders() {
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
final ConsumerMetadata<String, String> mockedConsumerMetadata = mock(ConsumerMetadata.class);
|
||||
|
||||
final Map<String, Integer> topicsStreamMap = new HashMap<String, Integer>();
|
||||
when(mockedConsumerMetadata.getTopicStreamMap()).thenReturn(topicsStreamMap);
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
final Decoder<String> mockedKeyDecoder = mock(Decoder.class);
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
final Decoder<String> mockedValueDecoder = mock(Decoder.class);
|
||||
|
||||
when(mockedConsumerMetadata.getKeyDecoder()).thenReturn(mockedKeyDecoder);
|
||||
when(mockedConsumerMetadata.getValueDecoder()).thenReturn(mockedValueDecoder);
|
||||
|
||||
final ConsumerConnectionProvider mockedConsumerConnectionProvider =
|
||||
mock(ConsumerConnectionProvider.class);
|
||||
final MessageLeftOverTracker mockedMessageLeftOverTracker = mock(MessageLeftOverTracker.class);
|
||||
final ConsumerConnector mockedConsumerConnector = mock(ConsumerConnector.class);
|
||||
|
||||
when(mockedConsumerConnectionProvider.getConsumerConnector()).thenReturn(mockedConsumerConnector);
|
||||
|
||||
final Map<String, List<KafkaStream<byte[], byte[]>>> messageStreams = new HashMap<String, List<KafkaStream<byte[],byte[]>>>();
|
||||
when(mockedConsumerConnector.createMessageStreams(topicsStreamMap)).thenReturn(messageStreams);
|
||||
|
||||
final ConsumerConfiguration consumerConfiguration = new ConsumerConfiguration(mockedConsumerMetadata,
|
||||
mockedConsumerConnectionProvider, mockedMessageLeftOverTracker);
|
||||
|
||||
consumerConfiguration.getConsumerMapWithMessageStreams();
|
||||
|
||||
verify(mockedConsumerMetadata, atLeast(1)).getTopicStreamMap();
|
||||
verify(mockedConsumerConnector, atMost(0)).createMessageStreams(topicsStreamMap);
|
||||
verify(mockedConsumerConnector, atLeast(1)).createMessageStreams(topicsStreamMap, mockedKeyDecoder, mockedValueDecoder);
|
||||
}
|
||||
|
||||
private boolean valueFound(final List<Object> l, final String value){
|
||||
for (final Object o : l){
|
||||
if (value.equals(o)){
|
||||
|
||||
Reference in New Issue
Block a user