|
|
|
|
@@ -37,6 +37,7 @@ import org.mockito.stubbing.Answer;
|
|
|
|
|
* @since 0.5
|
|
|
|
|
*/
|
|
|
|
|
public class ConsumerConfigurationTests<K,V> {
|
|
|
|
|
|
|
|
|
|
@Test
|
|
|
|
|
@SuppressWarnings("unchecked")
|
|
|
|
|
public void testReceiveMessageForSingleTopicFromSingleStream() {
|
|
|
|
|
@@ -72,9 +73,9 @@ public class ConsumerConfigurationTests<K,V> {
|
|
|
|
|
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");
|
|
|
|
|
Assert.assertEquals(1, messages.size());
|
|
|
|
|
Assert.assertEquals(1, messages.get("topic").size());
|
|
|
|
|
Assert.assertEquals("got message", messages.get("topic").get(1).get(0));
|
|
|
|
|
|
|
|
|
|
verify(stream, times(1)).iterator();
|
|
|
|
|
verify(iterator, times(1)).next();
|
|
|
|
|
@@ -150,7 +151,7 @@ public class ConsumerConfigurationTests<K,V> {
|
|
|
|
|
sum += l.size();
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
Assert.assertEquals(sum, 3);
|
|
|
|
|
Assert.assertEquals(3, sum);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@Test
|
|
|
|
|
@@ -226,7 +227,7 @@ public class ConsumerConfigurationTests<K,V> {
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
Assert.assertEquals(sum, 9);
|
|
|
|
|
Assert.assertEquals(9, sum);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
@@ -284,11 +285,11 @@ public class ConsumerConfigurationTests<K,V> {
|
|
|
|
|
final Map<String, List<KafkaStream<K,V>>> messageStreams = new HashMap<String, List<KafkaStream<K,V>>>();
|
|
|
|
|
messageStreams.put("topic1", streams);
|
|
|
|
|
when(consumerConfiguration.createMessageStreamsForTopic()).thenReturn(messageStreams);
|
|
|
|
|
final ConsumerIterator iterator = mock(ConsumerIterator.class);
|
|
|
|
|
final ConsumerIterator<K, V> iterator = mock(ConsumerIterator.class);
|
|
|
|
|
when(stream.iterator()).thenReturn(iterator);
|
|
|
|
|
final MessageAndMetadata messageAndMetadata = mock(MessageAndMetadata.class);
|
|
|
|
|
final MessageAndMetadata<K, V> messageAndMetadata = mock(MessageAndMetadata.class);
|
|
|
|
|
when(iterator.next()).thenReturn(messageAndMetadata);
|
|
|
|
|
when(messageAndMetadata.message()).thenReturn("got message");
|
|
|
|
|
when(messageAndMetadata.message()).thenReturn((V) "got message");
|
|
|
|
|
when(messageAndMetadata.topic()).thenReturn("topic1");
|
|
|
|
|
when(messageAndMetadata.partition()).thenReturn(1);
|
|
|
|
|
|
|
|
|
|
@@ -303,7 +304,7 @@ public class ConsumerConfigurationTests<K,V> {
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
}
|
|
|
|
|
Assert.assertEquals(sum, 5);
|
|
|
|
|
Assert.assertEquals(5, sum);
|
|
|
|
|
|
|
|
|
|
Assert.assertTrue(messages.containsKey("topic1"));
|
|
|
|
|
Assert.assertTrue(messages.containsKey("topic2"));
|
|
|
|
|
@@ -333,7 +334,8 @@ public class ConsumerConfigurationTests<K,V> {
|
|
|
|
|
when(mockedConsumerConnectionProvider.getConsumerConnector()).thenReturn(mockedConsumerConnector);
|
|
|
|
|
|
|
|
|
|
final Map<String, List<KafkaStream<K,V>>> messageStreams = new HashMap<String, List<KafkaStream<K,V>>>();
|
|
|
|
|
when((Map<String, List<KafkaStream<K,V>>>) (Object) mockedConsumerConnector.createMessageStreams(topicsStreamMap)).thenReturn(messageStreams);
|
|
|
|
|
when((Map<String, List<KafkaStream<K,V>>>)
|
|
|
|
|
(Object) mockedConsumerConnector.createMessageStreams(topicsStreamMap)).thenReturn(messageStreams);
|
|
|
|
|
|
|
|
|
|
final ConsumerConfiguration<K,V> consumerConfiguration = new ConsumerConfiguration<K,V>(mockedConsumerMetadata,
|
|
|
|
|
mockedConsumerConnectionProvider, mockedMessageLeftOverTracker);
|
|
|
|
|
@@ -374,51 +376,54 @@ public class ConsumerConfigurationTests<K,V> {
|
|
|
|
|
final Map<String, List<KafkaStream<byte[], byte[]>>> messageStreams = new HashMap<String, List<KafkaStream<byte[],byte[]>>>();
|
|
|
|
|
when(mockedConsumerConnector.createMessageStreams(topicsStreamMap)).thenReturn(messageStreams);
|
|
|
|
|
|
|
|
|
|
final ConsumerConfiguration<String, String> consumerConfiguration = new ConsumerConfiguration<String, String>(mockedConsumerMetadata,
|
|
|
|
|
mockedConsumerConnectionProvider, mockedMessageLeftOverTracker);
|
|
|
|
|
final ConsumerConfiguration<String, String> consumerConfiguration =
|
|
|
|
|
new ConsumerConfiguration<String, String>(mockedConsumerMetadata, mockedConsumerConnectionProvider,
|
|
|
|
|
mockedMessageLeftOverTracker);
|
|
|
|
|
|
|
|
|
|
consumerConfiguration.createMessageStreamsForTopic();
|
|
|
|
|
|
|
|
|
|
verify(mockedConsumerMetadata, atLeast(1)).getTopicStreamMap();
|
|
|
|
|
verify(mockedConsumerConnector, atMost(0)).createMessageStreams(topicsStreamMap);
|
|
|
|
|
verify(mockedConsumerConnector, atLeast(1)).createMessageStreams(topicsStreamMap, mockedKeyDecoder, mockedValueDecoder);
|
|
|
|
|
verify(mockedConsumerConnector, atLeast(1))
|
|
|
|
|
.createMessageStreams(topicsStreamMap, mockedKeyDecoder, mockedValueDecoder);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
@Test
|
|
|
|
|
@SuppressWarnings("unchecked")
|
|
|
|
|
public void testReceiveMessageForTopicFilterFromSingleStream() {
|
|
|
|
|
final ConsumerMetadata consumerMetadata = mock(ConsumerMetadata.class);
|
|
|
|
|
final ConsumerMetadata<byte[], String> consumerMetadata = mock(ConsumerMetadata.class);
|
|
|
|
|
final ConsumerConnectionProvider consumerConnectionProvider =
|
|
|
|
|
mock(ConsumerConnectionProvider.class);
|
|
|
|
|
final MessageLeftOverTracker messageLeftOverTracker = mock(MessageLeftOverTracker.class);
|
|
|
|
|
final MessageLeftOverTracker<byte[], String> messageLeftOverTracker = mock(MessageLeftOverTracker.class);
|
|
|
|
|
final ConsumerConnector consumerConnector = mock(ConsumerConnector.class);
|
|
|
|
|
|
|
|
|
|
when(consumerMetadata.getTopicFilterConfiguration()).thenReturn(new TopicFilterConfiguration(".*", 1, false));
|
|
|
|
|
|
|
|
|
|
when(consumerConnectionProvider.getConsumerConnector()).thenReturn(consumerConnector);
|
|
|
|
|
|
|
|
|
|
final ConsumerConfiguration consumerConfiguration = new ConsumerConfiguration(consumerMetadata,
|
|
|
|
|
consumerConnectionProvider, messageLeftOverTracker);
|
|
|
|
|
final ConsumerConfiguration<byte[], String> consumerConfiguration =
|
|
|
|
|
new ConsumerConfiguration<byte[], String>(consumerMetadata, consumerConnectionProvider,
|
|
|
|
|
messageLeftOverTracker);
|
|
|
|
|
consumerConfiguration.setMaxMessages(1);
|
|
|
|
|
|
|
|
|
|
final KafkaStream stream = mock(KafkaStream.class);
|
|
|
|
|
final List<KafkaStream<byte[], byte[]>> streams = new ArrayList<KafkaStream<byte[], byte[]>>();
|
|
|
|
|
final KafkaStream<byte[],String> stream = mock(KafkaStream.class);
|
|
|
|
|
final List<KafkaStream<byte[], String>> streams = new ArrayList<KafkaStream<byte[], String>>();
|
|
|
|
|
streams.add(stream);
|
|
|
|
|
|
|
|
|
|
when(consumerConfiguration.createMessageStreamsForTopicFilter()).thenReturn(streams);
|
|
|
|
|
final ConsumerIterator iterator = mock(ConsumerIterator.class);
|
|
|
|
|
final ConsumerIterator<byte[], String> iterator = mock(ConsumerIterator.class);
|
|
|
|
|
when(stream.iterator()).thenReturn(iterator);
|
|
|
|
|
final MessageAndMetadata messageAndMetadata = mock(MessageAndMetadata.class);
|
|
|
|
|
final MessageAndMetadata<byte[], String> 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");
|
|
|
|
|
Assert.assertEquals(1, messages.size());
|
|
|
|
|
Assert.assertEquals(1, messages.get("topic").size());
|
|
|
|
|
Assert.assertEquals("got message", messages.get("topic").get(1).get(0));
|
|
|
|
|
|
|
|
|
|
verify(stream, times(1)).iterator();
|
|
|
|
|
verify(iterator, times(1)).next();
|
|
|
|
|
@@ -429,10 +434,10 @@ public class ConsumerConfigurationTests<K,V> {
|
|
|
|
|
@Test
|
|
|
|
|
@SuppressWarnings("unchecked")
|
|
|
|
|
public void testReceiveMessageForTopicFilterFromMultipleStreams() {
|
|
|
|
|
final ConsumerMetadata consumerMetadata = mock(ConsumerMetadata.class);
|
|
|
|
|
final ConsumerMetadata<byte[], byte[]> consumerMetadata = mock(ConsumerMetadata.class);
|
|
|
|
|
final ConsumerConnectionProvider consumerConnectionProvider =
|
|
|
|
|
mock(ConsumerConnectionProvider.class);
|
|
|
|
|
final MessageLeftOverTracker messageLeftOverTracker = mock(MessageLeftOverTracker.class);
|
|
|
|
|
final MessageLeftOverTracker<byte[], byte[]> messageLeftOverTracker = mock(MessageLeftOverTracker.class);
|
|
|
|
|
|
|
|
|
|
when(consumerMetadata.getTopicFilterConfiguration()).thenReturn(new TopicFilterConfiguration(".*", 1, false));
|
|
|
|
|
|
|
|
|
|
@@ -440,48 +445,49 @@ public class ConsumerConfigurationTests<K,V> {
|
|
|
|
|
|
|
|
|
|
when(consumerConnectionProvider.getConsumerConnector()).thenReturn(consumerConnector);
|
|
|
|
|
|
|
|
|
|
final ConsumerConfiguration consumerConfiguration = new ConsumerConfiguration(consumerMetadata,
|
|
|
|
|
consumerConnectionProvider, messageLeftOverTracker);
|
|
|
|
|
final ConsumerConfiguration<byte[], byte[]> consumerConfiguration =
|
|
|
|
|
new ConsumerConfiguration<byte[], byte[]>(consumerMetadata, consumerConnectionProvider,
|
|
|
|
|
messageLeftOverTracker);
|
|
|
|
|
consumerConfiguration.setMaxMessages(3);
|
|
|
|
|
|
|
|
|
|
final KafkaStream stream1 = mock(KafkaStream.class);
|
|
|
|
|
final KafkaStream stream2 = mock(KafkaStream.class);
|
|
|
|
|
final KafkaStream stream3 = mock(KafkaStream.class);
|
|
|
|
|
final KafkaStream<byte[], byte[]> stream1 = mock(KafkaStream.class);
|
|
|
|
|
final KafkaStream<byte[], byte[]> stream2 = mock(KafkaStream.class);
|
|
|
|
|
final KafkaStream<byte[], byte[]> stream3 = mock(KafkaStream.class);
|
|
|
|
|
final List<KafkaStream<byte[], byte[]>> streams = new ArrayList<KafkaStream<byte[], byte[]>>();
|
|
|
|
|
streams.add(stream1);
|
|
|
|
|
streams.add(stream2);
|
|
|
|
|
streams.add(stream3);
|
|
|
|
|
|
|
|
|
|
when(consumerConfiguration.createMessageStreamsForTopicFilter()).thenReturn(streams);
|
|
|
|
|
final ConsumerIterator iterator1 = mock(ConsumerIterator.class);
|
|
|
|
|
final ConsumerIterator iterator2 = mock(ConsumerIterator.class);
|
|
|
|
|
final ConsumerIterator iterator3 = mock(ConsumerIterator.class);
|
|
|
|
|
final ConsumerIterator<byte[], byte[]> iterator1 = mock(ConsumerIterator.class);
|
|
|
|
|
final ConsumerIterator<byte[], byte[]> iterator2 = mock(ConsumerIterator.class);
|
|
|
|
|
final ConsumerIterator<byte[], byte[]> iterator3 = mock(ConsumerIterator.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);
|
|
|
|
|
final MessageAndMetadata<byte[], byte[]> messageAndMetadata1 = mock(MessageAndMetadata.class);
|
|
|
|
|
final MessageAndMetadata<byte[], byte[]> messageAndMetadata2 = mock(MessageAndMetadata.class);
|
|
|
|
|
final MessageAndMetadata<byte[], byte[]> messageAndMetadata3 = mock(MessageAndMetadata.class);
|
|
|
|
|
|
|
|
|
|
when(iterator1.next()).thenReturn(messageAndMetadata1);
|
|
|
|
|
when(iterator2.next()).thenReturn(messageAndMetadata2);
|
|
|
|
|
when(iterator3.next()).thenReturn(messageAndMetadata3);
|
|
|
|
|
|
|
|
|
|
when(messageAndMetadata1.message()).thenReturn("got message");
|
|
|
|
|
when(messageAndMetadata1.message()).thenReturn("got message".getBytes());
|
|
|
|
|
when(messageAndMetadata1.topic()).thenReturn("topic");
|
|
|
|
|
when(messageAndMetadata1.partition()).thenReturn(1);
|
|
|
|
|
|
|
|
|
|
when(messageAndMetadata2.message()).thenReturn("got message");
|
|
|
|
|
when(messageAndMetadata2.message()).thenReturn("got message".getBytes());
|
|
|
|
|
when(messageAndMetadata2.topic()).thenReturn("topic");
|
|
|
|
|
when(messageAndMetadata2.partition()).thenReturn(2);
|
|
|
|
|
|
|
|
|
|
when(messageAndMetadata3.message()).thenReturn("got message");
|
|
|
|
|
when(messageAndMetadata3.message()).thenReturn("got message".getBytes());
|
|
|
|
|
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);
|
|
|
|
|
Assert.assertEquals(1, messages.size());
|
|
|
|
|
int sum = 0;
|
|
|
|
|
|
|
|
|
|
final Map<Integer, List<Object>> values = messages.get("topic");
|
|
|
|
|
@@ -490,7 +496,7 @@ public class ConsumerConfigurationTests<K,V> {
|
|
|
|
|
sum += l.size();
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
Assert.assertEquals(sum, 3);
|
|
|
|
|
Assert.assertEquals(3, sum);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private boolean valueFound(final List<Object> l, final String value){
|
|
|
|
|
|