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 2706848f18..a74655c18b 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 @@ -87,7 +87,7 @@ ]]> - + + http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans.xsd"> @@ -18,7 +16,7 @@ order="3" topic="foo" message-key-expression="'bar'" - partition-id-expression="2"> + partition-id-expression="'2'"> 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 a201eca256..6787ea59e2 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 @@ -66,7 +66,7 @@ public class KafkaOutboundAdapterParserTests { assertThat(messageHandler.getOrder()).isEqualTo(3); 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, "partitionIdExpression.expression")).isEqualTo("'2'"); } diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/outbound/KafkaProducerMessageHandlerTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/outbound/KafkaProducerMessageHandlerTests.java index 5e42087f05..dcc4c3cfb8 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/outbound/KafkaProducerMessageHandlerTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/outbound/KafkaProducerMessageHandlerTests.java @@ -29,6 +29,7 @@ import org.junit.ClassRule; import org.junit.Test; import org.springframework.beans.factory.BeanFactory; +import org.springframework.expression.spel.standard.SpelExpressionParser; import org.springframework.integration.support.MessageBuilder; import org.springframework.kafka.core.ConsumerFactory; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; @@ -104,10 +105,12 @@ public class KafkaProducerMessageHandlerTests { assertThat(record).has(key((Integer) null)); assertThat(record).has(value("baz")); + handler.setPartitionIdExpression(new SpelExpressionParser().parseExpression("headers['kafka_partitionId']")); + message = MessageBuilder.withPayload(KafkaNull.INSTANCE) .setHeader(KafkaHeaders.TOPIC, topic1) .setHeader(KafkaHeaders.MESSAGE_KEY, 2) - .setHeader(KafkaHeaders.PARTITION_ID, 1) + .setHeader(KafkaHeaders.PARTITION_ID, "1") .build(); handler.handleMessage(message);