From 710ff2c2923551cb1c6d841c14e24958c2c5ecda Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Wed, 28 Mar 2018 13:32:37 -0400 Subject: [PATCH] Fix NPE in Kafka Streams binder * Fix NPE in Kafka Streams binder Fix NPE when user provided consumer properties are missing in kafka streams binder Resolves #343 * Addressing PR review comments --- .../KafkaStreamsStreamListenerSetupMethodOrchestrator.java | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java index 18a947c23..228e3763b 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java @@ -221,7 +221,7 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator implements StreamListene enableNativeDecodingForKTableAlways(parameterType, bindingProperties); StreamsConfig streamsConfig = null; //Retrieve the StreamsConfig created for this method if available. - //Otherwise, carete the StreamsBuilderFactory and get the underlying config. + //Otherwise, create the StreamsBuilderFactory and get the underlying config. if (!methodStreamsBuilderFactoryBeanMap.containsKey(method)) { streamsConfig = buildStreamsBuilderAndRetrieveConfig(method, applicationContext, bindingProperties); } @@ -291,7 +291,8 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator implements StreamListene Serde keySerde, Serde valueSerde) { KStream stream = streamsBuilder.stream(bindingServiceProperties.getBindingDestination(inboundName), Consumed.with(keySerde, valueSerde)); - if (bindingProperties.getConsumer().isUseNativeDecoding()){ + final boolean nativeDecoding = bindingServiceProperties.getConsumerProperties(inboundName).isUseNativeDecoding(); + if (nativeDecoding){ LOG.info("Native decoding is enabled for " + inboundName + ". Inbound deserialization done at the broker."); } else { @@ -300,7 +301,7 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator implements StreamListene stream = stream.map((key, value) -> { KeyValue keyValue; String contentType = bindingProperties.getContentType(); - if (!StringUtils.isEmpty(contentType) && !bindingProperties.getConsumer().isUseNativeDecoding()) { + if (!StringUtils.isEmpty(contentType) && !nativeDecoding) { Message message = MessageBuilder.withPayload(value) .setHeader(MessageHeaders.CONTENT_TYPE, contentType).build(); keyValue = new KeyValue<>(key, message);