diff --git a/docs/src/main/asciidoc/kafka-streams.adoc b/docs/src/main/asciidoc/kafka-streams.adoc index 5cc9aa268..53c6e60a9 100644 --- a/docs/src/main/asciidoc/kafka-streams.adoc +++ b/docs/src/main/asciidoc/kafka-streams.adoc @@ -1468,7 +1468,7 @@ Kafka Streams binder provides the following actuator endpoints for retrieving th `/actuator/kafkastreamstopology` -`/actuator/kafkastreamstopology/` +`/actuator/kafkastreamstopology/` You need to include the actuator and web dependencies from Spring Boot to access these endpoints. Further, you also need to add `kafkastreamstopology` to `management.endpoints.web.exposure.include` property. diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/endpoint/KafkaStreamsTopologyEndpoint.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/endpoint/KafkaStreamsTopologyEndpoint.java index 3ee1dcb88..598a0165f 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/endpoint/KafkaStreamsTopologyEndpoint.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/endpoint/KafkaStreamsTopologyEndpoint.java @@ -16,6 +16,7 @@ package org.springframework.cloud.stream.binder.kafka.streams.endpoint; +import java.util.ArrayList; import java.util.List; import org.springframework.boot.actuate.endpoint.annotation.Endpoint; @@ -46,13 +47,14 @@ public class KafkaStreamsTopologyEndpoint { } @ReadOperation - public String kafkaStreamsTopology() { + public List kafkaStreamsTopologies() { final List streamsBuilderFactoryBeans = this.kafkaStreamsRegistry.streamsBuilderFactoryBeans(); final StringBuilder topologyDescription = new StringBuilder(); + final List descs = new ArrayList<>(); streamsBuilderFactoryBeans.stream() .forEach(streamsBuilderFactoryBean -> - topologyDescription.append(streamsBuilderFactoryBean.getTopology().describe().toString())); - return topologyDescription.toString(); + descs.add(streamsBuilderFactoryBean.getTopology().describe().toString())); + return descs; } @ReadOperation diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountFunctionTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountFunctionTests.java index 2156ef264..e086bd286 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountFunctionTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountFunctionTests.java @@ -18,6 +18,7 @@ package org.springframework.cloud.stream.binder.kafka.streams.function; import java.util.Arrays; import java.util.Date; +import java.util.List; import java.util.Map; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; @@ -112,7 +113,8 @@ public class KafkaStreamsBinderWordCountFunctionTests { //Testing topology endpoint final KafkaStreamsRegistry kafkaStreamsRegistry = context.getBean(KafkaStreamsRegistry.class); final KafkaStreamsTopologyEndpoint kafkaStreamsTopologyEndpoint = new KafkaStreamsTopologyEndpoint(kafkaStreamsRegistry); - final String topology1 = kafkaStreamsTopologyEndpoint.kafkaStreamsTopology(); + final List topologies = kafkaStreamsTopologyEndpoint.kafkaStreamsTopologies(); + final String topology1 = topologies.get(0); final String topology2 = kafkaStreamsTopologyEndpoint.kafkaStreamsTopology("testKstreamWordCountFunction"); assertThat(topology1).isNotEmpty(); assertThat(topology1).isEqualTo(topology2);