diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderHealthIndicator.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderHealthIndicator.java index 2722759cc..0d1f06052 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderHealthIndicator.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderHealthIndicator.java @@ -115,7 +115,9 @@ public class KafkaStreamsBinderHealthIndicator extends AbstractHealthIndicator { } finally { // Close admin client immediately. - adminClient.close(Duration.ofSeconds(0)); + if (adminClient != null) { + adminClient.close(Duration.ofSeconds(0)); + } } } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/serde/CollectionSerde.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/serde/CollectionSerde.java index 91f9c92c5..01573d20e 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/serde/CollectionSerde.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/serde/CollectionSerde.java @@ -107,11 +107,12 @@ public class CollectionSerde implements Serde> { */ public CollectionSerde(Class targetTypeForJsonSerde, Class collectionsClass) { this.collectionClass = collectionsClass; - JsonSerde jsonSerde = new JsonSerde(targetTypeForJsonSerde); + try (JsonSerde jsonSerde = new JsonSerde(targetTypeForJsonSerde)) { - this.inner = Serdes.serdeFrom( - new CollectionSerializer<>(jsonSerde.serializer()), - new CollectionDeserializer<>(jsonSerde.deserializer(), collectionsClass)); + this.inner = Serdes.serdeFrom( + new CollectionSerializer<>(jsonSerde.serializer()), + new CollectionDeserializer<>(jsonSerde.deserializer(), collectionsClass)); + } } @Override @@ -204,8 +205,10 @@ public class CollectionSerde implements Serde> { final int records = dataInputStream.readInt(); for (int i = 0; i < records; i++) { final byte[] valueBytes = new byte[dataInputStream.readInt()]; - dataInputStream.read(valueBytes); - collection.add(valueDeserializer.deserialize(topic, valueBytes)); + final int read = dataInputStream.read(valueBytes); + if (read != -1) { + collection.add(valueDeserializer.deserialize(topic, valueBytes)); + } } } catch (IOException e) { diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderHealthIndicator.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderHealthIndicator.java index d44d18dc8..b967edf7a 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderHealthIndicator.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderHealthIndicator.java @@ -99,15 +99,16 @@ public class KafkaBinderHealthIndicator implements HealthIndicator, DisposableBe } } + private synchronized Consumer initMetadataConsumer() { + if (this.metadataConsumer == null) { + this.metadataConsumer = this.consumerFactory.createConsumer(); + } + return this.metadataConsumer; + } + private Health buildHealthStatus() { try { - if (this.metadataConsumer == null) { - synchronized (KafkaBinderHealthIndicator.this) { - if (this.metadataConsumer == null) { - this.metadataConsumer = this.consumerFactory.createConsumer(); - } - } - } + initMetadataConsumer(); synchronized (this.metadataConsumer) { Set downMessages = new HashSet<>(); final Map topicsInUse = KafkaBinderHealthIndicator.this.binder diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetrics.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetrics.java index 88e802bea..6385448ef 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetrics.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetrics.java @@ -174,31 +174,26 @@ public class KafkaBinderMetrics } } - private ConsumerFactory createConsumerFactory() { + private synchronized ConsumerFactory createConsumerFactory() { if (this.defaultConsumerFactory == null) { - synchronized (this) { - if (this.defaultConsumerFactory == null) { - Map props = new HashMap<>(); - props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, - ByteArrayDeserializer.class); - props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, - ByteArrayDeserializer.class); - Map mergedConfig = this.binderConfigurationProperties - .mergedConsumerConfiguration(); - if (!ObjectUtils.isEmpty(mergedConfig)) { - props.putAll(mergedConfig); - } - if (!props.containsKey(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG)) { - props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, - this.binderConfigurationProperties - .getKafkaConnectionString()); - } - this.defaultConsumerFactory = new DefaultKafkaConsumerFactory<>( - props); - } + Map props = new HashMap<>(); + props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, + ByteArrayDeserializer.class); + props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, + ByteArrayDeserializer.class); + Map mergedConfig = this.binderConfigurationProperties + .mergedConsumerConfiguration(); + if (!ObjectUtils.isEmpty(mergedConfig)) { + props.putAll(mergedConfig); } + if (!props.containsKey(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG)) { + props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, + this.binderConfigurationProperties + .getKafkaConnectionString()); + } + this.defaultConsumerFactory = new DefaultKafkaConsumerFactory<>( + props); } - return this.defaultConsumerFactory; }