GH-180: Add headerMapper() to DSL

Fixes https://github.com/spring-projects/spring-integration-kafka/issues/180
This commit is contained in:
Gary Russell
2017-10-20 11:31:42 -04:00
committed by Artem Bilan
parent 7aa04d3224
commit b65f59cd67
2 changed files with 24 additions and 1 deletions

View File

@@ -31,6 +31,8 @@ import org.springframework.integration.expression.ValueExpression;
import org.springframework.integration.kafka.outbound.KafkaProducerMessageHandler;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.core.ProducerFactory;
import org.springframework.kafka.support.DefaultKafkaHeaderMapper;
import org.springframework.kafka.support.KafkaHeaderMapper;
import org.springframework.kafka.support.LoggingProducerListener;
import org.springframework.kafka.support.ProducerListener;
import org.springframework.kafka.support.converter.RecordMessageConverter;
@@ -257,6 +259,24 @@ public class KafkaProducerMessageHandlerSpec<K, V, S extends KafkaProducerMessag
return _this();
}
/**
* Add a default header mapper to map spring messaging headers to Kafka headers.
* @return the spec.
*/
public S headerMapper() {
return headerMapper(new DefaultKafkaHeaderMapper());
}
/**
* Specify a header mapper to map spring messaging headers to Kafka headers.
* @param mapper the mapper.
* @return the spec.
*/
public S headerMapper(KafkaHeaderMapper mapper) {
this.target.setHeaderMapper(mapper);
return _this();
}
/**
* A {@link KafkaTemplate}-based {@link KafkaProducerMessageHandlerSpec} extension.
*

View File

@@ -19,6 +19,7 @@ package org.springframework.integration.kafka.dsl;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
import java.util.Collections;
import java.util.Map;
import java.util.stream.Stream;
@@ -123,7 +124,7 @@ public class KafkaDslTests {
this.kafkaProducer1.setPartitionIdExpression(new ValueExpression<>(0));
this.kafkaProducer2.setPartitionIdExpression(new ValueExpression<>(0));
this.sendToKafkaFlowInput.send(new GenericMessage<>("foo"));
this.sendToKafkaFlowInput.send(new GenericMessage<>("foo", Collections.singletonMap("foo", "bar")));
for (int i = 0; i < 100; i++) {
Message<?> receive = this.listeningFromKafkaResults1.receive(20000);
@@ -139,6 +140,7 @@ public class KafkaDslTests {
assertThat(headers.get(KafkaHeaders.OFFSET)).isEqualTo((long) i);
assertThat(headers.get(KafkaHeaders.TIMESTAMP_TYPE)).isEqualTo("CREATE_TIME");
assertThat(headers.get(KafkaHeaders.RECEIVED_TIMESTAMP)).isEqualTo(1487694048633L);
assertThat(headers.get("foo")).isEqualTo("bar");
}
for (int i = 0; i < 100; i++) {
@@ -257,6 +259,7 @@ public class KafkaDslTests {
.messageKey(m -> m
.getHeaders()
.get(IntegrationMessageHeaderAccessor.SEQUENCE_NUMBER))
.headerMapper()
.partitionId(m -> 10)
.topicExpression("headers[kafka_topic] ?: '" + topic + "'")
.configureKafkaTemplate(t -> t.id("kafkaTemplate:" + topic));