From 36c39974adeef319dd303254b81d16b475a5bbda Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Mon, 22 Jan 2018 13:33:40 -0500 Subject: [PATCH] Fix default partition header expression --- .../cloud/stream/binder/kafka/KafkaMessageChannelBinder.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 4bee676d4..59698bafc 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 @@ -701,7 +701,7 @@ public class KafkaMessageChannelBinder extends setBeanFactory(KafkaMessageChannelBinder.this.getBeanFactory()); if (producerProperties.isPartitioned()) { SpelExpressionParser parser = new SpelExpressionParser(); - setPartitionIdExpression(parser.parseExpression("headers." + BinderHeaders.PARTITION_HEADER)); + setPartitionIdExpression(parser.parseExpression("headers['" + BinderHeaders.PARTITION_HEADER + "']")); } if (producerProperties.getExtension().isSync()) { setSync(true);