diff --git a/spring-cloud-stream-docs/src/main/asciidoc/images/SCSt-groups.png b/spring-cloud-stream-docs/src/main/asciidoc/images/SCSt-groups.png new file mode 100644 index 000000000..931ac9727 Binary files /dev/null and b/spring-cloud-stream-docs/src/main/asciidoc/images/SCSt-groups.png differ diff --git a/spring-cloud-stream-docs/src/main/asciidoc/images/SCSt-partitioning.png b/spring-cloud-stream-docs/src/main/asciidoc/images/SCSt-partitioning.png new file mode 100644 index 000000000..833d2ac66 Binary files /dev/null and b/spring-cloud-stream-docs/src/main/asciidoc/images/SCSt-partitioning.png differ diff --git a/spring-cloud-stream-docs/src/main/asciidoc/images/SCSt-sensors.png b/spring-cloud-stream-docs/src/main/asciidoc/images/SCSt-sensors.png new file mode 100644 index 000000000..6a15c8729 Binary files /dev/null and b/spring-cloud-stream-docs/src/main/asciidoc/images/SCSt-sensors.png differ diff --git a/spring-cloud-stream-docs/src/main/asciidoc/images/SCSt-with-binder.png b/spring-cloud-stream-docs/src/main/asciidoc/images/SCSt-with-binder.png new file mode 100644 index 000000000..b4d66fd5e Binary files /dev/null and b/spring-cloud-stream-docs/src/main/asciidoc/images/SCSt-with-binder.png differ diff --git a/spring-cloud-stream-docs/src/main/asciidoc/spring-cloud-stream-overview.adoc b/spring-cloud-stream-docs/src/main/asciidoc/spring-cloud-stream-overview.adoc index 83b8c538c..d38d64985 100644 --- a/spring-cloud-stream-docs/src/main/asciidoc/spring-cloud-stream-overview.adoc +++ b/spring-cloud-stream-docs/src/main/asciidoc/spring-cloud-stream-overview.adoc @@ -1,5 +1,5 @@ -[[spring-cloud-stream-overview]] -== Spring Cloud Stream Overview +[[spring-cloud-stream-reference]] += Spring Cloud Stream Reference Manual [partintro] -- @@ -7,18 +7,16 @@ This section goes into more detail about how you can work with Spring Cloud Stre such as creating and running stream applications. -- -=== Introducing Spring Cloud Stream +== Introducing Spring Cloud Stream -The Spring Cloud Stream project allows a user to develop and run messaging microservices using Spring Integration. -Just add `@EnableBinding` and run your app as a Spring Boot app (single application context). -Spring Cloud Stream applications connect to the physical broker through bindings, which link Spring Integration -channels to physical broker destinations, for either input (consumer bindings) or output (producer bindings). -The creation of the bindings, and therefore their broker-specific implementation is handled by a binder, which is -another important abstraction of Spring Cloud Stream. Binders abstract out the broker-specific implementation details. -In order to connect to a specific type of broker (e.g. Rabbit or Kafka) you just need to have the relevant binder -implementation on the classpath. +Spring Cloud Stream is a framework for building message-driven microservices. +Spring Cloud Stream builds upon Spring Boot to create DevOps friendly microservice applications and Spring Integration to provide connectivity to message brokers. +Spring Cloud Stream provides an opinionated configuration of message brokers, introducing the concepts of persistent pub/sub semantics, consumer groups and partitions across several middleware vendors. +This opinionated configuration provides the basis to create stream processing application. -Here's a sample source app (output channel only): +By adding `@EnableBinding` to your main application, you get immediate connectivity to a message broker and by adding `@StreamListener` to a method, you will receive events for stream processing. + +Here's a sample sink application for receiving external messages: [source,java] ---- @@ -30,6 +28,344 @@ public class StreamApplication { } } +@EnableBinding(Sink.class) +public class TimerSource { + + ... + + @StreamListener(Sink.INPUT) + public void processVote(Vote vote) { + votingService.recordVote(vote); + } +} +---- + +`@EnableBinding` is parameterized by one or more interfaces (in this case a single `Sink` interface), which declares input and/or output channels. +The interfaces `Source`, `Sink` and `Processor` are provided but you can define others. +Here's the definition of `Source`: + +[source,java] +---- +public interface Sink { + String INPUT = "input"; + + @Input(Sink.INPUT) + SubscribableChannel input(); +} +---- + +The `@Input` annotation is used to identify input channels (messages entering the app), and `@Output` is used to identify output channels (messages leaving the app). +These annotations are optionally parameterized by a channel name. If the name is not provided then the method name is used instead. +An implementation of the interface is created for you and can be used in the application context by autowiring it, e.g. into a test case: + +[source,java] +---- +@RunWith(SpringJUnit4ClassRunner.class) +@SpringApplicationConfiguration(classes = StreamApplication.class) +@WebAppConfiguration +@DirtiesContext +public class StreamApplicationTests { + + @Autowired + private Sink sink; + + @Test + public void contextLoads() { + assertNotNull(this.sink.input()); + } +} +---- + +== Spring Cloud Stream Main Concepts + +Spring Cloud Stream provides a number of abstractions and primitives that simplify writing message-driven microservices. +In this section we will provide an overview of: + +* Spring Cloud Stream application model together with the Binder abstraction +* Persistent publish-subscribe and consumer group support +* Partitioning +* Pluggable Binder API + + +=== Application structure + +A Spring Cloud Stream application consists of a middleware-neutral core that communicates with the outside world through input and output channels. +The channels are managed and injected into it by the framework, and a `Binder` connects them to the external brokers. +Different `Binder` implementations exist for different types of middleware, such as https://github.com/spring-cloud/spring-cloud-stream/tree/master/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka[Kafka], https://github.com/spring-cloud/spring-cloud-stream/tree/master/spring-cloud-stream-binders/spring-cloud-stream-binder-rabbit[Rabbit MQ], https://github.com/spring-cloud/spring-cloud-stream-binder-redis[Redis] or https://github.com/spring-cloud/spring-cloud-stream-binder-gemfire[Gemfire], and an extensible API allows you to write your own `Binder`. There is also https://github.com/spring-cloud/spring-cloud-stream/blob/master/spring-cloud-stream-test-support/src/main/java/org/springframework/cloud/stream/test/binder/TestSupportBinder.java[TestSupportBinder] that leaves the channel as-is so a test author can interact with the channels directly and easily assert on what is received. + +.Spring Cloud Stream Application +image::SCSt-with-binder.png[width=300,scaledwidth="50%"] + +Spring Cloud Stream uses Spring Boot for configuration, and the `Binder` makes it possible for Spring Cloud Stream applications to be flexible in terms of how it connects to the middleware. +For example, deployers can dynamically choose the destinations that these channels connect to at runtime (e.g. Kafka topics or Rabbit MQ exchanges). +This can be done through external configuration properties in any form that is supported by Spring Boot (application arguments, environment variables, `application.yml` files, etc). +Taking the sink example from the previous section, providing the `spring.cloud.stream.bindings.input.destination=raw-sensor-data` property to the application will cause it to read from the `raw-sensor-data` Kafka topic, or from a queue bound to the `raw-sensor-data` exchange in Rabbit MQ. See <> for more information on the available binder properties you can configure. You are also able to configure middleware specific properties, see <> for more information. + +Spring Cloud Stream will automatically detect and use a binder that is found on the classpath, so you can easily use different types of middleware with the same code, just by including a different binder at build time. +For more complex use cases, Spring Cloud Stream also provides the ability of packaging multiple binders within the same application and choosing what type of binder should be used at runtime, and even if multiple binders should be used at runtime for different channels. + + +==== Fat JAR + +Spring Cloud Stream applications can be run in standalone mode from your IDE for testing. To run in production you can create an executable (or "fat") JAR using the standard Spring Boot tooling provided for Maven or Gradle. + +=== Persistent publish subscribe and consumer groups + +Communication between different applications follows a publish-subscribe pattern, with data being broadcast through shared topics. +This can be seen in the following picture, which shows a typical deployment for a set of interacting Spring Cloud Stream applications. + +.Spring Cloud Stream Application topologies +image::SCSt-with-binder.png[width=300,scaledwidth="50%"] + +Data reported by sensors to an HTTP endpoint is sent to a common destination named `raw-sensor-data`, from where it is independently processed by a microservice that computes time windowed averages, as well as by a microservice that ingests the raw data into HDFS. +In order to do so, both applications will declare the topic as their input at runtime. +The publish-subscribe communication model reduces the complexity of both the producer and the consumer, and allows adding new applications to the topology without disrupting the existing flow. +For example, downstream from the average calculator we can have a component that calculates the highest temperature values in order to display and monitor them. +Later on, we can add an application that interprets the very same flow of averages for fault detection. +The fact that all the communication is done through shared topics rather than point to point queues reduces the coupling between microservices. + +While the concept of publish-subscribe messaging is not new, Spring Cloud Stream takes the extra step of making it an opinionated choice for its application model. +It also makes it easy for users to work with it across different platform by using the native support of the middleware. + +[[consumer-groups]] +==== Consumer Groups +While the publish subscribe model ensures that it is easy to connect multiple application by sharing a topic, it is equally important to be able to scale up by creating multiple instances of a given application. +When doing so, the different instances would find themselves in a competing consumer relationship with each other: only one of the instances is expected to handle the message. +Spring Cloud Stream models this behavior through the concept of a consumer group, which is similar to (and inspired by) the notion of consumer groups in Kafka. +Each consumer binding can specify a group name such as `spring.cloud.stream.bindings.input.group=hdfsWrite` or `spring.cloud.stream.bindings.input.group=average`, as shown in the picture. +All groups that subscribe to a given destination will receive a copy of the published data, but only one member of the group will receive a given message from that destination. +By default, when a group is not specified, Spring Cloud Stream assigns the application to an anonymous, independent, single-member consumer group that will be in a publish-subscribe relationship with all the other consumer groups. + +.Spring Cloud Stream Consumer Groups +image::SCSt-groups.png[width=300,scaledwidth="50%"] + +[[durability]] +==== Durability + +Consistent with the opinionated application model of Spring Cloud Stream, consumer group subscriptions are durable. +This is to say that the binder implementation will ensure that group subscriptions are persistent and, once at least one subscription for a group has been created, that group will receive messages, even if they are sent while all the applications of the group were stopped. +Anonymous subscriptions are non-durable by nature. For some binder implementations (e.g. Rabbit) it is possible to have non-durable group subscriptions. + +In general, it is preferable to always specify a consumer group when binding an application to a given destination. +When scaling up a Spring Cloud Stream application, a consumer group must be specified for each of its input bindings, in order to prevent its instances from receiving duplicate messages (unless that behavior is desired, which is a less common use case). + +[[partitioning]] +=== Partitioning + +Spring Cloud Stream provides support for partitioning data between multiple instances of a given application. +In a partitioned scenario, one or more producer application instances will send data to multiple consumer application instances, ensuring that data with common characteristics is processed by the same consumer instance. +The physical communication medium (e.g. the broker topic) is viewed as structured into multiple partitions. +This happens regardless of whether the broker type is naturally partitioned (e.g. Kafka) or not (e.g. Rabbit), Spring Cloud Stream provides a common abstraction for implementing partitioned processing use cases in a uniform fashion. + +.Spring Cloud Stream Partitioning +image::SCSt-partitioning.png[width=300,scaledwidth="50%"] + +Partitioning is a critical concept in stateful processing, where ensuring that all the related data is processed together is critical for either performance or consistency. +For example, in the time-windowed average calculation example, it is important that measurements from the same sensor land in the same application instance. + +Setting up a partitioned processing scenario requires configuring both the data producing and the data consuming end. + +== Programming model + +This section will describe the programming model of Spring Cloud Stream, which consists from a number of predefined annotations that can be used to declare bound inputs and output channels, as well as how to listen to them. + +=== Declaring and binding channels + +==== Triggering binding via `@EnableBinding` + +A Spring application becomes a Spring Cloud Stream application when the `@EnableBinding` annotation is applied to one of its configuration classes. `@EnableBinding` itself is meta-annotated with `@Configuration`, and triggers the configuration of Spring Cloud Stream infrastructure as follows: + +[source,java] +---- +... +@Import(...) +@Configuration +@EnableIntegration +public @interface EnableBinding { + ... + Class[] value() default {}; +} +---- + +`@EnableBinding` can be parameterized with one or more interface classes, containing methods that represent bindable components (typically message channels). + +NOTE: As of version 1.0, the only supported bindable component is the Spring Messaging `MessageChannel` and its extensions `SubscribableChannel` and `PollableChannel`. +It is intended for future versions to extend support to other types of components, using the same mechanism. In this documentation, we will continue to refer to channels. + +==== `@Input` and `@Output` + +A Spring Cloud Stream application can have an arbitrary number of input and output channels defined as `@Input` and `@Output` methods in an interface, as follows: +[source,java] +---- +public interface Barista { + + @Input + SubscribableChannel orders(); + + @Output + MessageChannel hotDrinks(); + + @Output + MessageChannel coldDrinks(); +} +---- + +Using this interface as a parameter to `@EnableBinding`, as in the following example, will trigger the creation of three bound channels named `orders`, `hotDrinks` and `coldDrinks` respectively. + +[source,java] +---- +@EnableBinding(Barista.class) +public class CafeConfiguration { + + ... +} +---- + +===== Customizing channel names + +Both @Input and @Output allow specifying a customized name for the channel, as follows: + +[source,java] +---- +public interface Barista { + ... + @Input("inboundOrders") + SubscribableChannel orders(); +} +---- +In this case, the name of the bound channel being created will be `inboundOrders`. + +===== `Source`, `Sink`, and `Processor` + +For ease of addressing the most common use cases that involve either an input or an output channel, or both, out of the box Spring Cloud Stream provides three predefined interfaces. + +`Source` can be used for applications that have a single outbound channel. + +[source,java] +---- +public interface Source { + + String OUTPUT = "output"; + + @Output(Source.OUTPUT) + MessageChannel output(); + +} +---- + +`Sink` can be used for applications that have a single inbound channel. + +[source,java] +---- +public interface Sink { + + String INPUT = "input"; + + @Input(Sink.INPUT) + SubscribableChannel input(); + +} +---- + +`Processor` can be used for applications that have both an inbound and an outbound channel. + +[source,java] +---- +public interface Processor extends Source, Sink { +} +---- + +There is no special handling for either of these interfaces in Spring Cloud Stream, besides of the fact that they are provided out of the box. + +==== Accessing bound channels + +===== Injecting the bound interfaces + +For each of the bound interfaces, Spring Cloud Stream will generate a bean that implements it, and for which invoking an `@Input` or `@Output` annotated method will return the bound channel. +For example, the bean in the following example will send a message on the output channel every time its `hello` method is invoked, using the injected `Source` bean, and invoking `output()` to retrieve the target channel. + +[source,java] +---- +@Component +public class SendingBean { + + private Source source; + + @Autowired + public SendingBean(Source source) { + this.source = source; + } + + public void sayHello(String name) { + source.output().send(MessageBuilder.withPayload(body).build()); + } +} +---- + +===== Injecting channels directly + +Bound channels can be also injected directly. For example: + +[source, java] +---- +@Component +public class SendingBean { + + private MessageChannel output; + + @Autowired + public SendingBean(MessageChannel output) { + this.output = output; + } + + public void sayHello(String name) { + output.send(MessageBuilder.withPayload(body).build()); + } +} +---- + +Note that if the name of the channel is customized on the declaring annotation, that name should be used instead of the method name. Considering this declaration: + +[source,java] +---- +public interface CustomSource { + ... + @Output("customOutput") + MessageChannel output(); +} +---- + +The channel will be injected as follows: + +[source, java] +---- +@Component +public class SendingBean { + + @Autowired + private MessageChannel output; + + @Autowired @Qualifier("customOutput") + public SendingBean(MessageChannel output) { + this.output = output; + } + + public void sayHello(String name) { + customOutput.send(MessageBuilder.withPayload(body).build()); + } +} +---- + +==== Programming model + +Spring Cloud Stream allows you to write applications by either using Spring Integration annotations or Spring Cloud Stream's `@StreamListener` annotation which is modeled after other Spring Messaging annotations (e.g. `@MessageMapping`, `@JmsListener`, `@RabbitListener`, etc.) but add content type management and type coercion features. + +===== Native Spring Integration support + +Due to the fact that Spring Cloud Stream is Spring Integration based, it completely inherits its foundation and infrastructure, as well as the component. For example, the output channel of a `Source` can be attached to a `MessageSource`, as follows: + +[source, java] +---- @EnableBinding(Source.class) public class TimerSource { @@ -44,274 +380,65 @@ public class TimerSource { } ---- -`@EnableBinding` is parameterized by one or more interfaces (in this case a single `Source` interface), which declares -input and/or output channels. The interfaces `Source`, `Sink` and `Processor` are provided off the shelf, but you can -define others. Here's the definition of `Source`: +Or, the channels of a processor can be used in a transformer, as follows: [source,java] ---- -public interface Source { - String OUTPUT = "output"; - - @Output(Source.OUTPUT) - MessageChannel output(); +@EnableBinding(Processor.class) +public class TransformProcessor { + @Transformer(inputChannel = Processor.INPUT, outputChannel = Processor.OUTPUT) + public Object transform(String message) { + return message.toUpper(); + } } ---- -The `@Output` annotation is used to identify output channels (messages leaving the app), and `@Input` is used to -identify input channels (messages entering the app). It is optionally parameterized by a channel name - if the name is -not provided the method name is used instead. An implementation of the interface is created for you and can be used in -the application context by autowiring it, e.g. into a test case: +===== @StreamListener for automatic content type handling + +Complementary to the Spring Integration support, Spring Cloud Stream provides a `@StreamListener` annotation of its own modeled by the other similar Spring Messaging annotations (e.g. `@MessageMapping`, `@JmsListener`, `@RabbitListener`, etc.). +It provides a simpler model for handling inbound messages, especially for dealing with use cases that involve content type management and type coercion. +Spring Cloud Stream provides an extensible `MessageConverter` mechanism for handling data conversion by bound channels and, in this case, for dispatching to `@StreamListener` annotated methods. + +For example, an application that processes external `Vote` events can be declared as follows: [source,java] ---- -@RunWith(SpringJUnit4ClassRunner.class) -@SpringApplicationConfiguration(classes = StreamApplication.class) -@WebAppConfiguration -@DirtiesContext -public class StreamApplicationTests { +@EnableBinding(Sink.class) +public class VoteHandler { @Autowired - private Source source + VotingService votingService; - @Test - public void contextLoads() { - assertNotNull(this.source.output()); - } + @StreamListener(Sink.INPUT) + public void handle(Vote vote) { + votingService.record(vote); + } } ---- -NOTE: In this case there is only one `Source` in the application context so there is no need to qualify it when it is -autowired. If there is ambiguity, e.g. if you are composing one application from some others, you can use the -`@Bindings` qualifier to inject a specific channel set. The `@Bindings` qualifier takes a parameter which is the class -that carries the `@EnableBinding` annotation (in this case the `TimerSource`). +The distinction between this approach and a Spring Integration `@ServiceActivator` becomes relevant if one considers an inbound `Message` with a `String` payload and a `contentType` header of `application/json`. +For `@StreamListener`, the `MessageConverter` mechanism will use the `contentType` header to parse the `String` into a `Vote` object. -==== Multiple Input or Output Channels +Just as with the other Spring Messaging methods, method arguments can be annotated with `@Payload`, `@Headers` and `@Header`. +For methods that return data, `@SendTo` must be used for specifying the output binding destination for data returned by the methods as follows: -A stream app can have multiple input or output channels defined as `@Input` and `@Output` methods in an interface. -Instead of just one channel named "input" or "output", you can add multiple `MessageChannel` methods annotated with -`@Input` or `@Output`, and their names will be converted to external destination names on the broker. It is common to -specify the channel names at runtime in order to have multiple applications communicate over well known destination -names. Channel names can be specified as properties that consist of the channel names prefixed with -`spring.cloud.stream.bindings` (e.g. `spring.cloud.stream.bindings.input` or `spring.cloud.stream.bindings.output`). -These properties can be specified though environment variables, the application YAML file, or any of the other -mechanisms supported by Spring Boot. - -For example, you can have two `MessageChannels` called "default" and "tap" in an application with -`spring.cloud.stream.bindings.default.destination=foo` and `spring.cloud.stream.bindings.tap.destination=bar`, -and the result is 2 bindings to an external broker with destinations called "foo" and "bar". - -==== Inter-app Communication - -While Spring Cloud Stream makes it easy for individual boot apps to connect to messaging systems, the typical scenario -for Spring Cloud Stream is the creation of multi-app pipelines, where microservice apps are sending data to each other. -This can be achieved by correlating the input and output destinations of adjacent apps, as in the following example. - -Supposing that the design calls for the `time-source` app to send data to the `log-sink` app, we will use a -common destination named `ticktock` for bindings within both apps. `time-source` will set -`spring.cloud.stream.bindings.output.destination=ticktock`, and `log-sink` will set -`spring.cloud.stream.bindings.input.destination=ticktock`. - -==== Consumer Group Support - -Spring Cloud Stream is a library focusing on building message-driven microservices, and more specifically stream -processing applications. In such scenarios, communication between different logical applications follows a -publish-subscribe pattern, with data being broadcast through a shared topic, but at the same time, it is important to -be able to scale up by creating multiple instances of a given application, which are in a competing consumer -relationship with each other. - -Spring Cloud Stream models this behavior through the concept of a consumer group, which is similar to the notion of -consumer groups in Kafka. Each consumer binding can specify a group name such as -`spring.cloud.stream.bindings.input.group=foo` (the actual name of the binding may vary). Each consumer group bound to -a given destination will receive a copy of the published data, but within the group, only one application will receive -each specific message. - -If no consumer group is specified for a given binding, then the binding is treated as if belonging to an anonymous, -independent, single-member consumer group. Otherwise said, if no consumer group is specified for a binding, it will be -in a publish-subscribe relationship with any other consumer groups. - -In general, it is preferable to always specify a consumer group when binding an application to a given destination. -When scaling up a Spring Cloud Stream application, a consumer group must be specified for each of its input bindings, -in order to prevent its instances from receiving duplicate messages (unless that behavior is desired, which is a less -common use case). - -NOTE: This feature has been introduced since version 1.0.0.M4. - -==== Instance Index and Instance Count - -When scaling up Spring Cloud Stream applications, each instance can receive information about how many other instances -of the same application exist and what its own instance index is. This is done through the -`spring.cloud.stream.instanceCount` and `spring.cloud.stream.instanceIndex` properties. For example, if there are 3 -instances of the HDFS sink application, all three will have `spring.cloud.stream.instanceCount` set to 3, and the -applications will have `spring.cloud.stream.instanceIndex` set to 0, 1 and 2, respectively. When Spring Cloud Stream -applications are deployed via Spring Cloud Data Flow, these properties are configured automatically, but when Spring -Cloud Stream applications are launched independently, these properties must be set correctly. By default -`spring.cloud.stream.instanceCount` is 1, and `spring.cloud.stream.instanceIndex` is 0. - -Setting up the two properties correctly on scale up scenarios is important for addressing partitioning behavior in -general (see below), and they are always required by certain types of binders (e.g. the Kafka binder) in order to -ensure that data is split correctly across multiple consumer instances. - -==== Advanced Binding Properties - -The input and output destination names are the primary properties to set in order to have Spring Cloud Stream -applications communicate with each other as their channels are bound to an external message broker automatically. -However, there are a number of scenarios where it is required to configure other attributes besides the destination -name. This is done using the following naming scheme: -`spring.cloud.stream.bindings..=`. The `destination` attribute is one such -example: `spring.cloud.stream.bindings.input.destination=foo`. A shorthand equivalent can be used as follows: -`spring.cloud.stream.bindings.input=foo`, but that shorthand can only be used only when there are no other attributes -to set on the binding. In other words, -`spring.cloud.stream.bindings.input.destination=foo`,`spring.cloud.stream.bindings.input.partitioned=true` is a valid -setup, whereas `spring.cloud.stream.bindings.input=foo`,`spring.cloud.stream.bindings.input.partitioned=true` is not. - -===== Partitioning - -Spring Cloud Stream provides support for partitioning data between multiple instances of a given application. In a -partitioned scenario, one or more producer apps will send data to one or more consumer apps, ensuring that data with -common characteristics is processed by the same consumer instance. The physical communication medium (i.e. the broker -topic or queue) is viewed as structured into multiple partitions. Regardless of whether the broker type is naturally -partitioned (e.g. Kafka) or not (e.g. Rabbit), Spring Cloud Stream provides a common abstraction for implementing -partitioned processing use cases in a uniform fashion. - -Setting up a partitioned processing scenario requires configuring both the data producing and the data consuming end. - -====== Configuring Output Bindings for Partitioning - -An output binding is configured to send partitioned data, by setting one and only one of its `partitionKeyExpression` -or `partitionKeyExtractorClass` properties, as well as its `partitionCount` property. For example, setting -`spring.cloud.stream.bindings.output.partitionKeyExpression=payload.id`,`spring.cloud.stream.bindings.output.partitionCount=5` -is a valid and typical configuration. - -Based on this configuration, the data will be sent to the target partition using the following logic. A partition key's -value is calculated for each message sent to a partitioned output channel based on the `partitionKeyExpression`. The -`partitionKeyExpression` is a SpEL expression that is evaluated against the outbound message for extracting the -partitioning key. If a SpEL expression is not sufficient for your needs, you can instead calculate the partition key -value by setting the property `partitionKeyExtractorClass`. This class must implement the interface -`org.springframework.cloud.stream.binder.PartitionKeyExtractorStrategy`. While, in general, the SpEL expression should -suffice, more complex cases may use the custom implementation strategy. - -Once the message key is calculated, the partition selection process will determine the target partition as a value -between `0` and `partitionCount - 1`. The default calculation, applicable in most scenarios is based on the formula -`key.hashCode() % partitionCount`. This can be customized on the binding, either by setting a SpEL expression to be -evaluated against the key via the `partitionSelectorExpression` property, or by setting a -`org.springframework.cloud.stream.binder.PartitionSelectorStrategy` implementation via the `partitionSelectorClass` -property. - -Additional properties can be configured for more advanced scenarios, as described in the following section. - -====== Configuring Input Bindings for Partitioning - -An input binding is configured to receive partitioned data by setting its `partitioned` property, as well as the -instance index and instance count properties on the app itself, as follows: -`spring.cloud.stream.bindings.input.partitioned=true`,`spring.cloud.stream.instanceIndex=3`,`spring.cloud.stream.instanceCount=5`. -The instance count value represents the total number of app instances between which the data needs to be partitioned, -whereas instance index must be a unique value across the multiple instances, between `0` and `instanceCount - 1`. The -instance index helps each app instance to identify the unique partition (or in the case of Kafka, the partition set) -from which it receives data. It is important that both values are set correctly in order to ensure that all the data is -consumed, and that the app instances receive mutually exclusive datasets. - -While setting up multiple instances for partitioned data processing may be complex in the standalone case, Spring Cloud -Data Flow can simplify the process significantly, by populating both the input and output values correctly, as well as -relying on the runtime infrastructure to provide information about the instance index and instance count. - -=== Binder Selection - -Spring Cloud Stream relies on implementations of the Binder SPI to perform the task of connecting channels to message -brokers. Each Binder implementation typically connects to one type of messaging system. Spring Cloud Stream provides -out of the box binders for Kafka, RabbitMQ and Redis. - -====== Classpath Detection - -By default, Spring Cloud Stream relies on Spring Boot's auto-configuration to configure the binding process. If a -single binder implementation is found on the classpath, Spring Cloud Stream will use it automatically. So, for example, -a Spring Cloud Stream project that aims to bind only to RabbitMQ can simply add the following dependency: - -[source,xml] +[source,java] ---- - - org.springframework.cloud - spring-cloud-stream-binder-rabbit - +@EnableBinding(Processor.class) +public class TransformProcessor { + + @Autowired + VotingService votingService; + + @StreamListener(Processor.INPUT) + @SendTo(Processor.OUTPUT) + public VoteResult handle(Vote vote) { + return votingService.record(vote); + } +} ---- -====== Multiple Binders on the Classpath - -When multiple binders are present on the classpath, the application must indicate which binder is to be used for each -channel binding. Each binder configuration contains a `META-INF/spring.binders`, which is a simple properties file: - -[source] ----- -rabbit:\ -org.springframework.cloud.stream.binder.rabbit.config.RabbitServiceAutoConfiguration ----- - -Similar files exist for the other binder implementations (i.e. Kafka and Redis), and it is expected that custom binder -implementations will provide them, too. The key represents an identifying name for the binder implementation, whereas -the value is a comma-separated list of configuration classes that contain one and only one bean definition of the type -`org.springframework.cloud.stream.binder.Binder`. - -Selecting the binder can be done globally by either using the `spring.cloud.stream.defaultBinder` property, e.g. -`spring.cloud.stream.defaultBinder=rabbit`, or by individually configuring them on each channel binding. - -For instance, a processor app that reads from Kafka and writes to Rabbit can specify the following configuration: -`spring.cloud.stream.bindings.input.binder=kafka`,`spring.cloud.stream.bindings.output.binder=rabbit`. - -====== Connecting to Multiple Systems - -By default, binders share the Spring Boot auto-configuration of the application and create one instance of each binder -found on the classpath. In scenarios where an application should connect to more than one broker of the same type, -Spring Cloud Stream allows you to specify multiple binder configurations, with different environment settings. Please -note that turning on explicit binder configuration will disable the default binder configuration process altogether, so -all the binders in use must be included in the configuration. - -For example, this is the typical configuration for a processor that connects to two RabbitMQ broker instances: - -[source,yml] ----- -spring: - cloud: - stream: - bindings: - input: - destination: foo - binder: rabbit1 - output: - destination: bar - binder: rabbit2 - binders: - rabbit1: - type: rabbit - environment: - spring: - rabbitmq: - host: - rabbit2: - type: rabbit - environment: - spring: - rabbitmq: - host: ----- - - - -=== Managed vs Standalone - -Code using the Spring Cloud Stream library can be deployed as a standalone application or be used as a Spring Cloud -Data Flow module. In standalone mode, your application will run happily as a service or in any PaaS (Cloud Foundry, -Heroku, Azure, etc.). Spring Cloud Data Flow helps orchestrate the communication between instances, so the aspects of -configuration that deal with application interconnection will be configured transparently. - -==== Fat JAR - -You can run in standalone mode from your IDE for testing. To run in production you can create an executable (or "fat") -JAR using the standard Spring Boot tooling provided for Maven or Gradle. - -==== Health Indicator - -Spring Cloud Stream provides a health indicator for the binders, registered under the name of `binders`. It can be -enabled or disabled using the `management.health.binders.enabled` property. +NOTE: Content type headers can be either set by external applications in the case of Rabbit MQ, and they are supported as part of an extended internal protocol by Spring Cloud Stream for any type of transport (even the ones that do not support headers normally, like Kafka). === Binder SPI @@ -356,18 +483,405 @@ The RabbitMQ Binder implementation maps the destination to a `TopicExchange`, an will be bound to that `TopicExchange`. Each consumer instance that binds will trigger creation of a corresponding RabbitMQ `Consumer` instance for its group’s `Queue`. -==== Redis Binder +== Configuration options -.Redis Binder -image::redis-binder.png[width=300,scaledwidth="50%"] +Spring Cloud Stream supports general configuration options, as well as configuration for bindings and binders. Some binders allow additional properties for the bindings, supporting middleware-specific features. -NOTE: we recommend only using the Redis Binder for development +All configuration options can be provided to Spring Cloud Stream applications via all the mechanisms supported by Spring Boot: application arguments, environment variables, YML files etc. -The Redis Binder creates a `LIST` (which performs the role of a queue) for each consumer group. A consumer binding will -trigger `BRPOP` operations on its group's `LIST`. A producer binding will consult a `ZSET` to determine what groups -currently have active consumers, and then for each message being sent, an `LPUSH` operation will be executed on each of -those group's `LISTs`. +==== Spring Cloud Stream Properties -=== Samples +spring.cloud.stream.instanceCount:: + The number of deployed instances of the same application. Must be set for partitioning and with Kafka. Default value is `1`. +spring.cloud.stream.instanceIndex:: + The instance index of the application, a number from `0` to `instanceCount`-1. Used for partitioning and with Kafka. Automatically set in Cloud Foundry to match the instance index of the application. +spring.cloud.stream.dynamicDestinations:: + A list of destinations that can be bound dynamically, for example in a dynamic routing scenario. Only listed destinations can be bound if set. Default empty, allowing any destination to be bound. +spring.cloud.stream.defaultBinder:: + The default binder to use, if there are multiple binders configured. See <>. -For Spring Cloud Stream samples, please refer: https://github.com/spring-cloud/spring-cloud-stream-samples \ No newline at end of file +[[binding-properties]] +=== Binding properties + +Binding properties are supplied using the format `spring.cloud.stream.bindings..=`.`` represents the name of the channel being configured, e.g. `output` for a `Source`. +In what follows, we will indicate where the `spring.cloud.stream.bindings..` prefix is omitted and focus just on the property name, with the understanding that the prefix will be included at runtime. + +==== Properties for the use of Spring Cloud Stream + +The following binding properties are available for both input and output bindings and +must be prefixed with `spring.cloud.stream.bindings..` . + +destination:: + The target destination of channel on the bound middleware, e.g. Rabbit MQ exchange or + Kafka topic. If not set, the channel name will be used instead. +group:: + The consumer group of the channel. This property applies only to inbound bindings. + By default it is null, and indicates an anonymous consumer. See <>. +contentType:: + The content type of the channel. By default it is `null` and no type + coercion is performed. See <>. +binder:: + The binder used by this binding. By default, it is set to `null` and will + use the default binder, if one exists. See <> for details. + +==== Consumer properties + +The following binding properties are available for input bindings only and must be prefixed with `spring.cloud.stream.bindings..consumer`: + +concurrency:: + The concurrency of the inbound consumer. By default, set to `1`. +partitioned:: + Must be set to `true` if the consumer is receiving data from a partitioned + producer. By default it is set to `false`. +maxAttempts:: + The number of attempts of re-processing an inbound message. Default '3'. (Ignored by Kafka, currently). +backOffInitialInterval:: + The backoff initial interval on retry. Default `1000`.(Ignored by Kafka, currently). +backOffMaxInterval:: + The maximum backoff interval. Default `10000`.(Ignored by Kafka, currently). +backOffMultiplier:: + The backoff multiplier. Default `2.0`. + +==== Producer properties + +The following binding properties are available for output bindings only and must be prefixed with `spring.cloud.stream.bindings..producer`: + +partitionKeyExpression:: + A SpEL expression for partitioning outbound data. Default: `null`. If either this property is set or + `partitionKeyExtractorClass` is present, outbound data on this channel will be partitioned, + and `partitionCount` must be set to a value larger than 1 to be effective. + The two options are mutually exclusive. See <>. +partitionKeyExtractorClass:: + A `PartitionKeyExtractorStrategy` implementation. Default: `null`. If either this property is set or + `partitionKeyExpression` is present, outbound data on this channel will be partitioned, + and `partitionCount` must be set to a value larger than 1 to be effective. + The two options are mutually exclusive. See <>. +partitionSelectorClass:: + A `PartitionSelectorStrategy` implementation. Default `null`. Mutually exclusive with + `partitionSelectorExpression`. If none is set, the partition will be selected as the + `hashCode(key) % partitionCount`, where `key` is computed via either `partitionKeyExpression` + or `partitionKeyExtractorClass`. +partitionSelectorExpression:: + A SpEL expression for customizing partition selection. Default `null`. Mutually exclusive with + `partitionSelectorClass`. If none is set, the partition will be selected as the + `hashCode(key) % partitionCount`, where `key` is computed via either `partitionKeyExpression` + or `partitionKeyExtractorClass`. +partitionCount:: + The number of target partitions for the data, if partitioning is enabled. Default `1`. Must be + set to a value higher than `1` if the producer is partitioned. On Kafka it is interpreted as a + hint, and the larger of this and the partition count of the target topic will be used instead. +requiredGroups:: + A comma separated list of groups that the producer must ensure message delivery even if they + start after it has been created (e.g. by pre-creating durable queues in Rabbit MQ). + +[[binder-specific-configuration]] +== Binder-specific configuration + +This captures the binder, consumer and producer properties that are specific for several binder +implementations. + +=== Rabbit-specific settings + +==== Rabbit MQ Binder properties + +The binder supports the all Spring Boot properties for Rabbit MQ configuration. + +In addition to that, it also supports the following properties: + +spring.cloud.stream.binder.rabbit.addresses:: + A comma-separated list of RabbitMQ server addresses (used only for clustering and in conjunction with `nodes`). Default empty. +spring.cloud.stream.binder.rabbit.adminAddresses. Default empty. + A comma-separated list of RabbitMQ management plugin URLs - only used when nodes contains more than one entry. Entries in this list must correspond to the corresponding entry in addresses. Default empty. +spring.cloud.stream.binder.rabbit.nodes:: + A comma-separated list of RabbitMQ node names; when more than one entry, used to locate the server address where a queue is located. Entries in this list must correspond to the corresponding entry in addresses. Default empty. +spring.cloud.stream.binder.rabbit.username:: + The user name. Default `null`. +spring.cloud.stream.binder.rabbit.password:: + The password. Default `null`. +spring.cloud.stream.binder.rabbit.vhost:: + The virtual host. Default `null`. +spring.cloud.stream.binder.rabbit.useSSL:: + True if Rabbit MQ should use SSL. +spring.cloud.stream.binder.rabbit.sslPropertiesLocation:: + The location of the SSL properties file, when certificate exchange is used. +spring.cloud.stream.binder.rabbit.compressionLevel:: + Compression level for compressed bindings. Defaults to `1` (BEST_LEVEL). See `java.util.zip.Deflater`. + +==== Rabbit MQ Consumer Properties + +The following properties are available for Rabbit consumers only and +must be prefixed with `spring.cloud.stream.bindings..` + +acknowledgeMode:: + The acknowledge mode. Default `AUTO`. +autoBindDlq:: + Whether to automatically declare the DLQ and bind it to the binder DLX. Default `false`. +durableSubscription:: + Whether subscription should be durable. Only effective if `group` is also set. Default `true`. +maxConcurrency: + Default `1`. +prefetch: + Prefetch count. Default `1`. +prefix:: + A prefix to be added to the name of the `destination` and queues. Default "". +requeueRejected:: + Whether delivery failures should be requeued. Default `true`. +requestHeaderPatterns:: + The request headers to be transported. Default `[STANDARD_REQUEST_HEADERS,'*']`. +replyHeaderPatterns:: + The reply headers to be transported. Default `[STANDARD_REQUEST_HEADERS,'*']` +republishToDlq:: + By default, failed messages after retries are exhausted are rejected. If a dead-letter queue (DLQ) is configured, rabbitmq will route the failed message (unchanged) to the DLQ. Setting this property to true instructs the bus to republish failed messages to the DLQ, with additional headers, including the exception message and stack trace from the cause of the final failure. +transacted:: + Whether to use transacted channels. Default `false`. +txSize:: + The number of deliveries between acks. Default `1`. + +==== Rabbit Producer Properties + +The following properties are available for Rabbit producers only and +must be prefixed with `spring.cloud.stream.bindings..` + +autoBindDlq:: + Whether to automatically declare the DLQ and bind it to the binder DLX. Default `false`. +batchingEnabled:: + True to enable message batching by producers. Default `false`. +batchSize:: + The number of message to buffer when batching is enabled. Default `100`. +batchBufferLimit:: + Default `10000`. +batchTimeout:: + Default `5000`. +compress:: + Whether data should be compressed when sent. Default `false`. +deliveryMode:: + Delivery mode. Default `PERSISTENT`. +prefix:: + A prefix to be added to the name of the `destination` exchange. Default "". +requestHeaderPatterns:: + The request headers to be transported. Default `[STANDARD_REQUEST_HEADERS,'*']`. +replyHeaderPatterns:: + The reply headers to be transported. Default `[STANDARD_REQUEST_HEADERS,'*']` + +=== Kafka-specific settings + +==== Kafka binder properties + +spring.cloud.stream.binder.kafka.brokers:: + A list of brokers that the Kafka binder will connect to. Default `localhost`. +spring.cloud.stream.binder.kafka.defaultBrokerPort:: + The list of brokers allows to specify hosts with or without port information, i.e. `host1,host2:port2`. This configuration sets the default port when no port is configured in the broker list. Default `9092`. +spring.cloud.stream.binder.kafka.zkNodes:: + A list of Zookeeper nodes for the Kafka binder to connect to. Default `localhost`. +spring.cloud.stream.binder.kafka.defaultZkPort:: + The list of Zookeeper nodes allows to specify hosts with or without port information, i.e. `host1,host2:port2`. This configuration sets the default port when no port is configured in the node list. Default `2181`. +spring.cloud.stream.binder.kafka.headers:: + The list of custom that will be transported by the binder. Default empty. +spring.cloud.stream.binder.kafka.offsetUpdateTimeWindow:: + The frequency in milliseconds with which offsets are saved. Ignored if `0`. Default `10000`. +spring.cloud.stream.binder.kafka.offsetUpdateCount:: + The frequency in number of updates, which which consumed offsets are persisted. Ignored if `0`. Default `0`. Mutually exclusive with `offsetUpdateTimeWindow`. +spring.cloud.stream.binder.kafka.requiredAcks:: + The number of required acks on the broker. + +==== Kafka Consumer Properties + +The following properties are available for Kafka consumers only and +must be prefixed with `spring.cloud.stream.bindings..` + +autoCommitOffset:: + True to autocommit offsets when a message has been processed. If set to false, an `Acknowledgment` header will be available in the message headers for late acknowledgment. Default `true`. +mode:: + When set to `raw`, will disable header parsing on input. Useful when inbound data is coming from outside Spring Cloud Stream applications. Default `embeddedHeaders`. +resetOffsets:: + True to reset offsets on the consumer to the value provided by `startOffset`. Default `false`. +startOffset:: + The starting offset for new groups or when `resetOffsets` is `true`. Allowed values: `earliest`,`latest`. Defaults to null (equivalent to earliest). +minPartitionCount:: + The minimum number of partitions expected by the consumer if it creates the consumed topic automatically. Defaults to `1`. + +==== Kafka Producer Properties + +The following properties are available for Kafka producers only and +must be prefixed with `spring.cloud.stream.bindings..` + +bufferSize:: + This is an upper limit of how much data the Kafka Producer will attempt to batch before sending – specified in bytes. Default `16384`. +sync:: + Whether the producer is synchronous. Defaults to `false`. +batchTimeout:: + How long will the producer wait before sending in order to allow more messages to get accumulated in the same batch. Normally the producer will not wait at all, and simply send all the messages that accumulated while the previous send was in progress. A non-zero value may increase throughput at the expense of latency. Default `0`. +mode:: + When set to `raw`, disable header propagation on output. Useful when producing data for non-Spring Cloud Stream applications. Default `embeddedHeaders`. + +== Binder detection + +Spring Cloud Stream relies on implementations of the Binder SPI to perform the task of connecting channels to message +brokers. Each Binder implementation typically connects to one type of messaging system. Spring Cloud Stream provides +out of the box binders for Kafka, RabbitMQ and Redis. + +=== Classpath Detection + +By default, Spring Cloud Stream relies on Spring Boot's auto-configuration to configure the binding process. If a +single binder implementation is found on the classpath, Spring Cloud Stream will use it automatically. So, for example, +a Spring Cloud Stream project that aims to bind only to RabbitMQ can simply add the following dependency: + +[source,xml] +---- + + org.springframework.cloud + spring-cloud-stream-binder-rabbit + +---- + +[[multiple-binders]] +=== Multiple Binders on the Classpath + +When multiple binders are present on the classpath, the application must indicate which binder is to be used for each channel binding. Each binder configuration contains a `META-INF/spring.binders`, which is a simple properties file: + +[source] +---- +rabbit:\ +org.springframework.cloud.stream.binder.rabbit.config.RabbitServiceAutoConfiguration +---- + +Similar files exist for the other binder implementations (e.g. Kafka), and it is expected that custom binder +implementations will provide them, too. The key represents an identifying name for the binder implementation, whereas +the value is a comma-separated list of configuration classes that contain one and only one bean definition of the type +`org.springframework.cloud.stream.binder.Binder`. + +Selecting the binder can be done globally by either using the `spring.cloud.stream.defaultBinder` property, e.g. +`spring.cloud.stream.defaultBinder=rabbit`, or by individually configuring them on each channel binding. + +For instance, a processor app that reads from Kafka and writes to Rabbit can specify the following configuration: +`spring.cloud.stream.bindings.input.binder=kafka`,`spring.cloud.stream.bindings.output.binder=rabbit`. + +=== Connecting to Multiple Systems + +By default, binders share the Spring Boot auto-configuration of the application and create one instance of each binder +found on the classpath. In scenarios where an application should connect to more than one broker of the same type, +Spring Cloud Stream allows you to specify multiple binder configurations, with different environment settings. Please +note that turning on explicit binder configuration will disable the default binder configuration process altogether, so +all the binders in use must be included in the configuration. + +For example, this is the typical configuration for a processor that connects to two RabbitMQ broker instances: + +[source,yml] +---- +spring: + cloud: + stream: + bindings: + input: + destination: foo + binder: rabbit1 + output: + destination: bar + binder: rabbit2 + binders: + rabbit1: + type: rabbit + environment: + spring: + rabbitmq: + host: + rabbit2: + type: rabbit + environment: + spring: + rabbitmq: + host: +---- + +[[contenttypemanagement]] +== Content Type and Transformation + +Spring Cloud Stream allows to propagate information about the content type of the messages it produces by attaching by default a `contentType` header to outbound messages. +For middleware that does not directly support headers, Spring Cloud Stream provides its own mechanism of wrapping outbound messages in an envelope of its own, automatically. +For middleware that does support headers, Spring Cloud Stream applications may receive messages with a given content type from non-Spring Cloud Stream applications. + +Spring Cloud Stream can handle messages based on this information in two ways: + +* through its `contentType` settings on inbound and outbound channels; +* through its argument mapping done for `@StreamListener`-annotated methods. + +=== Type converting message channels + +=== @StreamListener and conversion + +== Inter-app Communication + +=== Connecting multiple application instances + +While Spring Cloud Stream makes it easy for individual boot apps to connect to messaging systems, the typical scenario for Spring Cloud Stream is the creation of multi-app pipelines, where microservice apps are sending data to each other. +This can be achieved by correlating the input and output destinations of adjacent apps, as in the following example. + +Supposing that the design calls for the `time-source` app to send data to the `log-sink` app, we will use a +common destination named `ticktock` for bindings within both apps. `time-source` will set +`spring.cloud.stream.bindings.output.destination=ticktock`, and `log-sink` will set +`spring.cloud.stream.bindings.input.destination=ticktock`. + +=== Instance Index and Instance Count + +When scaling up Spring Cloud Stream applications, each instance can receive information about how many other instances +of the same application exist and what its own instance index is. This is done through the +`spring.cloud.stream.instanceCount` and `spring.cloud.stream.instanceIndex` properties. For example, if there are 3 +instances of the HDFS sink application, all three will have `spring.cloud.stream.instanceCount` set to 3, and the +applications will have `spring.cloud.stream.instanceIndex` set to 0, 1 and 2, respectively. When Spring Cloud Stream +applications are deployed via Spring Cloud Data Flow, these properties are configured automatically, but when Spring +Cloud Stream applications are launched independently, these properties must be set correctly. By default +`spring.cloud.stream.instanceCount` is 1, and `spring.cloud.stream.instanceIndex` is 0. + +Setting up the two properties correctly on scale up scenarios is important for addressing partitioning behavior in +general (see below), and they are always required by certain types of binders (e.g. the Kafka binder) in order to +ensure that data is split correctly across multiple consumer instances. + +=== Partitioning + +==== Configuring Output Bindings for Partitioning + +An output binding is configured to send partitioned data, by setting one and only one of its `partitionKeyExpression` +or `partitionKeyExtractorClass` properties, as well as its `partitionCount` property. For example, setting +`spring.cloud.stream.bindings.output.partitionKeyExpression=payload.id`,`spring.cloud.stream.bindings.output.partitionCount=5` +is a valid and typical configuration. + +Based on this configuration, the data will be sent to the target partition using the following logic. A partition key's +value is calculated for each message sent to a partitioned output channel based on the `partitionKeyExpression`. The +`partitionKeyExpression` is a SpEL expression that is evaluated against the outbound message for extracting the +partitioning key. If a SpEL expression is not sufficient for your needs, you can instead calculate the partition key +value by setting the property `partitionKeyExtractorClass`. This class must implement the interface +`org.springframework.cloud.stream.binder.PartitionKeyExtractorStrategy`. While, in general, the SpEL expression should +suffice, more complex cases may use the custom implementation strategy. + +Once the message key is calculated, the partition selection process will determine the target partition as a value +between `0` and `partitionCount - 1`. The default calculation, applicable in most scenarios is based on the formula +`key.hashCode() % partitionCount`. This can be customized on the binding, either by setting a SpEL expression to be +evaluated against the key via the `partitionSelectorExpression` property, or by setting a +`org.springframework.cloud.stream.binder.PartitionSelectorStrategy` implementation via the `partitionSelectorClass` +property. + +Additional properties can be configured for more advanced scenarios, as described in the following section. + +===== Configuring Input Bindings for Partitioning + +An input binding is configured to receive partitioned data by setting its `partitioned` property, as well as the +instance index and instance count properties on the app itself, as follows: +`spring.cloud.stream.bindings.input.partitioned=true`,`spring.cloud.stream.instanceIndex=3`,`spring.cloud.stream.instanceCount=5`. +The instance count value represents the total number of app instances between which the data needs to be partitioned, +whereas instance index must be a unique value across the multiple instances, between `0` and `instanceCount - 1`. The +instance index helps each app instance to identify the unique partition (or in the case of Kafka, the partition set) +from which it receives data. It is important that both values are set correctly in order to ensure that all the data is +consumed, and that the app instances receive mutually exclusive datasets. + +While setting up multiple instances for partitioned data processing may be complex in the standalone case, Spring Cloud +Data Flow can simplify the process significantly, by populating both the input and output values correctly, as well as +relying on the runtime infrastructure to provide information about the instance index and instance count. + +== Health Indicator + +Spring Cloud Stream provides a health indicator for the binders, registered under the name of `binders`. It can be +enabled or disabled using the `management.health.binders.enabled` property. + +== Samples + +For Spring Cloud Stream samples, please refer: https://github.com/spring-cloud/spring-cloud-stream-samples