From fb40e521c7f7c9650ff7d8170edc9fb542999132 Mon Sep 17 00:00:00 2001 From: Jonathan Pearlin Date: Thu, 12 Jun 2014 18:07:19 -0400 Subject: [PATCH] INTEXT-106: Update to spring-integration 4.0 JIRA: https://jira.spring.io/browse/INTEXT-106 * Upgrade to SI 4.0 and other libs * Upgrade to Gradle 1.12 * Polisihng for `build.gradle` * Fix compile warnings --- .../KafkaHighLevelConsumerMessageSource.java | 2 +- .../outbound/KafkaProducerMessageHandler.java | 2 +- .../support/ConsumerConfigFactoryBean.java | 9 +- .../kafka/support/ConsumerConfiguration.java | 2 +- .../kafka/support/KafkaConsumerContext.java | 2 +- .../kafka/support/KafkaProducerContext.java | 5 +- .../kafka/support/ProducerConfiguration.java | 2 +- .../support/ConsumerConfigurationTests.java | 90 ++++++++++--------- .../support/KafkaConsumerContextTest.java | 2 +- .../support/ProducerConfigurationTests.java | 2 +- 10 files changed, 63 insertions(+), 55 deletions(-) diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaHighLevelConsumerMessageSource.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaHighLevelConsumerMessageSource.java index 8926e7e5e1..449a8ebeff 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaHighLevelConsumerMessageSource.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaHighLevelConsumerMessageSource.java @@ -15,10 +15,10 @@ */ package org.springframework.integration.kafka.inbound; -import org.springframework.integration.Message; import org.springframework.integration.context.IntegrationObjectSupport; import org.springframework.integration.core.MessageSource; import org.springframework.integration.kafka.support.KafkaConsumerContext; +import org.springframework.messaging.Message; import java.util.List; import java.util.Map; diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/outbound/KafkaProducerMessageHandler.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/outbound/KafkaProducerMessageHandler.java index ade5f1ae7d..a9062e7b06 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/outbound/KafkaProducerMessageHandler.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/outbound/KafkaProducerMessageHandler.java @@ -15,9 +15,9 @@ */ package org.springframework.integration.kafka.outbound; -import org.springframework.integration.Message; import org.springframework.integration.handler.AbstractMessageHandler; import org.springframework.integration.kafka.support.KafkaProducerContext; +import org.springframework.messaging.Message; /** * @author Soby Chacko diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ConsumerConfigFactoryBean.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ConsumerConfigFactoryBean.java index 6a39d49306..cbb81aa610 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ConsumerConfigFactoryBean.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ConsumerConfigFactoryBean.java @@ -15,13 +15,13 @@ */ package org.springframework.integration.kafka.support; -import kafka.consumer.ConsumerConfig; +import java.util.Properties; +import kafka.consumer.ConsumerConfig; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; -import org.springframework.beans.factory.FactoryBean; -import java.util.Properties; +import org.springframework.beans.factory.FactoryBean; /** * @author Soby Chacko @@ -43,7 +43,8 @@ public class ConsumerConfigFactoryBean implements FactoryBean consumerMetadata, + final ZookeeperConnect zookeeperConnect) { this(consumerMetadata, zookeeperConnect, null); } diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ConsumerConfiguration.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ConsumerConfiguration.java index 6ad6f2ff3e..2514a08650 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ConsumerConfiguration.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ConsumerConfiguration.java @@ -18,7 +18,7 @@ import kafka.message.MessageAndMetadata; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; -import org.springframework.integration.MessagingException; +import org.springframework.messaging.MessagingException; /** * @author Soby Chacko diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/KafkaConsumerContext.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/KafkaConsumerContext.java index 6822ccfa7f..5a0f4f2818 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/KafkaConsumerContext.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/KafkaConsumerContext.java @@ -19,9 +19,9 @@ import org.springframework.beans.BeansException; import org.springframework.beans.factory.BeanFactory; import org.springframework.beans.factory.BeanFactoryAware; import org.springframework.beans.factory.ListableBeanFactory; -import org.springframework.integration.Message; import org.springframework.integration.kafka.core.KafkaConsumerDefaults; import org.springframework.integration.support.MessageBuilder; +import org.springframework.messaging.Message; import java.util.Collection; import java.util.HashMap; diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/KafkaProducerContext.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/KafkaProducerContext.java index bf9b047c2b..8151e7d22b 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/KafkaProducerContext.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/KafkaProducerContext.java @@ -19,11 +19,12 @@ import java.util.Collection; import java.util.Map; import java.util.Properties; +import org.apache.commons.lang.StringUtils; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.springframework.beans.BeansException; import org.springframework.beans.factory.*; -import org.springframework.integration.Message; +import org.springframework.messaging.Message; /** * @author Soby Chacko @@ -52,7 +53,7 @@ public class KafkaProducerContext implements BeanFactoryAware { return producerConfiguration; } } - LOGGER.error("No is producer-configuration defined for topic " + topic + ". cannot send message"); + LOGGER.error("No producer-configuration defined for topic " + topic + ". cannot send message"); return null; } diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ProducerConfiguration.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ProducerConfiguration.java index e7e325eda9..93a4d7cc43 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ProducerConfiguration.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ProducerConfiguration.java @@ -16,7 +16,7 @@ import kafka.serializer.DefaultEncoder; import org.apache.commons.lang.builder.EqualsBuilder; import org.apache.commons.lang.builder.HashCodeBuilder; -import org.springframework.integration.Message; +import org.springframework.messaging.Message; /** * @author Soby Chacko diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/support/ConsumerConfigurationTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/support/ConsumerConfigurationTests.java index 9e810769f7..c23a2b676f 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/support/ConsumerConfigurationTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/support/ConsumerConfigurationTests.java @@ -37,6 +37,7 @@ import org.mockito.stubbing.Answer; * @since 0.5 */ public class ConsumerConfigurationTests { + @Test @SuppressWarnings("unchecked") public void testReceiveMessageForSingleTopicFromSingleStream() { @@ -72,9 +73,9 @@ public class ConsumerConfigurationTests { when(messageAndMetadata.partition()).thenReturn(1); final Map>> 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 { sum += l.size(); } - Assert.assertEquals(sum, 3); + Assert.assertEquals(3, sum); } @Test @@ -226,7 +227,7 @@ public class ConsumerConfigurationTests { } } - Assert.assertEquals(sum, 9); + Assert.assertEquals(9, sum); } @@ -284,11 +285,11 @@ public class ConsumerConfigurationTests { final Map>> messageStreams = new HashMap>>(); messageStreams.put("topic1", streams); when(consumerConfiguration.createMessageStreamsForTopic()).thenReturn(messageStreams); - final ConsumerIterator iterator = mock(ConsumerIterator.class); + final ConsumerIterator iterator = mock(ConsumerIterator.class); when(stream.iterator()).thenReturn(iterator); - final MessageAndMetadata messageAndMetadata = mock(MessageAndMetadata.class); + final MessageAndMetadata 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 { } } - 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 { when(mockedConsumerConnectionProvider.getConsumerConnector()).thenReturn(mockedConsumerConnector); final Map>> messageStreams = new HashMap>>(); - when((Map>>) (Object) mockedConsumerConnector.createMessageStreams(topicsStreamMap)).thenReturn(messageStreams); + when((Map>>) + (Object) mockedConsumerConnector.createMessageStreams(topicsStreamMap)).thenReturn(messageStreams); final ConsumerConfiguration consumerConfiguration = new ConsumerConfiguration(mockedConsumerMetadata, mockedConsumerConnectionProvider, mockedMessageLeftOverTracker); @@ -374,51 +376,54 @@ public class ConsumerConfigurationTests { final Map>> messageStreams = new HashMap>>(); when(mockedConsumerConnector.createMessageStreams(topicsStreamMap)).thenReturn(messageStreams); - final ConsumerConfiguration consumerConfiguration = new ConsumerConfiguration(mockedConsumerMetadata, - mockedConsumerConnectionProvider, mockedMessageLeftOverTracker); + final ConsumerConfiguration consumerConfiguration = + new ConsumerConfiguration(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 consumerMetadata = mock(ConsumerMetadata.class); final ConsumerConnectionProvider consumerConnectionProvider = mock(ConsumerConnectionProvider.class); - final MessageLeftOverTracker messageLeftOverTracker = mock(MessageLeftOverTracker.class); + final MessageLeftOverTracker 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 consumerConfiguration = + new ConsumerConfiguration(consumerMetadata, consumerConnectionProvider, + messageLeftOverTracker); consumerConfiguration.setMaxMessages(1); - final KafkaStream stream = mock(KafkaStream.class); - final List> streams = new ArrayList>(); + final KafkaStream stream = mock(KafkaStream.class); + final List> streams = new ArrayList>(); streams.add(stream); when(consumerConfiguration.createMessageStreamsForTopicFilter()).thenReturn(streams); - final ConsumerIterator iterator = mock(ConsumerIterator.class); + final ConsumerIterator iterator = mock(ConsumerIterator.class); when(stream.iterator()).thenReturn(iterator); - final MessageAndMetadata messageAndMetadata = mock(MessageAndMetadata.class); + 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>> 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 { @Test @SuppressWarnings("unchecked") public void testReceiveMessageForTopicFilterFromMultipleStreams() { - final ConsumerMetadata consumerMetadata = mock(ConsumerMetadata.class); + final ConsumerMetadata consumerMetadata = mock(ConsumerMetadata.class); final ConsumerConnectionProvider consumerConnectionProvider = mock(ConsumerConnectionProvider.class); - final MessageLeftOverTracker messageLeftOverTracker = mock(MessageLeftOverTracker.class); + final MessageLeftOverTracker messageLeftOverTracker = mock(MessageLeftOverTracker.class); when(consumerMetadata.getTopicFilterConfiguration()).thenReturn(new TopicFilterConfiguration(".*", 1, false)); @@ -440,48 +445,49 @@ public class ConsumerConfigurationTests { when(consumerConnectionProvider.getConsumerConnector()).thenReturn(consumerConnector); - final ConsumerConfiguration consumerConfiguration = new ConsumerConfiguration(consumerMetadata, - consumerConnectionProvider, messageLeftOverTracker); + final ConsumerConfiguration consumerConfiguration = + new ConsumerConfiguration(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 stream1 = mock(KafkaStream.class); + final KafkaStream stream2 = mock(KafkaStream.class); + final KafkaStream stream3 = mock(KafkaStream.class); final List> streams = new ArrayList>(); 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 iterator1 = mock(ConsumerIterator.class); + final ConsumerIterator iterator2 = mock(ConsumerIterator.class); + final ConsumerIterator 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 messageAndMetadata1 = mock(MessageAndMetadata.class); + final MessageAndMetadata messageAndMetadata2 = mock(MessageAndMetadata.class); + final MessageAndMetadata 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>> messages = consumerConfiguration.receive(); - Assert.assertEquals(messages.size(), 1); + Assert.assertEquals(1, messages.size()); int sum = 0; final Map> values = messages.get("topic"); @@ -490,7 +496,7 @@ public class ConsumerConfigurationTests { sum += l.size(); } - Assert.assertEquals(sum, 3); + Assert.assertEquals(3, sum); } private boolean valueFound(final List l, final String value){ diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/support/KafkaConsumerContextTest.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/support/KafkaConsumerContextTest.java index a5fffed09d..bb28b927be 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/support/KafkaConsumerContextTest.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/support/KafkaConsumerContextTest.java @@ -19,7 +19,7 @@ import org.junit.Assert; import org.junit.Test; import org.mockito.Mockito; import org.springframework.beans.factory.ListableBeanFactory; -import org.springframework.integration.Message; +import org.springframework.messaging.Message; import java.util.ArrayList; import java.util.HashMap; diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/support/ProducerConfigurationTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/support/ProducerConfigurationTests.java index 367e464186..cbdad2c771 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/support/ProducerConfigurationTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/support/ProducerConfigurationTests.java @@ -23,13 +23,13 @@ import org.junit.Assert; import org.junit.Test; import org.mockito.ArgumentCaptor; import org.mockito.Mockito; -import org.springframework.integration.Message; import org.springframework.integration.kafka.serializer.avro.AvroReflectDatumBackedKafkaEncoder; import org.springframework.integration.kafka.test.utils.NonSerializableTestKey; import org.springframework.integration.kafka.test.utils.NonSerializableTestPayload; import org.springframework.integration.kafka.test.utils.TestKey; import org.springframework.integration.kafka.test.utils.TestPayload; import org.springframework.integration.support.MessageBuilder; +import org.springframework.messaging.Message; import java.io.ByteArrayInputStream; import java.io.NotSerializableException;