diff --git a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java
index acb8f8959..7d17f2c95 100644
--- a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java
+++ b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java
@@ -23,6 +23,7 @@ import java.util.Map;
* @author Marius Bogoevici
* @author Ilayaperumal Gopinathan
* @author Soby Chacko
+ * @author Gary Russell
*
*
* Thanks to Laszlo Szabo for providing the initial patch for generic property support.
@@ -30,6 +31,28 @@ import java.util.Map;
*/
public class KafkaConsumerProperties {
+ public enum StartOffset {
+ earliest(-2L),
+ latest(-1L);
+
+ private final long referencePoint;
+
+ StartOffset(long referencePoint) {
+ this.referencePoint = referencePoint;
+ }
+
+ public long getReferencePoint() {
+ return this.referencePoint;
+ }
+ }
+
+ public enum StandardHeaders {
+ none,
+ id,
+ timestamp,
+ both
+ }
+
private boolean autoRebalanceEnabled = true;
private boolean autoCommitOffset = true;
@@ -48,6 +71,10 @@ public class KafkaConsumerProperties {
private String[] trustedPackages;
+ private StandardHeaders standardHeaders = StandardHeaders.none;
+
+ private String converterBeanName;
+
private Map configuration = new HashMap<>();
public boolean isAutoCommitOffset() {
@@ -98,21 +125,6 @@ public class KafkaConsumerProperties {
this.autoRebalanceEnabled = autoRebalanceEnabled;
}
- public enum StartOffset {
- earliest(-2L),
- latest(-1L);
-
- private final long referencePoint;
-
- StartOffset(long referencePoint) {
- this.referencePoint = referencePoint;
- }
-
- public long getReferencePoint() {
- return this.referencePoint;
- }
- }
-
public Map getConfiguration() {
return this.configuration;
}
@@ -144,4 +156,20 @@ public class KafkaConsumerProperties {
public void setDlqProducerProperties(KafkaProducerProperties dlqProducerProperties) {
this.dlqProducerProperties = dlqProducerProperties;
}
+ public StandardHeaders getStandardHeaders() {
+ return this.standardHeaders;
+ }
+
+ public void setStandardHeaders(StandardHeaders standardHeaders) {
+ this.standardHeaders = standardHeaders;
+ }
+
+ public String getConverterBeanName() {
+ return this.converterBeanName;
+ }
+
+ public void setConverterBeanName(String converterBeanName) {
+ this.converterBeanName = converterBeanName;
+ }
+
}
diff --git a/spring-cloud-stream-binder-kafka-docs/src/main/asciidoc/overview.adoc b/spring-cloud-stream-binder-kafka-docs/src/main/asciidoc/overview.adoc
index ea26c845f..2ebb3da4e 100644
--- a/spring-cloud-stream-binder-kafka-docs/src/main/asciidoc/overview.adoc
+++ b/spring-cloud-stream-binder-kafka-docs/src/main/asciidoc/overview.adoc
@@ -191,6 +191,16 @@ dlqProducerProperties::
All the properties available through kafka producer properties can be set through this property.
+
Default: Default Kafka producer properties.
+standardHeaders::
+ Indicates which standard headers are populated by the inbound channel adapter.
+ `none`, `id`, `timestamp` or `both`.
+ Useful if using native deserialization and the first component to receive a message needs an `id` (such as an aggregator that is configured to use a JDBC message store).
++
+Default: `none`
+converterBeanName::
+ The name of a bean that implements `RecordMessageConverter`; used in the inbound channel adapter to replace the default `MessagingMessageConverter`.
++
+Default: `null`
[[kafka-producer-properties]]
=== Kafka Producer Properties
diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java
index e1a65adf7..0c9db3b37 100644
--- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java
+++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java
@@ -43,6 +43,7 @@ import org.apache.kafka.common.serialization.ByteArrayDeserializer;
import org.apache.kafka.common.serialization.ByteArraySerializer;
import org.springframework.beans.factory.DisposableBean;
+import org.springframework.beans.factory.NoSuchBeanDefinitionException;
import org.springframework.cloud.stream.binder.AbstractMessageChannelBinder;
import org.springframework.cloud.stream.binder.BinderHeaders;
import org.springframework.cloud.stream.binder.ExtendedConsumerProperties;
@@ -51,6 +52,7 @@ import org.springframework.cloud.stream.binder.ExtendedPropertiesBinder;
import org.springframework.cloud.stream.binder.HeaderMode;
import org.springframework.cloud.stream.binder.kafka.properties.KafkaBinderConfigurationProperties;
import org.springframework.cloud.stream.binder.kafka.properties.KafkaConsumerProperties;
+import org.springframework.cloud.stream.binder.kafka.properties.KafkaConsumerProperties.StandardHeaders;
import org.springframework.cloud.stream.binder.kafka.properties.KafkaExtendedBindingProperties;
import org.springframework.cloud.stream.binder.kafka.properties.KafkaProducerProperties;
import org.springframework.cloud.stream.binder.kafka.provisioning.KafkaTopicProvisioner;
@@ -373,7 +375,25 @@ public class KafkaMessageChannelBinder extends
}
final KafkaMessageDrivenChannelAdapter, ?> kafkaMessageDrivenChannelAdapter = new KafkaMessageDrivenChannelAdapter<>(
messageListenerContainer);
- MessagingMessageConverter messageConverter = new MessagingMessageConverter();
+ MessagingMessageConverter messageConverter;
+ if (extendedConsumerProperties.getExtension().getConverterBeanName() == null) {
+ messageConverter = new MessagingMessageConverter();
+ StandardHeaders standardHeaders = extendedConsumerProperties.getExtension().getStandardHeaders();
+ messageConverter.setGenerateMessageId(StandardHeaders.id.equals(standardHeaders)
+ || StandardHeaders.both.equals(standardHeaders));
+ messageConverter.setGenerateTimestamp(StandardHeaders.timestamp.equals(standardHeaders)
+ || StandardHeaders.both.equals(standardHeaders));
+ }
+ else {
+ try {
+ messageConverter = getApplicationContext().getBean(
+ extendedConsumerProperties.getExtension().getConverterBeanName(),
+ MessagingMessageConverter.class);
+ }
+ catch (NoSuchBeanDefinitionException e) {
+ throw new IllegalStateException("Converter bean not present in application context", e);
+ }
+ }
KafkaHeaderMapper mapper = null;
if (this.configurationProperties.getHeaderMapperBeanName() != null) {
mapper = getApplicationContext().getBean(this.configurationProperties.getHeaderMapperBeanName(),
diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java
index ea9831778..06d506479 100644
--- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java
+++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java
@@ -74,11 +74,13 @@ import org.springframework.cloud.stream.binder.TestUtils;
import org.springframework.cloud.stream.binder.kafka.admin.KafkaAdminUtilsOperation;
import org.springframework.cloud.stream.binder.kafka.properties.KafkaBinderConfigurationProperties;
import org.springframework.cloud.stream.binder.kafka.properties.KafkaConsumerProperties;
+import org.springframework.cloud.stream.binder.kafka.properties.KafkaConsumerProperties.StandardHeaders;
import org.springframework.cloud.stream.binder.kafka.properties.KafkaProducerProperties;
import org.springframework.cloud.stream.binder.kafka.utils.KafkaTopicUtils;
import org.springframework.cloud.stream.config.BindingProperties;
import org.springframework.cloud.stream.provisioning.ProvisioningException;
import org.springframework.context.ApplicationContext;
+import org.springframework.context.ConfigurableApplicationContext;
import org.springframework.context.support.GenericApplicationContext;
import org.springframework.integration.IntegrationMessageHeaderAccessor;
import org.springframework.integration.channel.DirectChannel;
@@ -96,6 +98,7 @@ import org.springframework.kafka.support.Acknowledgment;
import org.springframework.kafka.support.KafkaHeaders;
import org.springframework.kafka.support.SendResult;
import org.springframework.kafka.support.TopicPartitionInitialOffset;
+import org.springframework.kafka.support.converter.MessagingMessageConverter;
import org.springframework.kafka.test.core.BrokerAddress;
import org.springframework.kafka.test.rule.KafkaEmbedded;
import org.springframework.kafka.test.utils.KafkaTestUtils;
@@ -1991,7 +1994,7 @@ public class KafkaBinderTests extends
}
@Test
- @SuppressWarnings("unchecked")
+ @SuppressWarnings({ "unchecked", "rawtypes" })
public void testNativeSerializationWithCustomSerializerDeserializer() throws Exception {
Binding> producerBinding = null;
Binding> consumerBinding = null;
@@ -2019,6 +2022,7 @@ public class KafkaBinderTests extends
consumerProperties.getExtension().setAutoRebalanceEnabled(false);
consumerProperties.getExtension().getConfiguration().put("value.deserializer",
"org.apache.kafka.common.serialization.IntegerDeserializer");
+ consumerProperties.getExtension().setStandardHeaders(StandardHeaders.both);
consumerBinding = binder.bindConsumer(testTopicName, "test", moduleInputChannel, consumerProperties);
// Let the consumer actually bind to the producer before sending a msg
binderBindUnbindLatency();
@@ -2027,6 +2031,8 @@ public class KafkaBinderTests extends
assertThat(inbound).isNotNull();
assertThat(inbound.getPayload()).isEqualTo(10);
assertThat(inbound.getHeaders()).doesNotContainKey("contentType");
+ assertThat(inbound.getHeaders().getId()).isNotNull();
+ assertThat(inbound.getHeaders().getTimestamp()).isNotNull();
}
finally {
if (producerBinding != null) {
@@ -2039,7 +2045,7 @@ public class KafkaBinderTests extends
}
@Test
- @SuppressWarnings("unchecked")
+ @SuppressWarnings({ "unchecked", "rawtypes" })
public void testNativeSerializationWithCustomSerializerDeserializerBytesPayload() throws Exception {
Binding> producerBinding = null;
Binding> consumerBinding = null;
@@ -2059,6 +2065,12 @@ public class KafkaBinderTests extends
invokeCreateTopic(zkUtils, testTopicName, 1, 1, new Properties());
configurationProperties.setAutoAddPartitions(true);
Binder binder = getBinder(configurationProperties);
+ ConfigurableApplicationContext context = TestUtils.getPropertyValue(binder, "binder.applicationContext",
+ ConfigurableApplicationContext.class);
+ MessagingMessageConverter converter = new MessagingMessageConverter();
+ converter.setGenerateMessageId(true);
+ converter.setGenerateTimestamp(true);
+ context.getBeanFactory().registerSingleton("testConverter", converter);
QueueChannel moduleInputChannel = new QueueChannel();
ExtendedProducerProperties producerProperties = createProducerProperties();
producerProperties.setUseNativeEncoding(true);
@@ -2071,6 +2083,7 @@ public class KafkaBinderTests extends
consumerProperties.getExtension()
.getConfiguration()
.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ByteArrayDeserializer.class.getName());
+ consumerProperties.getExtension().setConverterBeanName("testConverter");
consumerBinding = binder.bindConsumer(testTopicName, "test", moduleInputChannel, consumerProperties);
// Let the consumer actually bind to the producer before sending a msg
binderBindUnbindLatency();
@@ -2080,6 +2093,8 @@ public class KafkaBinderTests extends
assertThat(inbound.getPayload()).isEqualTo(new byte[1]);
assertThat(inbound.getHeaders()).containsKey("contentType");
assertThat(inbound.getHeaders().get(MessageHeaders.CONTENT_TYPE).toString()).isEqualTo("something/funky");
+ assertThat(inbound.getHeaders().getId()).isNotNull();
+ assertThat(inbound.getHeaders().getTimestamp()).isNotNull();
}
finally {
if (producerBinding != null) {
@@ -2270,7 +2285,7 @@ public class KafkaBinderTests extends
}
@Test
- @SuppressWarnings("unchecked")
+ @SuppressWarnings({ "unchecked", "rawtypes" })
public void testSendAndReceiveWithRawMode() throws Exception {
Binder binder = getBinder();