diff --git a/samples/kafka/build.gradle b/samples/kafka/build.gradle index 5b02cbd..53a8879 100644 --- a/samples/kafka/build.gradle +++ b/samples/kafka/build.gradle @@ -37,15 +37,8 @@ repositories { } dependencies { - compile([ - "org.apache.avro:avro:1.7.3", - "org.apache.avro:avro-compiler:1.7.3" - ]) - compile "org.springframework:spring-beans:3.1.3.RELEASE" - compile "org.springframework:spring-context:3.1.3.RELEASE" - compile "org.springframework:spring-expression:3.1.3.RELEASE" - compile "org.springframework.integration:spring-integration-stream:2.2.0.RELEASE" - compile("org.springframework.integration:spring-integration-kafka:0.5.0.BUILD-SNAPSHOT") { + compile "org.springframework.integration:spring-integration-stream:$springIntegrationVersion" + compile("org.springframework.integration:spring-integration-kafka:$springIntegrationKafkaVersion") { exclude module: 'log4j' exclude module: 'jms' exclude module: 'jmxtools' diff --git a/samples/kafka/gradle.properties b/samples/kafka/gradle.properties index d03e972..d1a04db 100644 --- a/samples/kafka/gradle.properties +++ b/samples/kafka/gradle.properties @@ -1,7 +1,6 @@ -springVersion = 3.1.3.RELEASE -springIntegrationVersion = 2.2.0.RELEASE -springIntegrationKafkaVersion = 0.5.0.BUILD-SNAPSHOT -version = 0.5.0.BUILD-SNAPSHOT +springIntegrationVersion = 4.0.3.RELEASE +springIntegrationKafkaVersion = 1.0.0.BUILD-SNAPSHOT +version = 1.0.0.BUILD-SNAPSHOT diff --git a/samples/kafka/gradle/wrapper/gradle-wrapper.jar b/samples/kafka/gradle/wrapper/gradle-wrapper.jar index 7f1e239..667288a 100644 Binary files a/samples/kafka/gradle/wrapper/gradle-wrapper.jar and b/samples/kafka/gradle/wrapper/gradle-wrapper.jar differ diff --git a/samples/kafka/gradle/wrapper/gradle-wrapper.properties b/samples/kafka/gradle/wrapper/gradle-wrapper.properties index 238e92a..30eb5f3 100644 --- a/samples/kafka/gradle/wrapper/gradle-wrapper.properties +++ b/samples/kafka/gradle/wrapper/gradle-wrapper.properties @@ -1,6 +1,6 @@ -#Mon Oct 28 16:21:39 EDT 2013 +#Tue Jul 22 17:46:28 EEST 2014 distributionBase=GRADLE_USER_HOME distributionPath=wrapper/dists zipStoreBase=GRADLE_USER_HOME zipStorePath=wrapper/dists -distributionUrl=http\://services.gradle.org/distributions/gradle-1.8-bin.zip +distributionUrl=http\://services.gradle.org/distributions/gradle-1.12-all.zip diff --git a/samples/kafka/gradlew b/samples/kafka/gradlew index e61422d..91a7e26 100755 --- a/samples/kafka/gradlew +++ b/samples/kafka/gradlew @@ -1,4 +1,4 @@ -#!/bin/bash +#!/usr/bin/env bash ############################################################################## ## @@ -61,9 +61,9 @@ while [ -h "$PRG" ] ; do fi done SAVED="`pwd`" -cd "`dirname \"$PRG\"`/" +cd "`dirname \"$PRG\"`/" >&- APP_HOME="`pwd -P`" -cd "$SAVED" +cd "$SAVED" >&- CLASSPATH=$APP_HOME/gradle/wrapper/gradle-wrapper.jar diff --git a/samples/kafka/src/main/java/org/springframework/integration/samples/kafka/outbound/CustomPartitioner.java b/samples/kafka/src/main/java/org/springframework/integration/samples/kafka/outbound/CustomPartitioner.java index a6f3318..884bd58 100644 --- a/samples/kafka/src/main/java/org/springframework/integration/samples/kafka/outbound/CustomPartitioner.java +++ b/samples/kafka/src/main/java/org/springframework/integration/samples/kafka/outbound/CustomPartitioner.java @@ -22,14 +22,14 @@ import kafka.producer.Partitioner; * * This class is for internal use only and therefore is at default access level */ -class CustomPartitioner implements Partitioner { +class CustomPartitioner implements Partitioner { /** * Uses the key to calculate a partition bucket id for routing * the data to the appropriate broker partition * @return an integer between 0 and numPartitions-1 */ @Override - public int partition(final T key, final int numPartitions) { + public int partition(final Object key, final int numPartitions) { final String s = (String) key; final Integer i = Integer.parseInt(s); return i % numPartitions; diff --git a/samples/kafka/src/main/java/org/springframework/integration/samples/kafka/outbound/OutboundRunner.java b/samples/kafka/src/main/java/org/springframework/integration/samples/kafka/outbound/OutboundRunner.java index f285c8d..53e7fc0 100644 --- a/samples/kafka/src/main/java/org/springframework/integration/samples/kafka/outbound/OutboundRunner.java +++ b/samples/kafka/src/main/java/org/springframework/integration/samples/kafka/outbound/OutboundRunner.java @@ -18,9 +18,9 @@ package org.springframework.integration.samples.kafka.outbound; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.springframework.context.support.ClassPathXmlApplicationContext; -import org.springframework.integration.MessageChannel; import org.springframework.integration.samples.kafka.user.User; import org.springframework.integration.support.MessageBuilder; +import org.springframework.messaging.MessageChannel; public class OutboundRunner { private static final String CONFIG = "kafkaOutboundAdapterParserTests-context.xml"; diff --git a/samples/kafka/src/main/java/org/springframework/integration/samples/kafka/outbound/PartitionlessTransformer.java b/samples/kafka/src/main/java/org/springframework/integration/samples/kafka/outbound/PartitionlessTransformer.java index 5846ae6..c4a4379 100644 --- a/samples/kafka/src/main/java/org/springframework/integration/samples/kafka/outbound/PartitionlessTransformer.java +++ b/samples/kafka/src/main/java/org/springframework/integration/samples/kafka/outbound/PartitionlessTransformer.java @@ -15,9 +15,9 @@ */ package org.springframework.integration.samples.kafka.outbound; -import org.springframework.integration.Message; import org.springframework.integration.support.MessageBuilder; import org.springframework.integration.transformer.Transformer; +import org.springframework.messaging.Message; import java.util.ArrayList; import java.util.Collection; @@ -36,13 +36,13 @@ public class PartitionlessTransformer implements Transformer { final Map>> origData = (Map>>) message.getPayload(); - final Map> nonPartitionedData = new HashMap>(); + final Map> nonPartitionedData = new HashMap<>(); for(final String topic : origData.keySet()) { final Map> partitionedData = origData.get(topic); final Collection> nonPartitionedDataFromTopic = partitionedData.values(); - final List mergedList = new ArrayList(); + final List mergedList = new ArrayList<>(); for (final List l : nonPartitionedDataFromTopic){ mergedList.addAll(l); diff --git a/spring-integration-kafka/build.gradle b/spring-integration-kafka/build.gradle index 153f8b1..a722984 100644 --- a/spring-integration-kafka/build.gradle +++ b/spring-integration-kafka/build.gradle @@ -19,7 +19,7 @@ sourceCompatibility = targetCompatibility = 1.6 ext { avroVersion = '1.7.6' jacocoVersion = '0.7.0.201403182114' - kafkaVersion = '0.8.0' + kafkaVersion = '0.8.1.1' metricsVersion = '2.2.0' scalaVersion = '2.10' springIntegrationVersion = '4.0.3.RELEASE' diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/DefaultPartitioner.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/DefaultPartitioner.java index 929ef89..007c116 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/DefaultPartitioner.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/DefaultPartitioner.java @@ -24,14 +24,15 @@ import kafka.utils.Utils; * * This class is for internal use only and therefore is at default access level */ -class DefaultPartitioner implements Partitioner { +class DefaultPartitioner implements Partitioner { /** * Uses the key to calculate a partition bucket id for routing * the data to the appropriate broker partition * @return an integer between 0 and numPartitions-1 */ @Override - public int partition(final T key, final int numPartitions) { + public int partition(final Object key, final int numPartitions) { return Utils.abs(key.hashCode()) % numPartitions; } + } diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ProducerFactoryBean.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ProducerFactoryBean.java index b622154..2dfbe6c 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ProducerFactoryBean.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ProducerFactoryBean.java @@ -71,7 +71,7 @@ public class ProducerFactoryBean implements FactoryBean> { LOGGER.info("Using producer properties => " + props); final ProducerConfig config = new ProducerConfig(props); final EventHandler eventHandler = new DefaultEventHandler(config, - producerMetadata.getPartitioner() == null ? new DefaultPartitioner() : producerMetadata.getPartitioner(), + producerMetadata.getPartitioner() == null ? new DefaultPartitioner() : producerMetadata.getPartitioner(), producerMetadata.getValueEncoder(), producerMetadata.getKeyEncoder(), new ProducerPool(config), new HashMap()); diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ProducerMetadata.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ProducerMetadata.java index e6dca25..fe6ba9a 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ProducerMetadata.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ProducerMetadata.java @@ -35,7 +35,7 @@ public class ProducerMetadata implements InitializingBean { private Class valueClassType; private final String topic; private String compressionCodec = "default"; - private Partitioner partitioner; + private Partitioner partitioner; private boolean async = false; private String batchNumMessages; @@ -94,11 +94,11 @@ public class ProducerMetadata implements InitializingBean { this.compressionCodec = compressionCodec; } - public Partitioner getPartitioner() { + public Partitioner getPartitioner() { return partitioner; } - public void setPartitioner(final Partitioner partitioner) { + public void setPartitioner(final Partitioner partitioner) { this.partitioner = partitioner; } 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 c23a2b6..1dab175 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 @@ -262,9 +262,25 @@ public class ConsumerConfigurationTests { topicStreamMap.put("topic1", 1); when(consumerMetadata.getTopicStreamMap()).thenReturn(topicStreamMap); 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); + + final MessageAndMetadata m1 = mock(MessageAndMetadata.class); + final MessageAndMetadata m2 = mock(MessageAndMetadata.class); + final MessageAndMetadata m3 = mock(MessageAndMetadata.class); + + when(m1.key()).thenReturn("key1"); + when(m1.message()).thenReturn("value1"); + when(m1.topic()).thenReturn("topic1"); + when(m1.partition()).thenReturn(1); + + when(m2.key()).thenReturn("key2"); + when(m2.message()).thenReturn("value2"); + when(m2.topic()).thenReturn("topic2"); + when(m2.partition()).thenReturn(1); + + when(m3.key()).thenReturn("key1"); + when(m3.message()).thenReturn("value3"); + when(m3.topic()).thenReturn("topic3"); + when(m3.partition()).thenReturn(1); final List> mList = new ArrayList>(); mList.add(m1);