diff --git a/gradle/wrapper/gradle-wrapper.jar b/gradle/wrapper/gradle-wrapper.jar index 5ccda13e..2c6137b8 100644 Binary files a/gradle/wrapper/gradle-wrapper.jar and b/gradle/wrapper/gradle-wrapper.jar differ diff --git a/gradle/wrapper/gradle-wrapper.properties b/gradle/wrapper/gradle-wrapper.properties index ab8d9dbe..cd2de496 100644 --- a/gradle/wrapper/gradle-wrapper.properties +++ b/gradle/wrapper/gradle-wrapper.properties @@ -1,6 +1,6 @@ -#Mon Mar 07 20:47:12 EST 2016 +#Wed Mar 30 18:31:11 EDT 2016 distributionBase=GRADLE_USER_HOME distributionPath=wrapper/dists zipStoreBase=GRADLE_USER_HOME zipStorePath=wrapper/dists -distributionUrl=https\://services.gradle.org/distributions/gradle-2.11-bin.zip +distributionUrl=https\://services.gradle.org/distributions/gradle-2.12-bin.zip diff --git a/gradlew b/gradlew index 91a7e269..9d82f789 100755 --- a/gradlew +++ b/gradlew @@ -42,11 +42,6 @@ case "`uname`" in ;; esac -# For Cygwin, ensure paths are in UNIX format before anything is touched. -if $cygwin ; then - [ -n "$JAVA_HOME" ] && JAVA_HOME=`cygpath --unix "$JAVA_HOME"` -fi - # Attempt to set APP_HOME # Resolve links: $0 may be a link PRG="$0" @@ -61,9 +56,9 @@ while [ -h "$PRG" ] ; do fi done SAVED="`pwd`" -cd "`dirname \"$PRG\"`/" >&- +cd "`dirname \"$PRG\"`/" >/dev/null APP_HOME="`pwd -P`" -cd "$SAVED" >&- +cd "$SAVED" >/dev/null CLASSPATH=$APP_HOME/gradle/wrapper/gradle-wrapper.jar @@ -114,6 +109,7 @@ fi if $cygwin ; then APP_HOME=`cygpath --path --mixed "$APP_HOME"` CLASSPATH=`cygpath --path --mixed "$CLASSPATH"` + JAVACMD=`cygpath --unix "$JAVACMD"` # We build the pattern for arguments to be converted via cygpath ROOTDIRSRAW=`find -L / -maxdepth 1 -mindepth 1 -type d 2>/dev/null` 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 75cd7ffb..766d67ab 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 @@ -21,6 +21,8 @@ import java.util.concurrent.Future; import org.apache.kafka.clients.producer.RecordMetadata; +import org.springframework.messaging.Message; + /** * The basic Kafka operations contract. * @@ -94,10 +96,19 @@ public interface KafkaOperations { */ Future send(String topic, int partition, K key, V data); + /** + * Send a message with routing information in message headers. + * @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); + // Sync methods - /** * Send the data to the default topic with no key or partition; * wait for result. @@ -180,6 +191,19 @@ public interface KafkaOperations { RecordMetadata syncSend(String topic, int partition, K key, V data) throws InterruptedException, ExecutionException; + /** + * Send a message with routing information in message headers. + * @param message the message to send. + * @return a Future for the {@link RecordMetadata}. + * @throws ExecutionException execution exception while awaiting result. + * @throws InterruptedException thread interrupted while awaiting result. + * @see org.springframework.kafka.support.KafkaHeaders#TOPIC + * @see org.springframework.kafka.support.KafkaHeaders#PARTITION_ID + * @see org.springframework.kafka.support.KafkaHeaders#MESSAGE_KEY + */ + RecordMetadata syncSend(Message message) + throws InterruptedException, ExecutionException; + /** * Flush the producer. */ 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 1a1594c9..ba1b11ae 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,9 +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.messaging.Message; +import org.springframework.messaging.MessageHeaders; /** @@ -88,28 +91,28 @@ public class KafkaTemplate implements KafkaOperations { } @Override - public Future send(V data) { + public Future send(V data) { return send(this.defaultTopic, data); } @Override - public Future send(K key, V data) { + public Future send(K key, V data) { return send(this.defaultTopic, key, data); } @Override - public Future send(int partition, K key, V data) { + public Future send(int partition, K key, V data) { return send(this.defaultTopic, partition, key, data); } @Override - public Future send(String topic, V data) { + public Future send(String topic, V data) { ProducerRecord producerRecord = new ProducerRecord<>(topic, data); return doSend(producerRecord); } @Override - public Future send(String topic, K key, V data) { + public Future send(String topic, K key, V data) { ProducerRecord producerRecord = new ProducerRecord<>(topic, key, data); return doSend(producerRecord); } @@ -121,11 +124,16 @@ public class KafkaTemplate implements KafkaOperations { } @Override - public Future send(String topic, int partition, K key, V data) { + public Future send(String topic, int partition, K key, V data) { ProducerRecord producerRecord = new ProducerRecord<>(topic, partition, key, data); return doSend(producerRecord); } + @Override + public Future send(Message message) { + ProducerRecord producerRecord = messageToProducerRecord(message); + return doSend(producerRecord); + } @Override public RecordMetadata syncSend(V data) throws InterruptedException, ExecutionException { @@ -180,6 +188,19 @@ public class KafkaTemplate implements KafkaOperations { return future.get(); } + @Override + public RecordMetadata syncSend(Message message) + throws InterruptedException, ExecutionException { + Future future = send(message); + flush(); + return future.get(); + } + + @Override + public void flush() { + this.producer.flush(); + } + /** * Send the producer record. * @param producerRecord the producer record. @@ -211,9 +232,14 @@ public class KafkaTemplate implements KafkaOperations { return future; } - @Override - public void flush() { - this.producer.flush(); + @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/support/KafkaHeaders.java b/spring-kafka/src/main/java/org/springframework/kafka/support/KafkaHeaders.java index ed128b6e..9a5ce0b5 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/support/KafkaHeaders.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/support/KafkaHeaders.java @@ -21,24 +21,24 @@ package org.springframework.kafka.support; * * @author Artem Bilan * @author Marius Bogoevici + * @author Gary Russell */ public abstract class KafkaHeaders { private static final String PREFIX = "kafka_"; /** - * The header for topic. + * The header containing the topic when sending data to Kafka. */ public static final String TOPIC = PREFIX + "topic"; - /** - * The header for message key. + * The header containing the message key when sending data to Kafka. */ public static final String MESSAGE_KEY = PREFIX + "messageKey"; /** - * The header for topic partition. + * The header containing the topic partition when sending data to Kafka. */ public static final String PARTITION_ID = PREFIX + "partitionId"; @@ -52,4 +52,19 @@ public abstract class KafkaHeaders { */ public static final String ACKNOWLEDGMENT = PREFIX + "acknowledgment"; + /** + * The header containing the topic from which the message was received. + */ + public static final String RECEIVED_TOPIC = PREFIX + "receivedTopic"; + + /** + * The header containing the message key for the received message. + */ + public static final String RECEIVED_MESSAGE_KEY = PREFIX + "receivedMessageKey"; + + /** + * The header containing the topic partition for the received message. + */ + public static final String RECEIVED_PARTITION_ID = PREFIX + "receivedPartitionId"; + } 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 cca74f41..3ae649c8 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 @@ -66,9 +66,9 @@ public class MessagingMessageConverter implements MessageConverter { KafkaMessageHeaders kafkaMessageHeaders = new KafkaMessageHeaders(this.generateMessageId, this.generateTimestamp); Map rawHeaders = kafkaMessageHeaders.getRawHeaders(); - rawHeaders.put(KafkaHeaders.MESSAGE_KEY, record.key()); - rawHeaders.put(KafkaHeaders.TOPIC, record.topic()); - rawHeaders.put(KafkaHeaders.PARTITION_ID, record.partition()); + rawHeaders.put(KafkaHeaders.RECEIVED_MESSAGE_KEY, record.key()); + rawHeaders.put(KafkaHeaders.RECEIVED_TOPIC, record.topic()); + rawHeaders.put(KafkaHeaders.RECEIVED_PARTITION_ID, record.partition()); rawHeaders.put(KafkaHeaders.OFFSET, record.offset()); if (acknowledgment != null) { 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 b9147639..ab206a1c 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 @@ -83,10 +83,12 @@ public class EnableKafkaIntegrationTests { assertThat(this.listener.latch1.await(10, TimeUnit.SECONDS)).isTrue(); waitListening("bar"); - template.send("annotated2", 0, "foo"); + template.send("annotated2", 0, 123, "foo"); template.flush(); assertThat(this.listener.latch2.await(10, TimeUnit.SECONDS)).isTrue(); + assertThat(this.listener.key).isEqualTo(123); assertThat(this.listener.partition).isNotNull(); + assertThat(this.listener.topic).isEqualTo("annotated2"); waitListening("baz"); template.send("annotated3", 0, "foo"); @@ -200,14 +202,23 @@ public class EnableKafkaIntegrationTests { private volatile Acknowledgment ack; + private Integer key; + + private String topic; + @KafkaListener(id = "foo", topics = "annotated1") public void listen1(String foo) { this.latch1.countDown(); } @KafkaListener(id = "bar", topicPattern = "annotated2") - public void listen2(@Payload String foo, @Header(KafkaHeaders.PARTITION_ID) int partitionHeader) { - this.partition = partitionHeader; + public void listen2(@Payload String foo, + @Header(KafkaHeaders.RECEIVED_MESSAGE_KEY) Integer key, + @Header(KafkaHeaders.RECEIVED_PARTITION_ID) int partition, + @Header(KafkaHeaders.RECEIVED_TOPIC) String topic) { + this.key = key; + this.partition = partition; + this.topic = topic; this.latch2.countDown(); } 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 c923c25a..98695840 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 @@ -35,9 +35,11 @@ import org.junit.Test; import org.springframework.kafka.listener.ContainerTestUtils; import org.springframework.kafka.listener.KafkaMessageListenerContainer; import org.springframework.kafka.listener.MessageListener; +import org.springframework.kafka.support.KafkaHeaders; import org.springframework.kafka.support.ProducerListenerAdapter; import org.springframework.kafka.test.rule.KafkaEmbedded; import org.springframework.kafka.test.utils.KafkaTestUtils; +import org.springframework.messaging.support.MessageBuilder; /** @@ -64,7 +66,6 @@ public class KafkaTemplateTests { @Override public void onMessage(ConsumerRecord record) { - System.out.println(record); records.add(record); } @@ -93,6 +94,23 @@ public class KafkaTemplateTests { assertThat(received).has(key((Integer) null)); assertThat(received).has(partition(0)); assertThat(received).has(value("qux")); + template.syncSend(MessageBuilder.withPayload("fiz") + .setHeader(KafkaHeaders.TOPIC, TEMPLATE_TOPIC) + .setHeader(KafkaHeaders.PARTITION_ID, 0) + .setHeader(KafkaHeaders.MESSAGE_KEY, 2) + .build()); + received = records.poll(10, TimeUnit.SECONDS); + assertThat(received).has(key(2)); + assertThat(received).has(partition(0)); + assertThat(received).has(value("fiz")); + template.syncSend(MessageBuilder.withPayload("buz") + .setHeader(KafkaHeaders.PARTITION_ID, 0) + .setHeader(KafkaHeaders.MESSAGE_KEY, 2) + .build()); + received = records.poll(10, TimeUnit.SECONDS); + assertThat(received).has(key(2)); + assertThat(received).has(partition(0)); + assertThat(received).has(value("buz")); } @Test diff --git a/src/reference/asciidoc/kafka.adoc b/src/reference/asciidoc/kafka.adoc index 766bb9d4..752a58b2 100644 --- a/src/reference/asciidoc/kafka.adoc +++ b/src/reference/asciidoc/kafka.adoc @@ -24,6 +24,7 @@ Future send(String topic, int partition, V data); Future send(String topic, int partition, K key, V data); +Future send(Message message); // Sync methods @@ -49,6 +50,9 @@ RecordMetadata syncSend(String topic, int partition, V data) RecordMetadata syncSend(String topic, int partition, K key, V data) throws InterruptedException, ExecutionException; +RecordMetadata syncSend(Message message) + throws InterruptedException, ExecutionException; + // Flush the producer. void flush(); @@ -81,6 +85,15 @@ The template can also be configured using standard `` definitions. Then, to use the template, simply invoke one of its methods. +When using the methods with a `Message` parameter, topic, partition and key information is provided in a message +header: + +- `KafkaHeaders.TOPIC` +- `KafkaHeaders.PARTITION_ID` +- `KafkaHeaders.MESSAGE_KEY` + +with the message payload being the data. + Optionally, you can configure the `KafkaTemplate` with a `ProducerListener` to get an async callback with the results of the send (success or failure) instead of waiting for the `Future` to complete. @@ -276,3 +289,16 @@ public void listen(String data, Acknowledgment ack) { ack.acknowledge(); } ---- + +Finally, metadata about the message is available from message headers: + +[source, java] +---- +@KafkaListener(id = "qux", topicPattern = "myTopic1") +public void listen(@Payload String foo, + @Header(KafkaHeaders.RECEIVED_MESSAGE_KEY) Integer key, + @Header(KafkaHeaders.RECEIVED_PARTITION_ID) int partition, + @Header(KafkaHeaders.RECEIVED_TOPIC) String topic) { + ... +} +----