From 94734e4ca3c9e1b73b7ad7ee3ccafff2f03d6827 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Tue, 16 Jul 2013 22:40:28 -0400 Subject: [PATCH] INTEXT-84 Kafka: Enhance Avro serialization support --- .../xml/KafkaConsumerContextParser.java | 4 + .../xml/KafkaProducerContextParser.java | 5 + .../KafkaHighLevelConsumerMessageSource.java | 6 +- .../outbound/KafkaProducerMessageHandler.java | 8 +- .../serializer/avro/AvroDatumSupport.java | 43 ++++++ ...> AvroReflectDatumBackedKafkaDecoder.java} | 28 ++-- ...> AvroReflectDatumBackedKafkaEncoder.java} | 27 ++-- .../kafka/serializer/avro/AvroSerializer.java | 11 +- .../AvroSpecificDatumBackedKafkaDecoder.java | 28 ++++ .../AvroSpecificDatumBackedKafkaEncoder.java | 28 ++++ .../avro/AvroSpecificDatumSerializer.java | 53 ------- .../serializer/common/StringEncoder.java | 2 +- .../support/ConsumerConfigFactoryBean.java | 6 +- .../kafka/support/ConsumerConfiguration.java | 66 ++++----- .../kafka/support/ConsumerMetadata.java | 16 ++- .../kafka/support/KafkaConsumerContext.java | 12 +- .../kafka/support/KafkaProducerContext.java | 19 +-- .../kafka/support/MessageLeftOverTracker.java | 8 +- ...afkaConsumerContextParserTests-context.xml | 2 +- .../xml/KafkaConsumerContextParserTests.java | 6 +- ...KafkaInboundAdapterParserTests-context.xml | 2 +- ...afkaOutboundAdapterParserTests-context.xml | 2 +- .../xml/KafkaOutboundAdapterParserTests.java | 7 +- ...afkaProducerContextParserTests-context.xml | 2 +- .../xml/KafkaProducerContextParserTests.java | 24 ++-- ...eflectDatumBackedKafkaSerializerTest.java} | 18 +-- ...pecificDatumBackedKafkaSerializerTest.java | 30 ++++ .../support/ConsumerConfigurationTests.java | 132 +++++++++--------- .../support/KafkaConsumerContextTest.java | 14 +- .../support/ProducerConfigurationTests.java | 50 ++++--- .../support/ProducerFactoryBeanTests.java | 6 +- .../integration/kafka/test/utils/User.java | 96 +++++++++++++ 32 files changed, 471 insertions(+), 290 deletions(-) create mode 100644 spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/avro/AvroDatumSupport.java rename spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/avro/{AvroBackedKafkaDecoder.java => AvroReflectDatumBackedKafkaDecoder.java} (58%) rename spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/avro/{AvroBackedKafkaEncoder.java => AvroReflectDatumBackedKafkaEncoder.java} (58%) create mode 100644 spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/avro/AvroSpecificDatumBackedKafkaDecoder.java create mode 100644 spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/avro/AvroSpecificDatumBackedKafkaEncoder.java delete mode 100644 spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/avro/AvroSpecificDatumSerializer.java rename spring-integration-kafka/src/test/java/org/springframework/integration/kafka/serializer/{AvroBackedKafkaSerializerTest.java => AvroReflectDatumBackedKafkaSerializerTest.java} (64%) create mode 100644 spring-integration-kafka/src/test/java/org/springframework/integration/kafka/serializer/AvroSpecificDatumBackedKafkaSerializerTest.java create mode 100644 spring-integration-kafka/src/test/java/org/springframework/integration/kafka/test/utils/User.java diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaConsumerContextParser.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaConsumerContextParser.java index 0ac442be62..e116eec6ed 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaConsumerContextParser.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaConsumerContextParser.java @@ -15,6 +15,7 @@ */ package org.springframework.integration.kafka.config.xml; +import kafka.serializer.DefaultDecoder; import org.springframework.beans.factory.config.BeanDefinition; import org.springframework.beans.factory.config.BeanDefinitionHolder; import org.springframework.beans.factory.support.AbstractBeanDefinition; @@ -32,7 +33,9 @@ import org.springframework.util.StringUtils; import org.springframework.util.xml.DomUtils; import org.w3c.dom.Element; +import java.util.ArrayList; import java.util.HashMap; +import java.util.List; import java.util.Map; /** @@ -54,6 +57,7 @@ public class KafkaConsumerContextParser extends AbstractSingleBeanDefinitionPars parseConsumerConfigurations(consumerConfigurations, parserContext, builder, element); } + @SuppressWarnings("unchecked") private void parseConsumerConfigurations(final Element consumerConfigurations, final ParserContext parserContext, final BeanDefinitionBuilder builder, final Element parentElem) { for (final Element consumerConfiguration : DomUtils.getChildElementsByTagName(consumerConfigurations, "consumer-configuration")) { diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaProducerContextParser.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaProducerContextParser.java index 5fc99379d2..dadfe47771 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaProducerContextParser.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaProducerContextParser.java @@ -30,11 +30,15 @@ import org.springframework.util.StringUtils; import org.springframework.util.xml.DomUtils; import org.w3c.dom.Element; +import java.util.HashMap; +import java.util.Map; + /** * @author Soby Chacko * @since 0.5 */ public class KafkaProducerContextParser extends AbstractSimpleBeanDefinitionParser { + @Override protected Class getBeanClass(final Element element) { return KafkaProducerContext.class; @@ -48,6 +52,7 @@ public class KafkaProducerContextParser extends AbstractSimpleBeanDefinitionPars parseProducerConfigurations(topics, parserContext); } + @SuppressWarnings("unchecked") private void parseProducerConfigurations(final Element topics, final ParserContext parserContext) { for (final Element producerConfiguration : DomUtils.getChildElementsByTagName(topics, "producer-configuration")){ final BeanDefinitionBuilder producerConfigurationBuilder = BeanDefinitionBuilder.genericBeanDefinition(ProducerConfiguration.class); 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 83e225e1ce..8926e7e5e1 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 @@ -28,11 +28,11 @@ import java.util.Map; * @since 0.5 * */ -public class KafkaHighLevelConsumerMessageSource extends IntegrationObjectSupport implements MessageSource>>> { +public class KafkaHighLevelConsumerMessageSource extends IntegrationObjectSupport implements MessageSource>>> { - private final KafkaConsumerContext kafkaConsumerContext; + private final KafkaConsumerContext kafkaConsumerContext; - public KafkaHighLevelConsumerMessageSource(final KafkaConsumerContext kafkaConsumerContext) { + public KafkaHighLevelConsumerMessageSource(final KafkaConsumerContext kafkaConsumerContext) { this.kafkaConsumerContext = kafkaConsumerContext; } 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 fdf918cb8c..ade5f1ae7d 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 @@ -23,15 +23,15 @@ import org.springframework.integration.kafka.support.KafkaProducerContext; * @author Soby Chacko * @since 0.5 */ -public class KafkaProducerMessageHandler extends AbstractMessageHandler { +public class KafkaProducerMessageHandler extends AbstractMessageHandler { - private final KafkaProducerContext kafkaProducerContext; + private final KafkaProducerContext kafkaProducerContext; - public KafkaProducerMessageHandler(final KafkaProducerContext kafkaProducerContext) { + public KafkaProducerMessageHandler(final KafkaProducerContext kafkaProducerContext) { this.kafkaProducerContext = kafkaProducerContext; } - public KafkaProducerContext getKafkaProducerContext() { + public KafkaProducerContext getKafkaProducerContext() { return kafkaProducerContext; } diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/avro/AvroDatumSupport.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/avro/AvroDatumSupport.java new file mode 100644 index 0000000000..6165b1d39b --- /dev/null +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/avro/AvroDatumSupport.java @@ -0,0 +1,43 @@ +package org.springframework.integration.kafka.serializer.avro; + +import org.apache.avro.io.DatumReader; +import org.apache.avro.io.DatumWriter; +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; + +import java.io.IOException; + +/** + * @author Soby Chacko + * @since 0.5 + */ +public abstract class AvroDatumSupport { + + private static final Log LOG = LogFactory.getLog(AvroDatumSupport.class); + + private final AvroSerializer avroSerializer; + + protected AvroDatumSupport() { + this.avroSerializer = new AvroSerializer(); + } + + @SuppressWarnings("unchecked") + public byte[] toBytes(final T source, final DatumWriter writer) { + try { + return avroSerializer.serialize(source, writer); + } catch (IOException e) { + LOG.error("Failed to encode source: " + e); + } + return null; + } + + @SuppressWarnings("unchecked") + public T fromBytes(final byte[] bytes, final DatumReader reader) { + try { + return avroSerializer.deserialize(bytes, reader); + } catch (IOException e) { + LOG.error("Failed to decode byte array: " + e); + } + return null; + } +} diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/avro/AvroBackedKafkaDecoder.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/avro/AvroReflectDatumBackedKafkaDecoder.java similarity index 58% rename from spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/avro/AvroBackedKafkaDecoder.java rename to spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/avro/AvroReflectDatumBackedKafkaDecoder.java index d702492f00..7e9b7c95c4 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/avro/AvroBackedKafkaDecoder.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/avro/AvroReflectDatumBackedKafkaDecoder.java @@ -15,12 +15,9 @@ */ package org.springframework.integration.kafka.serializer.avro; - import kafka.serializer.Decoder; -import org.apache.avro.Schema; -import org.apache.avro.reflect.ReflectData; - -import java.io.IOException; +import org.apache.avro.io.DatumReader; +import org.apache.avro.reflect.ReflectDatumReader; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; @@ -28,28 +25,19 @@ import org.apache.commons.logging.LogFactory; * @author Soby Chacko * @since 0.5 */ -public class AvroBackedKafkaDecoder implements Decoder { - private static final Log LOG = LogFactory.getLog(AvroBackedKafkaDecoder.class); +public class AvroReflectDatumBackedKafkaDecoder extends AvroDatumSupport implements Decoder { + private static final Log LOG = LogFactory.getLog(AvroReflectDatumBackedKafkaDecoder.class); - private final Class clazz; + private final DatumReader reader; - public AvroBackedKafkaDecoder(final Class clazz) { - this.clazz = clazz; + public AvroReflectDatumBackedKafkaDecoder(final Class clazz) { + this.reader = new ReflectDatumReader(clazz); } @Override @SuppressWarnings("unchecked") public T fromBytes(final byte[] bytes) { - final Schema schema = ReflectData.get().getSchema(clazz); - final AvroSerializer avroSerializer = new AvroSerializer(); - - try { - return (T) avroSerializer.deserialize(bytes, schema); - } catch (IOException e) { - LOG.error("Failed to decode byte array for schema: " + schema.getFullName(), e); - } - - return null; + return fromBytes(bytes, reader); } } diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/avro/AvroBackedKafkaEncoder.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/avro/AvroReflectDatumBackedKafkaEncoder.java similarity index 58% rename from spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/avro/AvroBackedKafkaEncoder.java rename to spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/avro/AvroReflectDatumBackedKafkaEncoder.java index 6ec8bfe517..6f1c0cc4df 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/avro/AvroBackedKafkaEncoder.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/avro/AvroReflectDatumBackedKafkaEncoder.java @@ -16,10 +16,8 @@ package org.springframework.integration.kafka.serializer.avro; import kafka.serializer.Encoder; -import org.apache.avro.Schema; -import org.apache.avro.reflect.ReflectData; - -import java.io.IOException; +import org.apache.avro.io.DatumWriter; +import org.apache.avro.reflect.ReflectDatumWriter; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; @@ -27,27 +25,18 @@ import org.apache.commons.logging.LogFactory; * @author Soby Chacko * @since 0.5 */ -public class AvroBackedKafkaEncoder implements Encoder { - private static final Log LOG = LogFactory.getLog(AvroBackedKafkaEncoder.class); +public class AvroReflectDatumBackedKafkaEncoder extends AvroDatumSupport implements Encoder { + private static final Log LOG = LogFactory.getLog(AvroReflectDatumBackedKafkaEncoder.class); - private final Class clazz; + private final DatumWriter writer; - public AvroBackedKafkaEncoder(final Class clazz) { - this.clazz = clazz; + public AvroReflectDatumBackedKafkaEncoder(final Class clazz) { + this.writer = new ReflectDatumWriter(clazz); } @Override @SuppressWarnings("unchecked") public byte[] toBytes(final T source) { - final Schema schema = ReflectData.get().getSchema(clazz); - final AvroSerializer avroSerializer = new AvroSerializer(); - - try { - return avroSerializer.serialize(source, schema); - } catch (IOException e) { - LOG.error("Failed to encode source for schema: " + schema.getFullName()); - } - - return null; + return toBytes(source, writer); } } diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/avro/AvroSerializer.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/avro/AvroSerializer.java index 7bc08f8f8a..2a98894c58 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/avro/AvroSerializer.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/avro/AvroSerializer.java @@ -15,15 +15,12 @@ */ package org.springframework.integration.kafka.serializer.avro; -import org.apache.avro.Schema; import org.apache.avro.io.DatumReader; import org.apache.avro.io.DatumWriter; import org.apache.avro.io.Decoder; import org.apache.avro.io.DecoderFactory; import org.apache.avro.io.Encoder; import org.apache.avro.io.EncoderFactory; -import org.apache.avro.reflect.ReflectDatumReader; -import org.apache.avro.reflect.ReflectDatumWriter; import java.io.ByteArrayOutputStream; import java.io.IOException; @@ -33,15 +30,13 @@ import java.io.IOException; * @since 0.5 */ public class AvroSerializer { - public T deserialize(final byte[] bytes, final Schema schema) throws IOException { - final Decoder decoder = DecoderFactory.get().binaryDecoder(bytes, null); - final DatumReader reader = new ReflectDatumReader(schema); + public T deserialize(final byte[] bytes, final DatumReader reader) throws IOException { + final Decoder decoder = DecoderFactory.get().binaryDecoder(bytes, null); return reader.read(null, decoder); } - public byte[] serialize(final T input, final Schema schema) throws IOException { - final DatumWriter writer = new ReflectDatumWriter(schema); + public byte[] serialize(final T input, final DatumWriter writer) throws IOException { final ByteArrayOutputStream stream = new ByteArrayOutputStream(); final Encoder encoder = EncoderFactory.get().binaryEncoder(stream, null); diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/avro/AvroSpecificDatumBackedKafkaDecoder.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/avro/AvroSpecificDatumBackedKafkaDecoder.java new file mode 100644 index 0000000000..622a9e38ca --- /dev/null +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/avro/AvroSpecificDatumBackedKafkaDecoder.java @@ -0,0 +1,28 @@ +package org.springframework.integration.kafka.serializer.avro; + +import kafka.serializer.Decoder; +import org.apache.avro.io.DatumReader; +import org.apache.avro.specific.SpecificDatumReader; +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; + +/** + * @author Soby Chacko + * @since 0.5 + */ +public class AvroSpecificDatumBackedKafkaDecoder extends AvroDatumSupport implements Decoder { + + private static final Log LOG = LogFactory.getLog(AvroSpecificDatumBackedKafkaDecoder.class); + + private final DatumReader reader; + + public AvroSpecificDatumBackedKafkaDecoder(final Class specificRecordBase) { + this.reader = new SpecificDatumReader(specificRecordBase); + } + + @Override + @SuppressWarnings("unchecked") + public T fromBytes(final byte[] bytes) { + return fromBytes(bytes, reader); + } +} diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/avro/AvroSpecificDatumBackedKafkaEncoder.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/avro/AvroSpecificDatumBackedKafkaEncoder.java new file mode 100644 index 0000000000..410dfd6047 --- /dev/null +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/avro/AvroSpecificDatumBackedKafkaEncoder.java @@ -0,0 +1,28 @@ +package org.springframework.integration.kafka.serializer.avro; + +import kafka.serializer.Encoder; +import org.apache.avro.io.DatumWriter; +import org.apache.avro.specific.SpecificDatumWriter; +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; + +/** + * @author Soby Chacko + * @since 0.5 + */ +public class AvroSpecificDatumBackedKafkaEncoder extends AvroDatumSupport implements Encoder { + + private static final Log LOG = LogFactory.getLog(AvroSpecificDatumBackedKafkaEncoder.class); + + private final DatumWriter writer; + + public AvroSpecificDatumBackedKafkaEncoder(final Class specificRecordClazz) { + this.writer = new SpecificDatumWriter(specificRecordClazz); + } + + @Override + @SuppressWarnings("unchecked") + public byte[] toBytes(final T source) { + return toBytes(source, writer); + } +} diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/avro/AvroSpecificDatumSerializer.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/avro/AvroSpecificDatumSerializer.java deleted file mode 100644 index b7b6064023..0000000000 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/avro/AvroSpecificDatumSerializer.java +++ /dev/null @@ -1,53 +0,0 @@ -/* - * Copyright 2002-2013 the original author or authors. - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.springframework.integration.kafka.serializer.avro; - -import org.apache.avro.Schema; -import org.apache.avro.io.DatumReader; -import org.apache.avro.io.DatumWriter; -import org.apache.avro.io.Decoder; -import org.apache.avro.io.DecoderFactory; -import org.apache.avro.io.Encoder; -import org.apache.avro.io.EncoderFactory; -import org.apache.avro.specific.SpecificDatumReader; -import org.apache.avro.specific.SpecificDatumWriter; - -import java.io.ByteArrayOutputStream; -import java.io.IOException; - -/** - * @author Soby Chacko - * @since 0.5 - */ -public class AvroSpecificDatumSerializer { - public T deserialize(final byte[] bytes, final Schema schema) throws IOException { - final Decoder decoder = DecoderFactory.get().binaryDecoder(bytes, null); - final DatumReader reader = new SpecificDatumReader(schema); - - return reader.read(null, decoder); - } - - public byte[] serialize(final T input, final Schema schema) throws IOException { - final DatumWriter writer = new SpecificDatumWriter(schema); - final ByteArrayOutputStream stream = new ByteArrayOutputStream(); - - final Encoder encoder = EncoderFactory.get().binaryEncoder(stream, null); - writer.write(input, encoder); - encoder.flush(); - - return stream.toByteArray(); - } -} diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/common/StringEncoder.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/common/StringEncoder.java index 7dc5105bb3..1131c4002c 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/common/StringEncoder.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/serializer/common/StringEncoder.java @@ -24,7 +24,7 @@ import java.util.Properties; * @author Soby Chacko * @since 0.5 */ -public class StringEncoder implements Encoder { +public class StringEncoder implements Encoder { private String encoding = "UTF8"; public void setEncoding(final String encoding){ 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 793162fbef..d8def7d99f 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 @@ -24,12 +24,12 @@ import java.util.Properties; * @author Soby Chacko * @since 0.5 */ -public class ConsumerConfigFactoryBean implements FactoryBean { +public class ConsumerConfigFactoryBean implements FactoryBean { - private final ConsumerMetadata consumerMetadata; + private final ConsumerMetadata consumerMetadata; private final ZookeeperConnect zookeeperConnect; - public ConsumerConfigFactoryBean(final ConsumerMetadata consumerMetadata, + public ConsumerConfigFactoryBean(final ConsumerMetadata consumerMetadata, final ZookeeperConnect zookeeperConnect){ this.consumerMetadata = consumerMetadata; this.zookeeperConnect = zookeeperConnect; 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 0e77998436..25fe67fe58 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 @@ -37,46 +37,46 @@ import java.util.concurrent.Future; * @author Soby Chacko * @since 0.5 */ -public class ConsumerConfiguration { +public class ConsumerConfiguration { private static final Log LOGGER = LogFactory.getLog(ConsumerConfiguration.class); - private final ConsumerMetadata consumerMetadata; + private final ConsumerMetadata consumerMetadata; private final ConsumerConnectionProvider consumerConnectionProvider; - private final MessageLeftOverTracker messageLeftOverTracker; + private final MessageLeftOverTracker messageLeftOverTracker; private ConsumerConnector consumerConnector; private volatile int count = 0; private int maxMessages = 1; private ExecutorService executorService = Executors.newCachedThreadPool(); - public ConsumerConfiguration(final ConsumerMetadata consumerMetadata, + public ConsumerConfiguration(final ConsumerMetadata consumerMetadata, final ConsumerConnectionProvider consumerConnectionProvider, - final MessageLeftOverTracker messageLeftOverTracker) { + final MessageLeftOverTracker messageLeftOverTracker) { this.consumerMetadata = consumerMetadata; this.consumerConnectionProvider = consumerConnectionProvider; this.messageLeftOverTracker = messageLeftOverTracker; } - public ConsumerMetadata getConsumerMetadata() { + public ConsumerMetadata getConsumerMetadata() { return consumerMetadata; } public Map>> receive() { count = messageLeftOverTracker.getCurrentCount(); - - final List>> tasks = new LinkedList>>(); final Object lock = new Object(); - final Map>> consumerMap = getConsumerMapWithMessageStreams(); - for (final List> streams : consumerMap.values()) { - for (final KafkaStream stream : streams) { - tasks.add(new Callable>() { + final List>>> tasks = new LinkedList>>>(); + + final Map>> consumerMap = getConsumerMapWithMessageStreams(); + for (final List> streams : consumerMap.values()) { + for (final KafkaStream stream : streams) { + tasks.add(new Callable>>() { @Override - public List call() throws Exception { - final List rawMessages = new ArrayList(); + public List> call() throws Exception { + final List> rawMessages = new ArrayList>(); try { while (count < maxMessages) { - final MessageAndMetadata messageAndMetadata = stream.iterator().next(); + final MessageAndMetadata messageAndMetadata = stream.iterator().next(); synchronized (lock) { if (count < maxMessages) { rawMessages.add(messageAndMetadata); @@ -94,17 +94,16 @@ public class ConsumerConfiguration { }); } } - return executeTasks(tasks); } - private Map>> executeTasks(final List>> tasks) { + private Map>> executeTasks(final List>>> tasks) { final Map>> messages = new ConcurrentHashMap>>(); messages.putAll(getLeftOverMessageMap()); try { - for (final Future> result : executorService.invokeAll(tasks)) { + for (final Future>> result : executorService.invokeAll(tasks)) { if (!result.get().isEmpty()) { final String topic = result.get().get(0).topic(); if (!messages.containsKey(topic)) { @@ -127,20 +126,21 @@ public class ConsumerConfiguration { return messages; } + @SuppressWarnings("unchecked") private Map>> getLeftOverMessageMap() { final Map>> messages = new ConcurrentHashMap>>(); - for (final MessageAndMetadata mamd : messageLeftOverTracker.getMessageLeftOverFromPreviousPoll()) { + for (final MessageAndMetadata mamd : messageLeftOverTracker.getMessageLeftOverFromPreviousPoll()) { final String topic = mamd.topic(); if (!messages.containsKey(topic)) { - final List l = new ArrayList(); + final List> l = new ArrayList>(); l.add(mamd); messages.put(topic, getPayload(l)); } else { final Map> existingPayloadMap = messages.get(topic); - final List l = new ArrayList(); + final List> l = new ArrayList>(); l.add(mamd); getPayload(l, existingPayloadMap); } @@ -149,10 +149,10 @@ public class ConsumerConfiguration { return messages; } - private Map> getPayload(final List messageAndMetadatas) { + private Map> getPayload(final List> messageAndMetadatas) { final Map> payloadMap = new ConcurrentHashMap>(); - for (final MessageAndMetadata messageAndMetadata : messageAndMetadatas) { + for (final MessageAndMetadata messageAndMetadata : messageAndMetadatas) { if (!payloadMap.containsKey(messageAndMetadata.partition())) { final List payload = new ArrayList(); payload.add(messageAndMetadata.message()); @@ -167,8 +167,8 @@ public class ConsumerConfiguration { return payloadMap; } - private void getPayload(final List messageAndMetadatas, final Map> existingPayloadMap) { - for (final MessageAndMetadata messageAndMetadata : messageAndMetadatas) { + private void getPayload(final List> messageAndMetadatas, final Map> existingPayloadMap) { + for (final MessageAndMetadata messageAndMetadata : messageAndMetadatas) { if (!existingPayloadMap.containsKey(messageAndMetadata.partition())) { final List payload = new ArrayList(); payload.add(messageAndMetadata.message()); @@ -181,16 +181,11 @@ public class ConsumerConfiguration { } @SuppressWarnings("unchecked") - public Map>> getConsumerMapWithMessageStreams() { - if (consumerMetadata.getValueDecoder() != null && - consumerMetadata.getKeyDecoder() != null) { - return getConsumerConnector().createMessageStreams( - consumerMetadata.getTopicStreamMap(), - consumerMetadata.getKeyDecoder(), - consumerMetadata.getValueDecoder()); - } - - return getConsumerConnector().createMessageStreams(consumerMetadata.getTopicStreamMap()); + public Map>> getConsumerMapWithMessageStreams() { + return getConsumerConnector().createMessageStreams( + consumerMetadata.getTopicStreamMap(), + consumerMetadata.getKeyDecoder(), + consumerMetadata.getValueDecoder()); } public int getMaxMessages() { @@ -205,7 +200,6 @@ public class ConsumerConfiguration { if (consumerConnector == null) { consumerConnector = consumerConnectionProvider.getConsumerConnector(); } - return consumerConnector; } } diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ConsumerMetadata.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ConsumerMetadata.java index d906bd7b65..476e66a461 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ConsumerMetadata.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ConsumerMetadata.java @@ -16,6 +16,8 @@ package org.springframework.integration.kafka.support; import kafka.serializer.Decoder; +import kafka.serializer.DefaultDecoder; +import org.springframework.beans.factory.InitializingBean; import org.springframework.integration.kafka.core.KafkaConsumerDefaults; import java.util.Map; @@ -24,7 +26,7 @@ import java.util.Map; * @author Soby Chacko * @since 0.5 */ -public class ConsumerMetadata { +public class ConsumerMetadata implements InitializingBean { //High level consumer defaults private String groupId = KafkaConsumerDefaults.GROUP_ID; @@ -172,4 +174,16 @@ public class ConsumerMetadata { public void setTopicStreamMap(final Map topicStreamMap) { this.topicStreamMap = topicStreamMap; } + + @Override + @SuppressWarnings("unchecked") + public void afterPropertiesSet() throws Exception { + if (valueDecoder == null) { + setValueDecoder((Decoder) new DefaultDecoder(null)); + } + + if (keyDecoder == null) { + setKeyDecoder((Decoder) getValueDecoder()); + } + } } 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 a75b6ab296..6822ccfa7f 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 @@ -32,24 +32,26 @@ import java.util.Map; * @author Soby Chacko * @since 0.5 */ -public class KafkaConsumerContext implements BeanFactoryAware { - private Map consumerConfigurations; +public class KafkaConsumerContext implements BeanFactoryAware { + private Map> consumerConfigurations; private String consumerTimeout = KafkaConsumerDefaults.CONSUMER_TIMEOUT; private ZookeeperConnect zookeeperConnect; - public Collection getConsumerConfigurations() { + public Collection> getConsumerConfigurations() { return consumerConfigurations.values(); } @Override + @SuppressWarnings("unchecked") public void setBeanFactory(final BeanFactory beanFactory) throws BeansException { - consumerConfigurations = ((ListableBeanFactory) beanFactory).getBeansOfType(ConsumerConfiguration.class); + consumerConfigurations = (Map>) + (Object) ((ListableBeanFactory) beanFactory).getBeansOfType(ConsumerConfiguration.class); } public Message>>> receive() { final Map>> consumedData = new HashMap>>(); - for (final ConsumerConfiguration consumerConfiguration : getConsumerConfigurations()) { + for (final ConsumerConfiguration consumerConfiguration : getConsumerConfigurations()) { final Map>> messages = consumerConfiguration.receive(); if (messages != null){ 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 f85207dbca..bdb7e1183f 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 @@ -28,12 +28,12 @@ import java.util.Map; * @author Soby Chacko * @since 0.5 */ -public class KafkaProducerContext implements BeanFactoryAware { - private Map topicsConfiguration; +public class KafkaProducerContext implements BeanFactoryAware { + private Map> topicsConfiguration; @SuppressWarnings("unchecked") public void send(final Message message) throws Exception { - final ProducerConfiguration producerConfiguration = + final ProducerConfiguration producerConfiguration = getTopicConfiguration(message.getHeaders().get("topic", String.class)); if (producerConfiguration != null) { @@ -41,10 +41,10 @@ public class KafkaProducerContext implements BeanFactoryAware { } } - private ProducerConfiguration getTopicConfiguration(final String topic){ - final Collection topics = topicsConfiguration.values(); + private ProducerConfiguration getTopicConfiguration(final String topic){ + final Collection> topics = topicsConfiguration.values(); - for (final ProducerConfiguration producerConfiguration : topics){ + for (final ProducerConfiguration producerConfiguration : topics){ if (producerConfiguration.getProducerMetadata().getTopic().equals(topic)){ return producerConfiguration; } @@ -53,12 +53,15 @@ public class KafkaProducerContext implements BeanFactoryAware { return null; } - public Map getTopicsConfiguration() { + public Map> getTopicsConfiguration() { return topicsConfiguration; } @Override + @SuppressWarnings("unchecked") public void setBeanFactory(final BeanFactory beanFactory) throws BeansException { - topicsConfiguration = ((ListableBeanFactory)beanFactory).getBeansOfType(ProducerConfiguration.class); + topicsConfiguration = + (Map>) (Object) + ((ListableBeanFactory)beanFactory).getBeansOfType(ProducerConfiguration.class); } } diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/MessageLeftOverTracker.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/MessageLeftOverTracker.java index 82f3461c09..bdeb044806 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/MessageLeftOverTracker.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/MessageLeftOverTracker.java @@ -24,14 +24,14 @@ import java.util.List; * @author Soby Chacko * @since 0.5 */ -public class MessageLeftOverTracker { - private final List messageLeftOverFromPreviousPoll = new ArrayList(); +public class MessageLeftOverTracker { + private final List> messageLeftOverFromPreviousPoll = new ArrayList>(); - public void addMessageAndMetadata(final MessageAndMetadata messageAndMetadata){ + public void addMessageAndMetadata(final MessageAndMetadata messageAndMetadata){ messageLeftOverFromPreviousPoll.add(messageAndMetadata); } - public List getMessageLeftOverFromPreviousPoll(){ + public List> getMessageLeftOverFromPreviousPoll(){ return messageLeftOverFromPreviousPoll; } diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaConsumerContextParserTests-context.xml b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaConsumerContextParserTests-context.xml index e0978bbf19..576dabd3c0 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaConsumerContextParserTests-context.xml +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaConsumerContextParserTests-context.xml @@ -10,7 +10,7 @@ zk-sync-time="2000"/> - + diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaConsumerContextParserTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaConsumerContextParserTests.java index 4b46b4ec6b..d91210c2a1 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaConsumerContextParserTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaConsumerContextParserTests.java @@ -31,7 +31,7 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; */ @RunWith(SpringJUnit4ClassRunner.class) @ContextConfiguration -public class KafkaConsumerContextParserTests { +public class KafkaConsumerContextParserTests { @Autowired private ApplicationContext appContext; @@ -39,10 +39,10 @@ public class KafkaConsumerContextParserTests { @Test @SuppressWarnings("unchecked") public void testConsumerContextConfiguration() { - final KafkaConsumerContext consumerContext = appContext.getBean("consumerContext", KafkaConsumerContext.class); + final KafkaConsumerContext consumerContext = appContext.getBean("consumerContext", KafkaConsumerContext.class); Assert.assertNotNull(consumerContext); - final ConsumerMetadata cm = appContext.getBean("consumerMetadata_default1", ConsumerMetadata.class); + final ConsumerMetadata cm = appContext.getBean("consumerMetadata_default1", ConsumerMetadata.class); Assert.assertNotNull(cm); } } diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaInboundAdapterParserTests-context.xml b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaInboundAdapterParserTests-context.xml index 70266d3684..b1ac3cfa4e 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaInboundAdapterParserTests-context.xml +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaInboundAdapterParserTests-context.xml @@ -20,7 +20,7 @@ - + diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaOutboundAdapterParserTests-context.xml b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaOutboundAdapterParserTests-context.xml index 52eef96071..89df6ff229 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaOutboundAdapterParserTests-context.xml +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaOutboundAdapterParserTests-context.xml @@ -23,7 +23,7 @@ - + diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaOutboundAdapterParserTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaOutboundAdapterParserTests.java index 504e2986e3..d2b0909cb6 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaOutboundAdapterParserTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaOutboundAdapterParserTests.java @@ -32,18 +32,19 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; */ @RunWith(SpringJUnit4ClassRunner.class) @ContextConfiguration -public class KafkaOutboundAdapterParserTests { +public class KafkaOutboundAdapterParserTests { @Autowired private ApplicationContext appContext; @Test + @SuppressWarnings("unchecked") public void testOutboundAdapterConfiguration(){ final PollingConsumer pollingConsumer = appContext.getBean("kafkaOutboundChannelAdapter", PollingConsumer.class); - final KafkaProducerMessageHandler messageHandler = appContext.getBean(KafkaProducerMessageHandler.class); + final KafkaProducerMessageHandler messageHandler = appContext.getBean(KafkaProducerMessageHandler.class); Assert.assertNotNull(pollingConsumer); Assert.assertNotNull(messageHandler); - final KafkaProducerContext producerContext = messageHandler.getKafkaProducerContext(); + final KafkaProducerContext producerContext = messageHandler.getKafkaProducerContext(); Assert.assertNotNull(producerContext); Assert.assertEquals(producerContext.getTopicsConfiguration().size(), 2); } diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaProducerContextParserTests-context.xml b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaProducerContextParserTests-context.xml index 8f9750af79..6316866b91 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaProducerContextParserTests-context.xml +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaProducerContextParserTests-context.xml @@ -20,7 +20,7 @@ - + diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaProducerContextParserTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaProducerContextParserTests.java index 411d02dbaa..fa9241bbac 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaProducerContextParserTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaProducerContextParserTests.java @@ -36,7 +36,7 @@ import java.util.Map; */ @RunWith(SpringJUnit4ClassRunner.class) @ContextConfiguration -public class KafkaProducerContextParserTests { +public class KafkaProducerContextParserTests { @Autowired private ApplicationContext appContext; @@ -44,34 +44,34 @@ public class KafkaProducerContextParserTests { @Test @SuppressWarnings("unchecked") public void testProducerContextConfiguration(){ - final KafkaProducerContext producerContext = appContext.getBean("producerContext", KafkaProducerContext.class); + final KafkaProducerContext producerContext = appContext.getBean("producerContext", KafkaProducerContext.class); Assert.assertNotNull(producerContext); - final Map topicConfigurations = producerContext.getTopicsConfiguration(); + final Map> topicConfigurations = producerContext.getTopicsConfiguration(); Assert.assertEquals(topicConfigurations.size(), 2); - final ProducerConfiguration producerConfigurationTest1 = topicConfigurations.get("producerConfiguration_test1"); + final ProducerConfiguration producerConfigurationTest1 = topicConfigurations.get("producerConfiguration_test1"); Assert.assertNotNull(producerConfigurationTest1); - final ProducerMetadata producerMetadataTest1 = producerConfigurationTest1.getProducerMetadata(); + final ProducerMetadata producerMetadataTest1 = producerConfigurationTest1.getProducerMetadata(); Assert.assertEquals(producerMetadataTest1.getTopic(), "test1"); Assert.assertEquals(producerMetadataTest1.getCompressionCodec(), "0"); Assert.assertEquals(producerMetadataTest1.getKeyClassType(), java.lang.String.class); Assert.assertEquals(producerMetadataTest1.getValueClassType(), java.lang.String.class); - final Encoder valueEncoder = appContext.getBean("valueEncoder", Encoder.class); + final Encoder valueEncoder = appContext.getBean("valueEncoder", Encoder.class); Assert.assertEquals(producerMetadataTest1.getValueEncoder(), valueEncoder); Assert.assertEquals(producerMetadataTest1.getKeyEncoder(), valueEncoder); - final Producer producerTest1 = appContext.getBean("prodFactory_test1", Producer.class); - Assert.assertEquals(producerConfigurationTest1, new ProducerConfiguration(producerMetadataTest1, producerTest1)); + final Producer producerTest1 = appContext.getBean("prodFactory_test1", Producer.class); + Assert.assertEquals(producerConfigurationTest1, new ProducerConfiguration(producerMetadataTest1, producerTest1)); - final ProducerConfiguration producerConfigurationTest2 = topicConfigurations.get("producerConfiguration_" + "test2"); + final ProducerConfiguration producerConfigurationTest2 = topicConfigurations.get("producerConfiguration_" + "test2"); Assert.assertNotNull(producerConfigurationTest2); - final ProducerMetadata producerMetadataTest2 = producerConfigurationTest2.getProducerMetadata(); + final ProducerMetadata producerMetadataTest2 = producerConfigurationTest2.getProducerMetadata(); Assert.assertEquals(producerMetadataTest2.getTopic(), "test2"); Assert.assertEquals(producerMetadataTest2.getCompressionCodec(), "0"); - final Producer producerTest2 = appContext.getBean("prodFactory_test2", Producer.class); - Assert.assertEquals(producerConfigurationTest2, new ProducerConfiguration(producerMetadataTest2, producerTest2)); + final Producer producerTest2 = appContext.getBean("prodFactory_test2", Producer.class); + Assert.assertEquals(producerConfigurationTest2, new ProducerConfiguration(producerMetadataTest2, producerTest2)); } } diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/serializer/AvroBackedKafkaSerializerTest.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/serializer/AvroReflectDatumBackedKafkaSerializerTest.java similarity index 64% rename from spring-integration-kafka/src/test/java/org/springframework/integration/kafka/serializer/AvroBackedKafkaSerializerTest.java rename to spring-integration-kafka/src/test/java/org/springframework/integration/kafka/serializer/AvroReflectDatumBackedKafkaSerializerTest.java index 6afea82fd8..224c42d48c 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/serializer/AvroBackedKafkaSerializerTest.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/serializer/AvroReflectDatumBackedKafkaSerializerTest.java @@ -17,19 +17,19 @@ package org.springframework.integration.kafka.serializer; import org.junit.Assert; import org.junit.Test; -import org.springframework.integration.kafka.serializer.avro.AvroBackedKafkaDecoder; -import org.springframework.integration.kafka.serializer.avro.AvroBackedKafkaEncoder; +import org.springframework.integration.kafka.serializer.avro.AvroReflectDatumBackedKafkaDecoder; +import org.springframework.integration.kafka.serializer.avro.AvroReflectDatumBackedKafkaEncoder; import org.springframework.integration.kafka.test.utils.TestObject; /** * @author Soby Chacko * @since 0.5 */ -public class AvroBackedKafkaSerializerTest { +public class AvroReflectDatumBackedKafkaSerializerTest { @Test @SuppressWarnings("unchecked") public void testDecodePlainSchema() { - final AvroBackedKafkaEncoder avroBackedKafkaEncoder = new AvroBackedKafkaEncoder(TestObject.class); + final AvroReflectDatumBackedKafkaEncoder avroBackedKafkaEncoder = new AvroReflectDatumBackedKafkaEncoder(TestObject.class); final TestObject testObject = new TestObject(); testObject.setTestData1("\"Test Data1\""); @@ -37,8 +37,8 @@ public class AvroBackedKafkaSerializerTest { final byte[] data = avroBackedKafkaEncoder.toBytes(testObject); - final AvroBackedKafkaDecoder avroBackedKafkaDecoder = new AvroBackedKafkaDecoder(TestObject.class); - final TestObject decodedFbu = (TestObject) avroBackedKafkaDecoder.fromBytes(data); + final AvroReflectDatumBackedKafkaDecoder avroReflectDatumBackedKafkaDecoder = new AvroReflectDatumBackedKafkaDecoder(TestObject.class); + final TestObject decodedFbu = avroReflectDatumBackedKafkaDecoder.fromBytes(data); Assert.assertEquals(testObject.getTestData1(), decodedFbu.getTestData1()); Assert.assertEquals(testObject.getTestData2(), decodedFbu.getTestData2()); @@ -47,12 +47,12 @@ public class AvroBackedKafkaSerializerTest { @Test @SuppressWarnings("unchecked") public void anotherTest() { - final AvroBackedKafkaEncoder avroBackedKafkaEncoder = new AvroBackedKafkaEncoder(java.lang.String.class); + final AvroReflectDatumBackedKafkaEncoder avroBackedKafkaEncoder = new AvroReflectDatumBackedKafkaEncoder(java.lang.String.class); final String testString = "Testing Avro"; final byte[] data = avroBackedKafkaEncoder.toBytes(testString); - final AvroBackedKafkaDecoder avroBackedKafkaDecoder = new AvroBackedKafkaDecoder(java.lang.String.class); - final String decodedS = (String) avroBackedKafkaDecoder.fromBytes(data); + final AvroReflectDatumBackedKafkaDecoder avroReflectDatumBackedKafkaDecoder = new AvroReflectDatumBackedKafkaDecoder(java.lang.String.class); + final String decodedS = avroReflectDatumBackedKafkaDecoder.fromBytes(data); Assert.assertEquals(testString, decodedS); } diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/serializer/AvroSpecificDatumBackedKafkaSerializerTest.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/serializer/AvroSpecificDatumBackedKafkaSerializerTest.java new file mode 100644 index 0000000000..8d18befd51 --- /dev/null +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/serializer/AvroSpecificDatumBackedKafkaSerializerTest.java @@ -0,0 +1,30 @@ +package org.springframework.integration.kafka.serializer; + +import org.junit.Assert; +import org.junit.Test; +import org.springframework.integration.kafka.serializer.avro.AvroSpecificDatumBackedKafkaDecoder; +import org.springframework.integration.kafka.serializer.avro.AvroSpecificDatumBackedKafkaEncoder; +import org.springframework.integration.kafka.test.utils.User; + +/** + * @author Soby Chacko + * @since 0.5 + */ +public class AvroSpecificDatumBackedKafkaSerializerTest { + + @Test + @SuppressWarnings("unchecked") + public void testEncodeDecodeFromSpecificDatumSchema() { + final AvroSpecificDatumBackedKafkaEncoder avroBackedKafkaEncoder = new AvroSpecificDatumBackedKafkaEncoder(User.class); + + final User user = new User("First", "Last"); + + final byte[] data = avroBackedKafkaEncoder.toBytes(user); + + final AvroSpecificDatumBackedKafkaDecoder avroSpecificDatumBackedKafkaDecoder = new AvroSpecificDatumBackedKafkaDecoder(User.class); + final User decodedUser = avroSpecificDatumBackedKafkaDecoder.fromBytes(data); + + Assert.assertEquals(user.getFirstName(), decodedUser.getFirstName().toString()); + Assert.assertEquals(user.getLastName(), decodedUser.getLastName().toString()); + } +} 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 f6f4f2d9b7..19c267a9aa 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 @@ -45,34 +45,34 @@ import org.mockito.stubbing.Answer; * @author Gunnar Hillert * @since 0.5 */ -public class ConsumerConfigurationTests { +public class ConsumerConfigurationTests { @Test @SuppressWarnings("unchecked") public void testReceiveMessageForSingleTopicFromSingleStream() { - 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(consumerConnectionProvider.getConsumerConnector()).thenReturn(consumerConnector); - final ConsumerConfiguration consumerConfiguration = new ConsumerConfiguration(consumerMetadata, + 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); - final Map>> messageStreams = new HashMap>>(); + final Map>> messageStreams = new HashMap>>(); messageStreams.put("topic", streams); when(consumerConfiguration.getConsumerMapWithMessageStreams()).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("topic"); when(messageAndMetadata.partition()).thenReturn(1); @@ -90,54 +90,54 @@ public class ConsumerConfigurationTests { @Test @SuppressWarnings("unchecked") public void testReceiveMessageForSingleTopicFromMultipleStreams() { - 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(consumerConnectionProvider.getConsumerConnector()).thenReturn(consumerConnector); - final ConsumerConfiguration consumerConfiguration = new ConsumerConfiguration(consumerMetadata, + 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 List> streams = new ArrayList>(); + 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); - final Map>> messageStreams = new HashMap>>(); + final Map>> messageStreams = new HashMap>>(); messageStreams.put("topic", streams); when(consumerConfiguration.getConsumerMapWithMessageStreams()).thenReturn(messageStreams); - 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((V)"got message"); when(messageAndMetadata1.topic()).thenReturn("topic"); when(messageAndMetadata1.partition()).thenReturn(1); - when(messageAndMetadata2.message()).thenReturn("got message"); + when(messageAndMetadata2.message()).thenReturn((V)"got message"); when(messageAndMetadata2.topic()).thenReturn("topic"); when(messageAndMetadata2.partition()).thenReturn(2); - when(messageAndMetadata3.message()).thenReturn("got message"); + when(messageAndMetadata3.message()).thenReturn((V)"got message"); when(messageAndMetadata3.topic()).thenReturn("topic"); when(messageAndMetadata3.partition()).thenReturn(3); @@ -157,56 +157,56 @@ public class ConsumerConfigurationTests { @Test @SuppressWarnings("unchecked") public void testReceiveMessageForMultipleTopicsFromMultipleStreams() { - 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(consumerConnectionProvider.getConsumerConnector()).thenReturn(consumerConnector); - final ConsumerConfiguration consumerConfiguration = new ConsumerConfiguration(consumerMetadata, + final ConsumerConfiguration consumerConfiguration = new ConsumerConfiguration(consumerMetadata, consumerConnectionProvider, messageLeftOverTracker); consumerConfiguration.setMaxMessages(9); - final KafkaStream stream1 = mock(KafkaStream.class); - final KafkaStream stream2 = mock(KafkaStream.class); - final KafkaStream stream3 = mock(KafkaStream.class); - final List> streams = new ArrayList>(); + 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); - final Map>> messageStreams = new HashMap>>(); + final Map>> messageStreams = new HashMap>>(); messageStreams.put("topic1", streams); messageStreams.put("topic2", streams); messageStreams.put("topic3", streams); when(consumerConfiguration.getConsumerMapWithMessageStreams()).thenReturn(messageStreams); - 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 message1"); + when(messageAndMetadata1.message()).thenReturn((V)"got message1"); when(messageAndMetadata1.topic()).thenReturn("topic1"); when(messageAndMetadata1.partition()).thenAnswer(getAnswer()); - when(messageAndMetadata2.message()).thenReturn("got message2"); + when(messageAndMetadata2.message()).thenReturn((V)"got message2"); when(messageAndMetadata2.topic()).thenReturn("topic2"); when(messageAndMetadata1.partition()).thenAnswer(getAnswer()); - when(messageAndMetadata3.message()).thenReturn("got message3"); + when(messageAndMetadata3.message()).thenReturn((V)"got message3"); when(messageAndMetadata3.topic()).thenReturn("topic3"); when(messageAndMetadata1.partition()).thenAnswer(getAnswer()); @@ -244,39 +244,39 @@ public class ConsumerConfigurationTests { @Test @SuppressWarnings("unchecked") public void testReceiveMessageAndVerifyMessageLeftoverFromPreviousPollAreTakenFirst() { - 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(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 = 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 List mList = new ArrayList(); + final List> mList = new ArrayList>(); mList.add(m1); mList.add(m2); mList.add(m3); - when(messageLeftOverTracker.getMessageLeftOverFromPreviousPoll()).thenReturn(mList); + when((List>) (Object) messageLeftOverTracker.getMessageLeftOverFromPreviousPoll()).thenReturn(mList); when(consumerConnectionProvider.getConsumerConnector()).thenReturn(consumerConnector); - final ConsumerConfiguration consumerConfiguration = new ConsumerConfiguration(consumerMetadata, + final ConsumerConfiguration consumerConfiguration = new ConsumerConfiguration(consumerMetadata, consumerConnectionProvider, messageLeftOverTracker); consumerConfiguration.setMaxMessages(5); - 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); - final Map>> messageStreams = new HashMap>>(); + final Map>> messageStreams = new HashMap>>(); messageStreams.put("topic1", streams); when(consumerConfiguration.getConsumerMapWithMessageStreams()).thenReturn(messageStreams); final ConsumerIterator iterator = mock(ConsumerIterator.class); - when(stream.iterator()).thenReturn(iterator); + when(stream.iterator()).thenReturn((ConsumerIterator) iterator); final MessageAndMetadata messageAndMetadata = mock(MessageAndMetadata.class); when(iterator.next()).thenReturn(messageAndMetadata); when(messageAndMetadata.message()).thenReturn("got message"); @@ -306,9 +306,10 @@ public class ConsumerConfigurationTests { } @Test + @SuppressWarnings("unchecked") public void testGetConsumerMapWithMessageStreamsWithNullDecoders() { - final ConsumerMetadata mockedConsumerMetadata = mock(ConsumerMetadata.class); + final ConsumerMetadata mockedConsumerMetadata = mock(ConsumerMetadata.class); assertNull(mockedConsumerMetadata.getKeyDecoder()); assertNull(mockedConsumerMetadata.getValueDecoder()); @@ -317,25 +318,26 @@ public class ConsumerConfigurationTests { when(mockedConsumerMetadata.getTopicStreamMap()).thenReturn(topicsStreamMap); final ConsumerConnectionProvider mockedConsumerConnectionProvider = mock(ConsumerConnectionProvider.class); - final MessageLeftOverTracker mockedMessageLeftOverTracker = mock(MessageLeftOverTracker.class); + final MessageLeftOverTracker mockedMessageLeftOverTracker = mock(MessageLeftOverTracker.class); final ConsumerConnector mockedConsumerConnector = mock(ConsumerConnector.class); when(mockedConsumerConnectionProvider.getConsumerConnector()).thenReturn(mockedConsumerConnector); - final Map>> messageStreams = new HashMap>>(); - when(mockedConsumerConnector.createMessageStreams(topicsStreamMap)).thenReturn(messageStreams); + final Map>> messageStreams = new HashMap>>(); + when((Map>>) (Object) mockedConsumerConnector.createMessageStreams(topicsStreamMap)).thenReturn(messageStreams); - final ConsumerConfiguration consumerConfiguration = new ConsumerConfiguration(mockedConsumerMetadata, + 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); + verify(mockedConsumerConnector, atLeast(1)).createMessageStreams(topicsStreamMap, null, null); + //verify(mockedConsumerConnector, atMost(0)).createMessageStreams(topicsStreamMap, null, null); } @Test + @SuppressWarnings("unchecked") public void testGetConsumerMapWithMessageStreamsWithDecoders() { @SuppressWarnings("unchecked") @@ -355,7 +357,7 @@ public class ConsumerConfigurationTests { final ConsumerConnectionProvider mockedConsumerConnectionProvider = mock(ConsumerConnectionProvider.class); - final MessageLeftOverTracker mockedMessageLeftOverTracker = mock(MessageLeftOverTracker.class); + final MessageLeftOverTracker mockedMessageLeftOverTracker = mock(MessageLeftOverTracker.class); final ConsumerConnector mockedConsumerConnector = mock(ConsumerConnector.class); when(mockedConsumerConnectionProvider.getConsumerConnector()).thenReturn(mockedConsumerConnector); @@ -363,7 +365,7 @@ public class ConsumerConfigurationTests { final Map>> messageStreams = new HashMap>>(); when(mockedConsumerConnector.createMessageStreams(topicsStreamMap)).thenReturn(messageStreams); - final ConsumerConfiguration consumerConfiguration = new ConsumerConfiguration(mockedConsumerMetadata, + final ConsumerConfiguration consumerConfiguration = new ConsumerConfiguration(mockedConsumerMetadata, mockedConsumerConnectionProvider, mockedMessageLeftOverTracker); consumerConfiguration.getConsumerMapWithMessageStreams(); 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 456bca8346..a5fffed09d 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 @@ -30,20 +30,22 @@ import java.util.Map; * @author Soby Chacko * @since 0.5 */ -public class KafkaConsumerContextTest { +public class KafkaConsumerContextTest { @Test + @SuppressWarnings("unchecked") public void testMergeResultsFromMultipleConsumerConfiguration() { - final KafkaConsumerContext kafkaConsumerContext = new KafkaConsumerContext(); + final KafkaConsumerContext kafkaConsumerContext = new KafkaConsumerContext(); final ListableBeanFactory beanFactory = Mockito.mock(ListableBeanFactory.class); - final ConsumerConfiguration consumerConfiguration1 = Mockito.mock(ConsumerConfiguration.class); - final ConsumerConfiguration consumerConfiguration2 = Mockito.mock(ConsumerConfiguration.class); + final ConsumerConfiguration consumerConfiguration1 = Mockito.mock(ConsumerConfiguration.class); + final ConsumerConfiguration consumerConfiguration2 = Mockito.mock(ConsumerConfiguration.class); - final Map map = new HashMap(); + final Map> map = new HashMap>(); map.put("config1", consumerConfiguration1); map.put("config2", consumerConfiguration2); - Mockito.when(beanFactory.getBeansOfType(ConsumerConfiguration.class)).thenReturn(map); + Mockito.when((Map>) (Object) beanFactory.getBeansOfType(ConsumerConfiguration.class)).thenReturn( + map); kafkaConsumerContext.setBeanFactory(beanFactory); final Map>> result1 = new 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 6feefe0a3f..367e464186 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 @@ -24,7 +24,7 @@ import org.junit.Test; import org.mockito.ArgumentCaptor; import org.mockito.Mockito; import org.springframework.integration.Message; -import org.springframework.integration.kafka.serializer.avro.AvroBackedKafkaEncoder; +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; @@ -39,7 +39,7 @@ import java.io.ObjectInputStream; * @author Soby Chacko * @since 0.5 */ -public class ProducerConfigurationTests { +public class ProducerConfigurationTests { @Test @SuppressWarnings("unchecked") public void testSendMessageWithNonDefaultKeyAndValueEncoders() throws Exception { @@ -61,10 +61,12 @@ public class ProducerConfigurationTests { Mockito.verify(producer, Mockito.times(1)).send(Mockito.any(KeyedMessage.class)); - final ArgumentCaptor argument = ArgumentCaptor.forClass(KeyedMessage.class); + final ArgumentCaptor> argument = + (ArgumentCaptor>) (Object) + ArgumentCaptor.forClass(KeyedMessage.class); Mockito.verify(producer).send(argument.capture()); - final KeyedMessage capturedKeyMessage = argument.getValue(); + final KeyedMessage capturedKeyMessage = argument.getValue(); Assert.assertEquals(capturedKeyMessage.key(), "key"); Assert.assertEquals(capturedKeyMessage.message(), "test message"); @@ -93,12 +95,14 @@ public class ProducerConfigurationTests { Mockito.verify(producer, Mockito.times(1)).send(Mockito.any(KeyedMessage.class)); - final ArgumentCaptor argument = ArgumentCaptor.forClass(KeyedMessage.class); + final ArgumentCaptor> argument = + (ArgumentCaptor>) (Object) + ArgumentCaptor.forClass(KeyedMessage.class); Mockito.verify(producer).send(argument.capture()); - final KeyedMessage capturedKeyMessage = argument.getValue(); + final KeyedMessage capturedKeyMessage = argument.getValue(); - final byte[] keyBytes = (byte[])capturedKeyMessage.key(); + final byte[] keyBytes = capturedKeyMessage.key(); final ByteArrayInputStream keyInputStream = new ByteArrayInputStream (keyBytes); final ObjectInputStream keyObjectInputStream = new ObjectInputStream (keyInputStream); @@ -109,7 +113,7 @@ public class ProducerConfigurationTests { Assert.assertEquals(tk.getKeyPart1(), "compositePart1"); Assert.assertEquals(tk.getKeyPart2(), "compositePart2"); - final byte[] messageBytes = (byte[])capturedKeyMessage.message(); + final byte[] messageBytes = capturedKeyMessage.message(); final ByteArrayInputStream messageInputStream = new ByteArrayInputStream (messageBytes); final ObjectInputStream messageObjectInputStream = new ObjectInputStream (messageInputStream); @@ -130,7 +134,7 @@ public class ProducerConfigurationTests { @SuppressWarnings("unchecked") public void testSendMessageWithDefaultKeyEncoderAndNonDefaultValueEncoderAndCorrespondingData() throws Exception { final ProducerMetadata producerMetadata = new ProducerMetadata("test"); - final AvroBackedKafkaEncoder encoder = new AvroBackedKafkaEncoder(TestPayload.class); + final AvroReflectDatumBackedKafkaEncoder encoder = new AvroReflectDatumBackedKafkaEncoder(TestPayload.class); producerMetadata.setValueEncoder(encoder); producerMetadata.setKeyEncoder(new DefaultEncoder(null)); producerMetadata.setValueClassType(TestPayload.class); @@ -147,12 +151,14 @@ public class ProducerConfigurationTests { Mockito.verify(producer, Mockito.times(1)).send(Mockito.any(KeyedMessage.class)); - final ArgumentCaptor argument = ArgumentCaptor.forClass(KeyedMessage.class); + final ArgumentCaptor> argument = + (ArgumentCaptor>) (Object) + ArgumentCaptor.forClass(KeyedMessage.class); Mockito.verify(producer).send(argument.capture()); - final KeyedMessage capturedKeyMessage = argument.getValue(); + final KeyedMessage capturedKeyMessage = argument.getValue(); - final byte[] keyBytes = (byte[])capturedKeyMessage.key(); + final byte[] keyBytes = capturedKeyMessage.key(); final ByteArrayInputStream keyInputStream = new ByteArrayInputStream (keyBytes); final ObjectInputStream keyObjectInputStream = new ObjectInputStream (keyInputStream); @@ -171,7 +177,7 @@ public class ProducerConfigurationTests { @SuppressWarnings("unchecked") public void testSendMessageWithNonDefaultKeyEncoderAndDefaultValueEncoderAndCorrespondingData() throws Exception { final ProducerMetadata producerMetadata = new ProducerMetadata("test"); - final AvroBackedKafkaEncoder encoder = new AvroBackedKafkaEncoder(TestKey.class); + final AvroReflectDatumBackedKafkaEncoder encoder = new AvroReflectDatumBackedKafkaEncoder(TestKey.class); producerMetadata.setKeyEncoder(encoder); producerMetadata.setValueEncoder(new DefaultEncoder(null)); producerMetadata.setKeyClassType(TestKey.class); @@ -187,14 +193,16 @@ public class ProducerConfigurationTests { Mockito.verify(producer, Mockito.times(1)).send(Mockito.any(KeyedMessage.class)); - final ArgumentCaptor argument = ArgumentCaptor.forClass(KeyedMessage.class); + final ArgumentCaptor> argument = + (ArgumentCaptor>) (Object) + ArgumentCaptor.forClass(KeyedMessage.class); Mockito.verify(producer).send(argument.capture()); - final KeyedMessage capturedKeyMessage = argument.getValue(); + final KeyedMessage capturedKeyMessage = argument.getValue(); Assert.assertEquals(capturedKeyMessage.key(), tk); - final byte[] payloadBytes = (byte[])capturedKeyMessage.message(); + final byte[] payloadBytes = capturedKeyMessage.message(); final ByteArrayInputStream payloadBis = new ByteArrayInputStream (payloadBytes); final ObjectInputStream payloadOis = new ObjectInputStream (payloadBis); @@ -226,11 +234,13 @@ public class ProducerConfigurationTests { Mockito.verify(producer, Mockito.times(1)).send(Mockito.any(KeyedMessage.class)); - final ArgumentCaptor argument = ArgumentCaptor.forClass(KeyedMessage.class); + final ArgumentCaptor> argument = + (ArgumentCaptor>) (Object) + ArgumentCaptor.forClass(KeyedMessage.class); Mockito.verify(producer).send(argument.capture()); - final KeyedMessage capturedKeyMessage = argument.getValue(); - final byte[] keyBytes = (byte[])capturedKeyMessage.key(); + final KeyedMessage capturedKeyMessage = argument.getValue(); + final byte[] keyBytes = capturedKeyMessage.key(); final ByteArrayInputStream keyBis = new ByteArrayInputStream (keyBytes); final ObjectInputStream keyOis = new ObjectInputStream (keyBis); @@ -238,7 +248,7 @@ public class ProducerConfigurationTests { Assert.assertEquals("key", keyObj); - final byte[] payloadBytes = (byte[])capturedKeyMessage.message(); + final byte[] payloadBytes = capturedKeyMessage.message(); final ByteArrayInputStream payloadBis = new ByteArrayInputStream (payloadBytes); final ObjectInputStream payloadOis = new ObjectInputStream (payloadBis); diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/support/ProducerFactoryBeanTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/support/ProducerFactoryBeanTests.java index 683a2699c0..e1dc411fbd 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/support/ProducerFactoryBeanTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/support/ProducerFactoryBeanTests.java @@ -24,14 +24,14 @@ import org.mockito.Mockito; * @author Soby Chacko * @since 0.5 */ -public class ProducerFactoryBeanTests { +public class ProducerFactoryBeanTests { @Test public void createProducerWithDefaultMetadata() throws Exception { final ProducerMetadata producerMetadata = new ProducerMetadata("test"); final ProducerMetadata tm = Mockito.spy(producerMetadata); final ProducerFactoryBean producerFactoryBean = new ProducerFactoryBean(tm, "localhost:9092"); - final Producer producer = producerFactoryBean.getObject(); + final Producer producer = producerFactoryBean.getObject(); Assert.assertTrue(producer != null); @@ -50,7 +50,7 @@ public class ProducerFactoryBeanTests { producerMetadata.setBatchNumMessages("300"); final ProducerMetadata tm = Mockito.spy(producerMetadata); final ProducerFactoryBean producerFactoryBean = new ProducerFactoryBean(tm, "localhost:9092"); - final Producer producer = producerFactoryBean.getObject(); + final Producer producer = producerFactoryBean.getObject(); Assert.assertTrue(producer != null); diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/test/utils/User.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/test/utils/User.java new file mode 100644 index 0000000000..c6ad7ddaff --- /dev/null +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/test/utils/User.java @@ -0,0 +1,96 @@ +package org.springframework.integration.kafka.test.utils; + +import org.apache.avro.specific.SpecificRecord; + +/** + * @author Soby Chacko + * @since 0.5 + *

+ * This class is copied (partly) from an Avro generated class for necessary testing. + * Please use caution when modify. + */ +public class User extends org.apache.avro.specific.SpecificRecordBase implements SpecificRecord { + + public static final org.apache.avro.Schema SCHEMA$ = new org.apache.avro.Schema.Parser().parse("{\"type\":\"record\",\"name\":\"User\",\"namespace\":\"org.springframework.integration.samples.kafka.user\",\"fields\":[{\"name\":\"firstName\",\"type\":\"string\"},{\"name\":\"lastName\",\"type\":\"string\"}]}"); + public java.lang.CharSequence firstName; + public java.lang.CharSequence lastName; + + /** + * Default constructor. + */ + public User() { + } + + /** + * All-args constructor. + */ + public User(java.lang.CharSequence firstName, java.lang.CharSequence lastName) { + this.firstName = firstName; + this.lastName = lastName; + } + + public org.apache.avro.Schema getSchema() { + return SCHEMA$; + } + + // Used by DatumWriter. Applications should not call. + public java.lang.Object get(int field$) { + switch (field$) { + case 0: + return firstName; + case 1: + return lastName; + default: + throw new org.apache.avro.AvroRuntimeException("Bad index"); + } + } + + // Used by DatumReader. Applications should not call. + @SuppressWarnings(value = "unchecked") + public void put(int field$, java.lang.Object value$) { + switch (field$) { + case 0: + firstName = (java.lang.CharSequence) value$; + break; + case 1: + lastName = (java.lang.CharSequence) value$; + break; + default: + throw new org.apache.avro.AvroRuntimeException("Bad index"); + } + } + + /** + * Gets the value of the 'firstName' field. + */ + public java.lang.CharSequence getFirstName() { + return firstName; + } + + /** + * Sets the value of the 'firstName' field. + * + * @param value the value to set. + */ + public void setFirstName(java.lang.CharSequence value) { + this.firstName = value; + } + + /** + * Gets the value of the 'lastName' field. + */ + public java.lang.CharSequence getLastName() { + return lastName; + } + + /** + * Sets the value of the 'lastName' field. + * + * @param value the value to set. + */ + public void setLastName(java.lang.CharSequence value) { + this.lastName = value; + } +} + +