diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaOutboundChannelAdapterParser.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaOutboundChannelAdapterParser.java index 657b3c3b16..29ef3e8740 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaOutboundChannelAdapterParser.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaOutboundChannelAdapterParser.java @@ -67,9 +67,6 @@ public class KafkaOutboundChannelAdapterParser extends AbstractOutboundChannelAd kafkaProducerMessageHandlerBuilder.addPropertyValue("partitionIdExpression", partitionIdExpressionDef); } - IntegrationNamespaceUtils.setValueIfAttributeDefined(kafkaProducerMessageHandlerBuilder, element, - "enable-header-routing"); - return kafkaProducerMessageHandlerBuilder.getBeanDefinition(); } 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 f8120e6090..f082e2c161 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 @@ -24,6 +24,7 @@ import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.support.KafkaHeaders; import org.springframework.messaging.Message; import org.springframework.util.Assert; +import org.springframework.util.StringUtils; /** * Kafka Message Handler. @@ -43,8 +44,6 @@ public class KafkaProducerMessageHandler extends AbstractMessageHandler { private EvaluationContext evaluationContext; - private boolean enableHeaderRouting = true; - private volatile Expression topicExpression; private volatile Expression messageKeyExpression; @@ -56,19 +55,6 @@ public class KafkaProducerMessageHandler extends AbstractMessageHandler { this.kafkaTemplate = kafkaTemplate; } - /** - * Enable the use of headers for determining the target topic and partition of outbound messages. By default it is - * set to true, but it can be disabled when those values are produced by upstream components that read messages - * from Kafka sources themselves. - * @param enableHeaderRouting whether the topic and destination headers should be considered - * @since 1.3 - * @see KafkaHeaders#TOPIC - * @see KafkaHeaders#PARTITION_ID - */ - public void setEnableHeaderRouting(boolean enableHeaderRouting) { - this.enableHeaderRouting = enableHeaderRouting; - } - public void setTopicExpression(Expression topicExpression) { this.topicExpression = topicExpression; } @@ -77,16 +63,6 @@ public class KafkaProducerMessageHandler extends AbstractMessageHandler { this.messageKeyExpression = messageKeyExpression; } - /** - * Set the partition expression. - * @param partitionExpression an expression that returns a partition id - * @deprecated as of 1.3, {@link #setPartitionIdExpression(Expression)} should be used instead - */ - @Deprecated - public void setPartitionExpression(Expression partitionExpression) { - setPartitionIdExpression(partitionExpression); - } - public void setPartitionIdExpression(Expression partitionIdExpression) { this.partitionIdExpression = partitionIdExpression; } @@ -106,12 +82,13 @@ public class KafkaProducerMessageHandler extends AbstractMessageHandler { protected void handleMessageInternal(final Message message) throws Exception { String topic = this.topicExpression != null ? this.topicExpression.getValue(this.evaluationContext, message, String.class) - //TODO revise the headers fallback behavior in favor of just expression - : (this.enableHeaderRouting ? message.getHeaders().get(KafkaHeaders.TOPIC, String.class) : null); + : message.getHeaders().get(KafkaHeaders.TOPIC, String.class); + + Assert.state(StringUtils.hasText(topic), "The 'topic' can not be empty or null"); Integer partitionId = this.partitionIdExpression != null ? this.partitionIdExpression.getValue(this.evaluationContext, message, Integer.class) - : (this.enableHeaderRouting ? message.getHeaders().get(KafkaHeaders.PARTITION_ID, Integer.class) : null); + : message.getHeaders().get(KafkaHeaders.PARTITION_ID, Integer.class); Object messageKey = this.messageKeyExpression != null ? this.messageKeyExpression.getValue(this.evaluationContext, message) diff --git a/spring-integration-kafka/src/main/resources/org/springframework/integration/kafka/config/spring-integration-kafka-2.0.xsd b/spring-integration-kafka/src/main/resources/org/springframework/integration/kafka/config/spring-integration-kafka-2.0.xsd index 96208456b8..8cd912deb2 100644 --- a/spring-integration-kafka/src/main/resources/org/springframework/integration/kafka/config/spring-integration-kafka-2.0.xsd +++ b/spring-integration-kafka/src/main/resources/org/springframework/integration/kafka/config/spring-integration-kafka-2.0.xsd @@ -96,14 +96,6 @@ ]]> - - - - - 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 e0acb89455..6d32da9126 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 @@ -52,7 +52,6 @@ public class KafkaOutboundAdapterParserTests { assertThat(TestUtils.getPropertyValue(messageHandler, "topicExpression.literalValue")).isEqualTo("foo"); assertThat(TestUtils.getPropertyValue(messageHandler, "messageKeyExpression.expression")).isEqualTo("'bar'"); assertThat(TestUtils.getPropertyValue(messageHandler, "partitionIdExpression.expression")).isEqualTo("2"); - assertThat(TestUtils.getPropertyValue(messageHandler, "enableHeaderRouting")).isEqualTo(Boolean.TRUE); } }