diff --git a/build.gradle b/build.gradle index 3ef04359..1a2adff8 100644 --- a/build.gradle +++ b/build.gradle @@ -70,16 +70,14 @@ subprojects { subproject -> ext { assertjVersion = '3.3.0' - avroVersion = '1.7.6' gsCollectionsVersion = '5.0.0' hamcrestVersion = '1.3' + jacksonVersion = '2.3.2' junitVersion = '4.12' kafkaVersion = '0.9.0.1' log4jVersion = '1.2.17' mockitoVersion = '1.9.5' - // metricsVersion = '2.2.0' scalaVersion = '2.11' - reactor2Version = '2.0.6.RELEASE' springRetryVersion = '1.1.2.RELEASE' springVersion = '4.2.5.RELEASE' @@ -143,12 +141,9 @@ project ('spring-kafka') { dependencies { compile "org.springframework:spring-messaging:$springVersion" -// compile ("org.apache.avro:avro:$avroVersion", optional) -// compile ("org.apache.avro:avro-compiler:$avroVersion", optional) -// compile "com.goldmansachs:gs-collections:$gsCollectionsVersion" -// compile "io.projectreactor:reactor-core:$reactor2Version" - compile "org.apache.kafka:kafka-clients:$kafkaVersion" + compile ("com.fasterxml.jackson.core:jackson-core:$jacksonVersion", optional) + compile ("com.fasterxml.jackson.core:jackson-databind:$jacksonVersion", optional) testCompile project (":spring-kafka-test") testCompile "org.assertj:assertj-core:$assertjVersion" diff --git a/spring-kafka-test/src/main/java/org/springframework/kafka/test/utils/KafkaTestUtils.java b/spring-kafka-test/src/main/java/org/springframework/kafka/test/utils/KafkaTestUtils.java index 01cfccef..e8aa1cff 100644 --- a/spring-kafka-test/src/main/java/org/springframework/kafka/test/utils/KafkaTestUtils.java +++ b/spring-kafka-test/src/main/java/org/springframework/kafka/test/utils/KafkaTestUtils.java @@ -105,6 +105,8 @@ public final class KafkaTestUtils { * Poll the consumer, expecting a single record for the specified topic. * @param consumer the consumer. * @param topic the topic. + * @param the key type. + * @param the value type. * @return the record. * @throws org.junit.ComparisonFailure if exactly one record is not received. */ @@ -117,6 +119,8 @@ public final class KafkaTestUtils { /** * Poll the consumer for records. * @param consumer the consumer. + * @param the key type. + * @param the value type. * @return the records. */ public static ConsumerRecords getRecords(Consumer consumer) { diff --git a/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaException.java b/spring-kafka/src/main/java/org/springframework/kafka/KafkaException.java similarity index 96% rename from spring-kafka/src/main/java/org/springframework/kafka/core/KafkaException.java rename to spring-kafka/src/main/java/org/springframework/kafka/KafkaException.java index cf55a6e3..be919be7 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaException.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/KafkaException.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.kafka.core; +package org.springframework.kafka; import org.springframework.core.NestedRuntimeException; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/config/AbstractKafkaListenerContainerFactory.java b/spring-kafka/src/main/java/org/springframework/kafka/config/AbstractKafkaListenerContainerFactory.java index 6c8c3849..fda96fe2 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/config/AbstractKafkaListenerContainerFactory.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/config/AbstractKafkaListenerContainerFactory.java @@ -23,6 +23,7 @@ import org.springframework.kafka.core.ConsumerFactory; import org.springframework.kafka.listener.AbstractMessageListenerContainer; import org.springframework.kafka.listener.AbstractMessageListenerContainer.AckMode; import org.springframework.kafka.listener.ErrorHandler; +import org.springframework.kafka.support.converter.MessageConverter; /** * Base {@link KafkaListenerContainerFactory} for Spring's base container implementation. @@ -54,6 +55,8 @@ public abstract class AbstractKafkaListenerContainerFactory } @Override - public void setupListenerContainer(MessageListenerContainer listenerContainer) { - setupMessageListener(listenerContainer); + public void setupListenerContainer(MessageListenerContainer listenerContainer, MessageConverter messageConverter) { + setupMessageListener(listenerContainer, messageConverter); } /** * Create a {@link MessageListener} that is able to serve this endpoint for the * specified container. * @param container the {@link MessageListenerContainer} to create a {@link MessageListener}. + * @param messageConverter the message converter - may be null. * @return a a {@link MessageListener} instance. */ - protected abstract MessageListener createMessageListener(MessageListenerContainer container); + protected abstract MessageListener createMessageListener(MessageListenerContainer container, + MessageConverter messageConverter); - private void setupMessageListener(MessageListenerContainer container) { - MessageListener messageListener = createMessageListener(container); + private void setupMessageListener(MessageListenerContainer container, MessageConverter messageConverter) { + MessageListener messageListener = createMessageListener(container, messageConverter); Assert.state(messageListener != null, "Endpoint [" + this + "] must provide a non null message listener"); container.setupMessageListener(messageListener); } diff --git a/spring-kafka/src/main/java/org/springframework/kafka/config/KafkaListenerEndpoint.java b/spring-kafka/src/main/java/org/springframework/kafka/config/KafkaListenerEndpoint.java index 40f6aca0..d30b0f61 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/config/KafkaListenerEndpoint.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/config/KafkaListenerEndpoint.java @@ -22,6 +22,7 @@ import java.util.regex.Pattern; import org.apache.kafka.common.TopicPartition; import org.springframework.kafka.listener.MessageListenerContainer; +import org.springframework.kafka.support.converter.MessageConverter; /** * Model for a Kafka listener endpoint. Can be used against a @@ -75,7 +76,8 @@ public interface KafkaListenerEndpoint { * use but an implementation may override any default setting that * was already set. * @param listenerContainer the listener container to configure + * @param messageConverter the message converter - can be null */ - void setupListenerContainer(MessageListenerContainer listenerContainer); + void setupListenerContainer(MessageListenerContainer listenerContainer, MessageConverter messageConverter); } diff --git a/spring-kafka/src/main/java/org/springframework/kafka/config/MethodKafkaListenerEndpoint.java b/spring-kafka/src/main/java/org/springframework/kafka/config/MethodKafkaListenerEndpoint.java index febbd347..1decc744 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/config/MethodKafkaListenerEndpoint.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/config/MethodKafkaListenerEndpoint.java @@ -21,6 +21,7 @@ import java.lang.reflect.Method; import org.springframework.kafka.listener.MessageListenerContainer; import org.springframework.kafka.listener.adapter.HandlerAdapter; import org.springframework.kafka.listener.adapter.MessagingMessageListenerAdapter; +import org.springframework.kafka.support.converter.MessageConverter; import org.springframework.messaging.handler.annotation.support.MessageHandlerMethodFactory; import org.springframework.messaging.handler.invocation.InvocableHandlerMethod; import org.springframework.util.Assert; @@ -88,11 +89,15 @@ public class MethodKafkaListenerEndpoint extends AbstractKafkaListenerEndp } @Override - protected MessagingMessageListenerAdapter createMessageListener(MessageListenerContainer container) { + protected MessagingMessageListenerAdapter createMessageListener(MessageListenerContainer container, + MessageConverter messageConverter) { Assert.state(this.messageHandlerMethodFactory != null, "Could not create message listener - MessageHandlerMethodFactory not set"); MessagingMessageListenerAdapter messageListener = createMessageListenerInstance(); messageListener.setHandlerMethod(configureListenerAdapter(messageListener)); + if (messageConverter != null) { + messageListener.setMessageConverter(messageConverter); + } return messageListener; } @@ -112,7 +117,7 @@ public class MethodKafkaListenerEndpoint extends AbstractKafkaListenerEndp * @return the {@link MessagingMessageListenerAdapter} instance. */ protected MessagingMessageListenerAdapter createMessageListenerInstance() { - return new MessagingMessageListenerAdapter(); + return new MessagingMessageListenerAdapter(this.method); } @Override diff --git a/spring-kafka/src/main/java/org/springframework/kafka/config/SimpleKafkaListenerEndpoint.java b/spring-kafka/src/main/java/org/springframework/kafka/config/SimpleKafkaListenerEndpoint.java index 1149a3e9..4933bbc8 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/config/SimpleKafkaListenerEndpoint.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/config/SimpleKafkaListenerEndpoint.java @@ -19,6 +19,7 @@ package org.springframework.kafka.config; import org.springframework.kafka.listener.MessageListener; import org.springframework.kafka.listener.MessageListenerContainer; +import org.springframework.kafka.support.converter.MessageConverter; /** * A {@link KafkaListenerEndpoint} simply providing the {@link MessageListener} to @@ -56,7 +57,8 @@ public class SimpleKafkaListenerEndpoint extends AbstractKafkaListenerEndp @Override - protected MessageListener createMessageListener(MessageListenerContainer container) { + protected MessageListener createMessageListener(MessageListenerContainer container, + MessageConverter messageConverter) { return getMessageListener(); } diff --git a/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaOperations.java b/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaOperations.java index 766d67ab..935a24d8 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaOperations.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaOperations.java @@ -97,14 +97,15 @@ public interface KafkaOperations { Future send(String topic, int partition, K key, V data); /** - * Send a message with routing information in message headers. + * Send a message with routing information in message headers. The message payload + * may be converted before sending. * @param message the message to send. * @return a Future for the {@link RecordMetadata}. * @see org.springframework.kafka.support.KafkaHeaders#TOPIC * @see org.springframework.kafka.support.KafkaHeaders#PARTITION_ID * @see org.springframework.kafka.support.KafkaHeaders#MESSAGE_KEY */ - Future send(Message message); + Future convertAndSend(Message message); // Sync methods @@ -192,7 +193,8 @@ public interface KafkaOperations { throws InterruptedException, ExecutionException; /** - * Send a message with routing information in message headers. + * Send a message with routing information in message headers. The message payload + * may be converted before sending. * @param message the message to send. * @return a Future for the {@link RecordMetadata}. * @throws ExecutionException execution exception while awaiting result. @@ -201,7 +203,7 @@ public interface KafkaOperations { * @see org.springframework.kafka.support.KafkaHeaders#PARTITION_ID * @see org.springframework.kafka.support.KafkaHeaders#MESSAGE_KEY */ - RecordMetadata syncSend(Message message) + RecordMetadata syncConvertAndSend(Message message) throws InterruptedException, ExecutionException; /** diff --git a/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaTemplate.java b/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaTemplate.java index ba1b11ae..61a86183 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaTemplate.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaTemplate.java @@ -25,12 +25,12 @@ import org.apache.kafka.clients.producer.Producer; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.clients.producer.RecordMetadata; -import org.springframework.kafka.support.KafkaHeaders; import org.springframework.kafka.support.LoggingProducerListener; import org.springframework.kafka.support.ProducerListener; import org.springframework.kafka.support.ProducerListenerInvokingCallback; +import org.springframework.kafka.support.converter.MessageConverter; +import org.springframework.kafka.support.converter.MessagingMessageConverter; import org.springframework.messaging.Message; -import org.springframework.messaging.MessageHeaders; /** @@ -48,6 +48,8 @@ public class KafkaTemplate implements KafkaOperations { private final ProducerFactory producerFactory; + private MessageConverter messageConverter = new MessagingMessageConverter(); + private volatile Producer producer; private volatile String defaultTopic; @@ -90,6 +92,22 @@ public class KafkaTemplate implements KafkaOperations { this.producerListener = producerListener; } + /** + * Return the message converter. + * @return the message converter. + */ + public MessageConverter getMessageConverter() { + return this.messageConverter; + } + + /** + * Set the message converter to use. + * @param messageConverter the message converter. + */ + public void setMessageConverter(MessageConverter messageConverter) { + this.messageConverter = messageConverter; + } + @Override public Future send(V data) { return send(this.defaultTopic, data); @@ -129,10 +147,11 @@ public class KafkaTemplate implements KafkaOperations { return doSend(producerRecord); } + @SuppressWarnings("unchecked") @Override - public Future send(Message message) { - ProducerRecord producerRecord = messageToProducerRecord(message); - return doSend(producerRecord); + public Future convertAndSend(Message message) { + ProducerRecord producerRecord = this.messageConverter.fromMessage(message, this.defaultTopic); + return doSend((ProducerRecord) producerRecord); } @Override @@ -189,9 +208,9 @@ public class KafkaTemplate implements KafkaOperations { } @Override - public RecordMetadata syncSend(Message message) + public RecordMetadata syncConvertAndSend(Message message) throws InterruptedException, ExecutionException { - Future future = send(message); + Future future = convertAndSend(message); flush(); return future.get(); } @@ -232,14 +251,4 @@ public class KafkaTemplate implements KafkaOperations { return future; } - @SuppressWarnings({ "rawtypes", "unchecked" }) - private ProducerRecord messageToProducerRecord(Message message) { - MessageHeaders headers = message.getHeaders(); - String topic = headers.get(KafkaHeaders.TOPIC, String.class); - Integer partition = headers.get(KafkaHeaders.PARTITION_ID, Integer.class); - Object key = headers.get(KafkaHeaders.MESSAGE_KEY); - Object payload = message.getPayload(); - return new ProducerRecord(topic == null ? this.defaultTopic : topic, partition, key, payload); - } - } diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/ListenerExecutionFailedException.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/ListenerExecutionFailedException.java index 8a0ef0c9..c5da0581 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/ListenerExecutionFailedException.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/ListenerExecutionFailedException.java @@ -16,7 +16,7 @@ package org.springframework.kafka.listener; -import org.springframework.kafka.core.KafkaException; +import org.springframework.kafka.KafkaException; /** * The listener specific {@link KafkaException} extension. diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/DelegatingInvocableHandler.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/DelegatingInvocableHandler.java index c5f62bf4..dbe56255 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/DelegatingInvocableHandler.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/DelegatingInvocableHandler.java @@ -24,7 +24,7 @@ import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import org.springframework.core.MethodParameter; -import org.springframework.kafka.core.KafkaException; +import org.springframework.kafka.KafkaException; import org.springframework.messaging.Message; import org.springframework.messaging.handler.annotation.Payload; import org.springframework.messaging.handler.invocation.InvocableHandlerMethod; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/MessagingMessageListenerAdapter.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/MessagingMessageListenerAdapter.java index 7d6fbf0a..1c537436 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/MessagingMessageListenerAdapter.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/MessagingMessageListenerAdapter.java @@ -16,8 +16,14 @@ package org.springframework.kafka.listener.adapter; +import java.lang.reflect.Method; +import java.lang.reflect.ParameterizedType; +import java.lang.reflect.Type; +import java.lang.reflect.WildcardType; + import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.springframework.core.MethodParameter; import org.springframework.kafka.listener.ListenerExecutionFailedException; import org.springframework.kafka.support.Acknowledgment; import org.springframework.kafka.support.converter.MessageConverter; @@ -25,6 +31,8 @@ import org.springframework.kafka.support.converter.MessagingMessageConverter; import org.springframework.messaging.Message; import org.springframework.messaging.MessagingException; import org.springframework.messaging.converter.MessageConversionException; +import org.springframework.messaging.handler.annotation.Payload; + /** * A {@link org.springframework.kafka.listener.MessageListener MessageListener} @@ -45,9 +53,16 @@ import org.springframework.messaging.converter.MessageConversionException; */ public class MessagingMessageListenerAdapter extends AbstractAdaptableMessageListener { + private final Type inferredType; + private HandlerAdapter handlerMethod; - private MessageConverter messageConverter = new MessagingMessageConverter<>(); + private MessageConverter messageConverter = new MessagingMessageConverter(); + + + public MessagingMessageListenerAdapter(Method method) { + this.inferredType = determineInferredType(method); + } /** * Set the {@link HandlerAdapter} to use to invoke the method @@ -62,7 +77,7 @@ public class MessagingMessageListenerAdapter extends AbstractAdaptableMess * Set the MessageConverter. * @param messageConverter the converter. */ - public void setMessageConverter(MessageConverter messageConverter) { + public void setMessageConverter(MessageConverter messageConverter) { this.messageConverter = messageConverter; } @@ -72,7 +87,7 @@ public class MessagingMessageListenerAdapter extends AbstractAdaptableMess * @return the {@link MessagingMessageConverter} for this listener, * being able to convert {@link org.springframework.messaging.Message}. */ - protected final MessageConverter getMessageConverter() { + protected final MessageConverter getMessageConverter() { return this.messageConverter; } @@ -87,7 +102,7 @@ public class MessagingMessageListenerAdapter extends AbstractAdaptableMess } protected Message toMessagingMessage(ConsumerRecord record, Acknowledgment acknowledgment) { - return getMessageConverter().toMessage(record, acknowledgment); + return getMessageConverter().toMessage(record, acknowledgment, this.inferredType); } /** @@ -124,4 +139,61 @@ public class MessagingMessageListenerAdapter extends AbstractAdaptableMess + "Bean [" + this.handlerMethod.getBean() + "]"; } + private Type determineInferredType(Method method) { + if (method == null) { + return null; + } + + Type genericParameterType = null; + + for (int i = 0; i < method.getParameterTypes().length; i++) { + MethodParameter methodParameter = new MethodParameter(method, i); + /* + * We're looking for a single non-annotated parameter, or one annotated with @Payload. + * We ignore parameters with type Message because they are not involved with conversion. + */ + if (eligibleParameter(methodParameter) + && (methodParameter.getParameterAnnotations().length == 0 + || methodParameter.hasParameterAnnotation(Payload.class))) { + if (genericParameterType == null) { + genericParameterType = methodParameter.getGenericParameterType(); + if (genericParameterType instanceof ParameterizedType) { + ParameterizedType parameterizedType = (ParameterizedType) genericParameterType; + if (parameterizedType.getRawType().equals(Message.class)) { + genericParameterType = ((ParameterizedType) genericParameterType) + .getActualTypeArguments()[0]; + } + } + } + else { + if (this.logger.isDebugEnabled()) { + this.logger.debug("Ambiguous parameters for target payload for method " + method + + "; no inferred type available"); + } + return null; + } + } + } + + return genericParameterType; + } + + /* + * Don't consider parameter types that are available after conversion. + * Acknowledgment, ConsumerRecord and Message. + */ + private boolean eligibleParameter(MethodParameter methodParameter) { + Type parameterType = methodParameter.getGenericParameterType(); + if (parameterType.equals(Acknowledgment.class) || parameterType.equals(ConsumerRecord.class)) { + return false; + } + if (parameterType instanceof ParameterizedType) { + ParameterizedType parameterizedType = (ParameterizedType) parameterType; + if (parameterizedType.getRawType().equals(Message.class)) { + return !(parameterizedType.getActualTypeArguments()[0] instanceof WildcardType); + } + } + return !parameterType.equals(Message.class); // could be Message without a generic type + } + } diff --git a/spring-kafka/src/main/java/org/springframework/kafka/support/converter/ConversionException.java b/spring-kafka/src/main/java/org/springframework/kafka/support/converter/ConversionException.java new file mode 100644 index 00000000..dcf67aff --- /dev/null +++ b/spring-kafka/src/main/java/org/springframework/kafka/support/converter/ConversionException.java @@ -0,0 +1,34 @@ +/* + * Copyright 2016 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.kafka.support.converter; + +import org.springframework.kafka.KafkaException; + +/** + * Exception for conversions. + * + * @author Gary Russell + * + */ +@SuppressWarnings("serial") +public class ConversionException extends KafkaException { + + public ConversionException(String message, Throwable cause) { + super(message, cause); + } + +} diff --git a/spring-kafka/src/main/java/org/springframework/kafka/support/converter/MessageConverter.java b/spring-kafka/src/main/java/org/springframework/kafka/support/converter/MessageConverter.java index 4138d6a5..a4f21bf7 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/support/converter/MessageConverter.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/support/converter/MessageConverter.java @@ -16,21 +16,36 @@ package org.springframework.kafka.support.converter; +import java.lang.reflect.Type; + import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.apache.kafka.clients.producer.ProducerRecord; import org.springframework.kafka.support.Acknowledgment; import org.springframework.messaging.Message; /** - * The Kafka specific {@link Message} converter strategy. - * - * @param the key type. - * @param the value type. + * A Kafka-specific {@link Message} converter strategy. * * @author Gary Russell */ -public interface MessageConverter { +public interface MessageConverter { - Message toMessage(ConsumerRecord record, Acknowledgment acknowledgment); + /** + * Convert a {@link ConsumerRecord} to a {@link Message}. + * @param record the record. + * @param acknowledgment the acknowledgment. + * @param payloadType the required payload type. + * @return the message. + */ + Message toMessage(ConsumerRecord record, Acknowledgment acknowledgment, Type payloadType); + + /** + * Convert a message to a producer record. + * @param message the message. + * @param defaultTopic the default topic to use if no header found. + * @return the producer record. + */ + ProducerRecord fromMessage(Message message, String defaultTopic); } diff --git a/spring-kafka/src/main/java/org/springframework/kafka/support/converter/MessagingMessageConverter.java b/spring-kafka/src/main/java/org/springframework/kafka/support/converter/MessagingMessageConverter.java index 3ae649c8..35e7d0f0 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/support/converter/MessagingMessageConverter.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/support/converter/MessagingMessageConverter.java @@ -16,9 +16,11 @@ package org.springframework.kafka.support.converter; +import java.lang.reflect.Type; import java.util.Map; import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.apache.kafka.clients.producer.ProducerRecord; import org.springframework.kafka.support.Acknowledgment; import org.springframework.kafka.support.KafkaHeaders; @@ -31,13 +33,10 @@ import org.springframework.messaging.support.MessageBuilder; *

* Populates {@link KafkaHeaders} based on the {@link ConsumerRecord} onto the returned message. * - * @param the key type. - * @param the value type. - * * @author Marius Bogoevici * @author Gary Russell */ -public class MessagingMessageConverter implements MessageConverter { +public class MessagingMessageConverter implements MessageConverter { private boolean generateMessageId = false; @@ -62,8 +61,9 @@ public class MessagingMessageConverter implements MessageConverter { } @Override - public Message toMessage(ConsumerRecord record, Acknowledgment acknowledgment) { - KafkaMessageHeaders kafkaMessageHeaders = new KafkaMessageHeaders(this.generateMessageId, this.generateTimestamp); + public Message toMessage(ConsumerRecord record, Acknowledgment acknowledgment, Type type) { + KafkaMessageHeaders kafkaMessageHeaders = new KafkaMessageHeaders(this.generateMessageId, + this.generateTimestamp); Map rawHeaders = kafkaMessageHeaders.getRawHeaders(); rawHeaders.put(KafkaHeaders.RECEIVED_MESSAGE_KEY, record.key()); @@ -75,15 +75,36 @@ public class MessagingMessageConverter implements MessageConverter { rawHeaders.put(KafkaHeaders.ACKNOWLEDGMENT, acknowledgment); } - return MessageBuilder.createMessage(extractAndConvertValue(record), kafkaMessageHeaders); + return MessageBuilder.createMessage(extractAndConvertValue(record, type), kafkaMessageHeaders); + } + + @SuppressWarnings({ "unchecked", "rawtypes" }) + @Override + public ProducerRecord fromMessage(Message message, String defaultTopic) { + MessageHeaders headers = message.getHeaders(); + String topic = headers.get(KafkaHeaders.TOPIC, String.class); + Integer partition = headers.get(KafkaHeaders.PARTITION_ID, Integer.class); + Object key = headers.get(KafkaHeaders.MESSAGE_KEY); + Object payload = convertPayload(message); + return new ProducerRecord(topic == null ? defaultTopic : topic, partition, key, payload); + } + + /** + * Subclasses can convert the payload; by default, it's sent unchanged to Kafka. + * @param message the message. + * @return the payload. + */ + protected Object convertPayload(Message message) { + return message.getPayload(); } /** * Subclasses can convert the value; by default, it's returned as provided by Kafka. * @param record the record. + * @param type the required type. * @return the value. */ - protected V extractAndConvertValue(ConsumerRecord record) { + protected Object extractAndConvertValue(ConsumerRecord record, Type type) { return record.value(); } diff --git a/spring-kafka/src/main/java/org/springframework/kafka/support/converter/StringJsonMessageConverter.java b/spring-kafka/src/main/java/org/springframework/kafka/support/converter/StringJsonMessageConverter.java new file mode 100644 index 00000000..cbc63605 --- /dev/null +++ b/spring-kafka/src/main/java/org/springframework/kafka/support/converter/StringJsonMessageConverter.java @@ -0,0 +1,78 @@ +/* + * Copyright 2016 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.kafka.support.converter; + +import java.io.IOException; +import java.lang.reflect.Type; + +import org.apache.kafka.clients.consumer.ConsumerRecord; + +import org.springframework.messaging.Message; + +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.databind.JavaType; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.fasterxml.jackson.databind.type.TypeFactory; + +/** + * JSON Message converter - String on output, String or byte[] on input. + * + * @author Gary Russell + * + */ +public class StringJsonMessageConverter extends MessagingMessageConverter { + + private final ObjectMapper objectMapper = new ObjectMapper(); + + + @Override + protected Object convertPayload(Message message) { + try { + return this.objectMapper.writeValueAsString(message.getPayload()); + } + catch (JsonProcessingException e) { + throw new ConversionException("Failed to convert to JSON", e); + } + } + + + @Override + protected Object extractAndConvertValue(ConsumerRecord record, Type type) { + JavaType javaType = TypeFactory.defaultInstance().constructType(type); + Object value = record.value(); + if (value instanceof String) { + try { + return this.objectMapper.readValue((String) value, javaType); + } + catch (IOException e) { + throw new ConversionException("Failed to convert from JSON", e); + } + } + else if (value instanceof byte[]) { + try { + return this.objectMapper.readValue((byte[]) value, javaType); + } + catch (IOException e) { + throw new ConversionException("Failed to convert from JSON", e); + } + } + else { + throw new IllegalStateException("Only String or byte[] supported"); + } + } + +} diff --git a/spring-kafka/src/test/java/org/springframework/kafka/annotation/EnableKafkaIntegrationTests.java b/spring-kafka/src/test/java/org/springframework/kafka/annotation/EnableKafkaIntegrationTests.java index 3f6acc6f..641c9ad2 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/annotation/EnableKafkaIntegrationTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/annotation/EnableKafkaIntegrationTests.java @@ -44,10 +44,12 @@ import org.springframework.kafka.listener.AbstractMessageListenerContainer.AckMo import org.springframework.kafka.listener.ConcurrentMessageListenerContainer; import org.springframework.kafka.support.Acknowledgment; import org.springframework.kafka.support.KafkaHeaders; +import org.springframework.kafka.support.converter.StringJsonMessageConverter; import org.springframework.kafka.test.rule.KafkaEmbedded; import org.springframework.kafka.test.utils.KafkaTestUtils; import org.springframework.messaging.handler.annotation.Header; import org.springframework.messaging.handler.annotation.Payload; +import org.springframework.messaging.support.MessageBuilder; import org.springframework.test.annotation.DirtiesContext; import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; @@ -70,7 +72,7 @@ public class EnableKafkaIntegrationTests { @ClassRule public static KafkaEmbedded embeddedKafka = new KafkaEmbedded(1, true, "annotated1", "annotated2", "annotated3", - "annotated4", "annotated5", "annotated6", "annotated7", "annotated8", "annotated9"); + "annotated4", "annotated5", "annotated6", "annotated7", "annotated8", "annotated9", "annotated10"); @Autowired public IfaceListenerImpl ifaceListener; @@ -81,6 +83,9 @@ public class EnableKafkaIntegrationTests { @Autowired public KafkaTemplate template; + @Autowired + public KafkaTemplate kafkaJsonTemplate; + @Autowired public KafkaListenerEndpointRegistry registry; @@ -137,6 +142,19 @@ public class EnableKafkaIntegrationTests { assertThat(this.ifaceListener.getLatch2().await(20, TimeUnit.SECONDS)).isTrue(); } + @Test + public void testJson() throws Exception { + Foo foo = new Foo(); + foo.setBar("bar"); + kafkaJsonTemplate.convertAndSend(MessageBuilder.withPayload(foo) + .setHeader(KafkaHeaders.TOPIC, "annotated10") + .setHeader(KafkaHeaders.PARTITION_ID, 0) + .setHeader(KafkaHeaders.MESSAGE_KEY, 2) + .build()); + assertThat(this.listener.latch6.await(20, TimeUnit.SECONDS)).isTrue(); + assertThat(this.listener.foo.getBar()).isEqualTo("bar"); + } + @Configuration @EnableKafka @EnableTransactionManagement(proxyTargetClass = true) @@ -149,12 +167,21 @@ public class EnableKafkaIntegrationTests { @Bean public KafkaListenerContainerFactory> - kafkaListenerContainerFactory() { + kafkaListenerContainerFactory() { SimpleKafkaListenerContainerFactory factory = new SimpleKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory()); return factory; } + @Bean + public KafkaListenerContainerFactory> + kafkaJsonListenerContainerFactory() { + SimpleKafkaListenerContainerFactory factory = new SimpleKafkaListenerContainerFactory<>(); + factory.setConsumerFactory(consumerFactory()); + factory.setMessageConverter(new StringJsonMessageConverter()); + return factory; + } + @Bean public KafkaListenerContainerFactory> kafkaManualAckListenerContainerFactory() { @@ -209,10 +236,17 @@ public class EnableKafkaIntegrationTests { } @Bean - public KafkaTemplate kafkaTemplate() { + public KafkaTemplate template() { return new KafkaTemplate(producerFactory()); } + @Bean + public KafkaTemplate kafkaJsonTemplate() { + KafkaTemplate kafkaTemplate = new KafkaTemplate(producerFactory()); + kafkaTemplate.setMessageConverter(new StringJsonMessageConverter()); + return kafkaTemplate; + } + } static class Listener { @@ -227,6 +261,8 @@ public class EnableKafkaIntegrationTests { private final CountDownLatch latch5 = new CountDownLatch(1); + private final CountDownLatch latch6 = new CountDownLatch(1); + private volatile Integer partition; private volatile ConsumerRecord record; @@ -237,6 +273,8 @@ public class EnableKafkaIntegrationTests { private String topic; + private Foo foo; + @KafkaListener(id = "foo", topics = "annotated1") public void listen1(String foo) { this.latch1.countDown(); @@ -275,6 +313,12 @@ public class EnableKafkaIntegrationTests { this.latch5.countDown(); } + @KafkaListener(id = "buz", topics = "annotated10", containerFactory = "kafkaJsonListenerContainerFactory") + public void listen6(Foo foo) { + this.foo = foo; + this.latch6.countDown(); + } + } interface IfaceListener { @@ -325,5 +369,18 @@ public class EnableKafkaIntegrationTests { } + public static class Foo { + + private String bar; + + public String getBar() { + return this.bar; + } + + public void setBar(String bar) { + this.bar = bar; + } + + } } diff --git a/spring-kafka/src/test/java/org/springframework/kafka/core/KafkaTemplateTests.java b/spring-kafka/src/test/java/org/springframework/kafka/core/KafkaTemplateTests.java index 33ea33bd..46eac982 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/core/KafkaTemplateTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/core/KafkaTemplateTests.java @@ -79,7 +79,7 @@ public class KafkaTemplateTests { assertThat(received).has(key((Integer) null)); assertThat(received).has(partition(0)); assertThat(received).has(value("qux")); - template.syncSend(MessageBuilder.withPayload("fiz") + template.syncConvertAndSend(MessageBuilder.withPayload("fiz") .setHeader(KafkaHeaders.TOPIC, TEMPLATE_TOPIC) .setHeader(KafkaHeaders.PARTITION_ID, 0) .setHeader(KafkaHeaders.MESSAGE_KEY, 2) @@ -88,7 +88,7 @@ public class KafkaTemplateTests { assertThat(received).has(key(2)); assertThat(received).has(partition(0)); assertThat(received).has(value("fiz")); - template.syncSend(MessageBuilder.withPayload("buz") + template.syncConvertAndSend(MessageBuilder.withPayload("buz") .setHeader(KafkaHeaders.PARTITION_ID, 0) .setHeader(KafkaHeaders.MESSAGE_KEY, 2) .build());