11.
Beginning With Coroutine Flow
Written by Filip Babić
Coroutines are excellent when it comes to bridging the synchronous and asynchronous worlds, returning values and communicating between threads. Usually, that’s what you want and need. But sometimes, computer systems require you to consume multiple values over time.
And there are two ways you can do this: using sequences and streams. However, there are certain limitations to both approaches. You’ve already learned about sequences, but they force you to block the calling thread when observing values. So see what streams have to offer and how they behave in code.
Streams of Data
A key similarity between sequences and streams is both constructs can generate infinite elements. Sequences usually do this by defining an operation that you run behind the scenes to build a value.
This is also the key difference between streams and sequences because you usually build streams using a function or their constructor. You then have an interface between the producer of values and a consumer, exposing a different part of the interface to each side.
Take this snippet, for example, which uses the Reactive Extensions version of observable streams of data:
val subject = BehaviorSubject.create<String>()
subject.subscribe(observer)
subject.onNext("one")
subject.onNext("two")
subject.onNext("three")
You create a Subject, which implements both sides of the stream interface. The provider can use functions such as offer, onNext, and send to fill the queue for the stream with values to consume. In this case, it uses onNext from Rx.
Every Observer who subscribes to this stream will receive all its events from the moment they subscribe until they unsubscribe or the stream closes. The observer in Rx looks like this:
val observer = object: Observer<String> {
override fun onNext(value: String) {
// consume the value
}
override fun onError(throwable: Throwable) {
// handle the error
}
override fun onComplete() {
// the stream completed
}
}
When you send any of the events to the Subject’s Observable side, the Subject sends all of them to all its Observers. It acts as a data relay from a central point to multiple observing nodes. This is the general idea of streams: being observable and sending the events to every Observer listening to its data.
But, depending on the implementation of streams, you might have a different setup. Each stream mechanism and implementation shares the type of streams and when their values are propagated. As such, there are hot and cold streams of data. Let’s consume them one at a time.
Hot Streams
Hot streams behave like TV channels or radio stations. They keep sending events and emitting their data even though no one may be listening or watching the show. It’s why they’re called hot. Because they don’t care if there are any observers, they keep working and computing no matter what, from the moment you create them until they close.
This is good when you want values computed in the background fast, preparing them for multiple observers you already have waiting. But if you’re going to add observers after the fact, you could lose the data a hot stream might emit between the computation and the observers starting to listen to events.
Additionally, if the producer of values is hot, it can keep producing values even though there are no consumers. This effectively wastes resources, and you have to close the stream manually if you stop using it.
If you were to use coroutines to build such hot streams, you’d use the Channel API. They’re hot by default and support coroutines, making them less leak-prone. But even if you used structured concurrency and coroutines within the Channels, you could leak some resources until the CoroutineScope cancels.
This is why the idea of having a cold stream is important.
Note: A good portion of the
Channels API is either experimental or being deprecated because the team initially had one plan for the API but built something else. Take care when using these APIs because they’re likely to change or be removed in the future.
Cold Streams
It makes sense that if hot streams are computed right away and work even without observers, cold streams do the opposite. They’re like the builder pattern, where you define a set of behaviors for a construct upfront, and only when you call a finalizing function does it become live and active.
Given that, cold streams are like a social event. You can prepare everything, think of each specific detail you have to fulfill and organize, but only when you’re certain that people are coming does the event happen. Following the analogy, cold streams won’t produce or send values until they have an active Observer to whom they can emit the events.
This is much better because if there are no observers, there’s no need to execute a potentially heavy operation to produce a value. But if there’s at least one observer, you compute the value and pass it down the stream to any amount of consumers.
It sounds too good to be true, and there’s a reason why hot and cold streams aren’t used everywhere and for every occasion. It’s because streams have a lot of internal limitations, and a good stream should support a lot of features to be versatile.
Let’s explore some of these constraints.
Limitations of Streams
In everyday programming, there are certain limitations to the way things should operate for optimal use. You don’t want to waste resources, freeze the UI, lose data, etc. Some of these concepts apply to streams, as well. These limitations revolve around the same problem: the speed of producing and consuming the values.
If your producer sends too many values and the consumer can’t process them quickly enough, you’re bound to lose some data. To effectively process the values, you have to apply backpressure. This is the technical term for eliminating the bottleneck in the producer-consumer pair.
When the producer queue fills up and the consumer can’t process the values fast enough, you become bottlenecked from the consumer side. But if the consumer eats the values too quickly and waits for more to be produced, you’re bottlenecked on the producer side. Either way, one side has to halt, or block, until the pair balances again.
Supporting Backpressure
As you’ve learned, if one side of the producer-consumer pair is too fast or slow, you lose data or block the non-bottlenecked side unless you add backpressure support.
You can achieve backpressure in different ways, such as using buffered underlying streams with a fixed capacity. This is the easiest solution but also the most error-prone because you can easily use a lot of computer memory or even overflow the buffer. This causes a bottleneck, and you lose data. You could have the capacity of unlimited, but then you risk overflowing the memory.
Another way is to build a synchronization mechanism, where you’d pause and resume threads as bottlenecks occur. But this could be even worse because you could be freezing threads for a long time, which wastes resources. This is why it’s important to avoid blocking threads when building streams with backpressure. Because of this design requirement, the Flow API is a fresh new take on streams.
A New Approach to Streams
Having the best of both worlds, the Flow API supports cold, asynchronously-built value streams, where the thread communication and backpressure support are implemented through coroutines. It’s the perfect combination.
Having coroutines allows for backpressure by design. If your producer is overflowing the consumer, you can suspend the producer until you free up the events queue you need to process. On the other hand, if your consumer is fast and you need to slow it down — introduce a delay or a debounce period, all you have to do is apply the same logic — suspend the consumer until it meets your conditions.
The other happy coincidence of coroutines is built-in context switching. By abstracting away threading and dispatching, through the use of CoroutineContexts, you can easily switch the consumption of events from one thread to another by passing in a different CoroutineContext from Dispatchers. And it’s performant because you don’t have to worry about thread allocation. That’s because coroutines use predefined thread pools.
This seems a bit too good to be true, right? It feels like the API will be quite complicated because it has to handle all those details a regular stream cannot intrinsically implement.
Well, that’s where the fun kicks in. Flow works based on only two interfaces: the Flow and the FlowCollector. For comparison, if you’re coming from a reactive-driven world, the Flow would be like an Observable, whereas the FlowCollector would be something similar to an Observer — a subscriber to events.
Let’s examine how they work.
Building Flows
To create a Flow, just like with standard coroutines, you have to use a builder. But first, open Main.kt in the starter project. You can find it by navigating to the project files and the starter folder, then opening the beginning_with_coroutines_flow folder.
Next, find main. It should be empty, but you’re about to add the code it needs to build a Flow. Add the following snippet to main so it doesn’t look empty. Conveniently enough, Flow’s builder function is called flow:
val flowOfStrings = flow {
for (number in 0..100) {
emit("Emitting: $number")
}
}
This snippet of code builds a Flow<String>, which calls emit 100 times, sending a String value to every observer that listens to the data. And to do that, you must call collect. Add the following snippet under the flow call:
GlobalScope.launch {
flowOfStrings.collect { value ->
println(value)
}
}
Thread.sleep(1000)
collect is a suspending function and needs to be called from a coroutine or another suspending function. From within, you have access to every single value you emit from within the Flow builder. In this case, you’re consuming each value by printing it out.
Build and run. You should see the following output:
Emitting: 0
Emitting: 1
Emitting: 2
Emitting: 3
....
Emitting: 99
Emitting: 100
As mentioned before, it’s 100 values printed one by one.
Now, to understand how Flows work from within, check the builder definition:
public fun <T> flow(@BuilderInference block: suspend FlowCollector<T>.() -> Unit): Flow<T>
= SafeFlow(block)
You create a Flow<T> with BuilderInference. This means that, like with producers and actors, you devise the generic type within the function constructor. Further, the lambda block is of the type FlowCollector<T>.() -> Unit, meaning the internal scope of the lambda will be a FlowCollector. This is great because you can create a Flow and emit values directly to the collector. You have the entire API connected in one place, making it simple and clean to use.
In one of the previous snippets, you collected the values from a Flow. But you can do much more with the Flow before consuming the data.
Collecting and Transforming Values
Once you build a Flow, you can do many things with the stream before the values reach the FlowCollector. Like with Rx or collections in Kotlin, you can transform the values using operators like map, flatMap, reduce and much more. Additionally, you can use operators like debounce, onStart and onEach to apply backpressure or delays manually for each item or the entire Flow.
Take the following snippet, for example:
GlobalScope.launch {
flowOfStrings
.map { it.split(" ") }
.map { it.last() }
.onEach { delay(100) }
.collect { value ->
println(value)
}
}
If you replace the previous way of consuming the Flow and run main again, you now see the values are mapped back to the actual numbers after the String is split. Furthermore, you print each value with a slight delay, ultimately suspending the Flow until the consumer is ready.
Note: All the operators above are marked with
suspend, so you have to call them from within a coroutine or another suspending function. This keeps the API uniform becauseFlows are built on coroutines.
Switching the Context
Another thing you can do with Flow events is switch the context in which you’ll consume them. To do that, you have to call flowOn(context: CoroutineContext) like this:
GlobalScope.launch {
flowOfStrings
.map { it.split(" ") }
.map { it.last() }
.flowOn(Dispatchers.IO)
.onEach { delay(100) }
.flowOn(Dispatchers.Default)
.collect { value ->
println(value)
}
}
In this snippet, you call flowOn twice: first after defining the mapping operations, and then the second time after delaying every item for 100 milliseconds. The real power of applying context switching is that you can do it as many times as you want for each operator you’re calling on the Flow. However, whenever you call flowOn, you’re applying the context switch only on the preceding operators, as the documentation states:
/**
* This operator retains a sequential nature of flow if changing the context does
* not call for changing the dispatcher. Otherwise, if changing dispatcher is required,
* it collects flow emissions in one coroutine that's run using a specified context
* and emits them from another coroutines with the original collector's context...
**/
This preserves the sequential nature of a Flow, and the context isn’t leaked into the downstream flow. The rest of the Flow operators and chained calls don’t know about the context switch, nor can they abuse the previous CoroutineContext.
Additionally, if the context elements are the same, the first operator takes precedence, like so:
flow.map { ... } // Will use Dispatchers.IO in the end
.flowOn(Dispatchers.IO) // The first operator's context takes precedence
.flowOn(Dispatchers.Default + customContext) // Contexts are merged with upstream
Ultimately, it’s important to know the final consumption of events can happen only in the original context. This means that no matter how many context switches you apply to the Flow, the last context will be the same as the original one.
So if you create a Flow on the main thread, you have to consume the events on it. You have to be careful about this. Otherwise, you’ll get an exception if you try to produce values in a different context than the one you’re consuming events in.
Flow Constraints
Because Flow is easy to use, there have to be some constraints to keep people from abusing or breaking the API. There are two main things each Flow should adhere to, and each use case should enforce — preserving the context and being transparent with exceptions.
Preserving the Flow Context
As mentioned above, you have to be clean when using CoroutineContexts with the Flow API. The producing and consuming contexts have to be the same. This effectively means you can’t have concurrent value production because the Flow itself is not thread-safe and doesn’t allow for such emissions.
So if you try to run the following snippet:
val flowOfStrings = flow {
for (number in 0..100) {
GlobalScope.launch {
emit("Emitting: $number")
}
}
}
GlobalScope.launch {
flowOfStrings.collect { println(it) }
}
You receive an exception saying you can’t change the Flow concurrently.
If you want coroutines to be synchronized and able to concurrently produce values in the Flow, use channelFlow instead and trySend or send to emit the values to the FlowCollector.
Changing the code to the following snippet will work:
val flowOfStrings = channelFlow {
for (number in 0..100) {
withContext(Dispatchers.IO) {
trySend("Emitting: $number")
}
}
}
GlobalScope.launch {
flowOfStrings.collect { println(it) }
}
If not, you should create the Flow values in a non-concurrent way and then use flowOn to switch the Flow to any CoroutineContext you want to avoid using channelFlow.
Additionally, the Flow‘s CoroutineContext can’t be bound to a Job. Therefore, you shouldn’t combine any Jobs with the context you’re trying to switch the Flow to. This is because the Flow shouldn’t be lifecycle-aware and able to be canceled, particularly because you can effectively mix multiple CoroutineContexts using flowOn and introducing a Job can only break things or make them unsafe.
Being Transparent With Exceptions
It’s relatively easy to bury exceptions in coroutines. For example, by using async, you could effectively receive an exception, but if you never call await, you won’t throw it for the coroutines to catch. Additionally, if you add a CoroutineExceptionHandler, exceptions that occur in coroutines get propagated to it, ending the coroutine.
This is why Flow exposes a convenient function that behaves like flowOn. You can use catch, providing a lambda which will catch any exception you produce in the stream and any of its previous operators. Examine the snippet below:
flowOfStrings
.map { it.split(" ") }
.map { it[1] }
.catch { it.printStackTrace() }
.flowOn(Dispatchers.Default)
.collect { println(it) }
Instead of mapping to it.last, you’re using indices. If you receive an empty string, this will cause an IndexOutOfBoundsException. But because you’re calling catch after map, you’ll catch an exception that occurs and print its stack trace. This way, you’ll be able to handle any exceptions from the original stream and the operators alike.
Change the way you build Flow, to this:
val flowOfStrings = flow {
emit("")
for (number in 0..100) {
emit("Emitting: $number")
}
}
You’ll now cause an exception to be thrown, but you’ll see the program doesn’t crash. This is because catch will stop the exception from throwing all the way up to cause your app to crash.
You’ll still get a stack trace from the exception. Add this line of code under collect and within the coroutine to be sure the program continues normally:
println("The code still works!")
Run the code again. You should now see an exception’s stack trace, and right after that, The code still works!. This means catch stopped the exception from one of the stream operators from breaking the entire program. And the rest of the coroutine still runs and works like a charm.
In case you’d want to continue emitting values if an exception occurs, you have access to the original FlowCollector within catch and it’s advised to simply call emitAll, with the fallback values or emit with a single value that notifies the user of a mistake or an error. Change the GlobalScope.launch code to the following:
flowOfStrings
.map { it.split(" ") }
.map { it[1] }
.catch {
it.printStackTrace()
// send the fallback value or values
emit("Fallback")
}
.flowOn(Dispatchers.Default)
.collect { println(it) }
println("The code still works!")
You should now see the exception stack trace printed out, as well as Fallback and The code still works!. This shows you can catch exceptions in streams, handle them correctly and continue the stream with some fallback values. You won’t break the outer coroutine, even if an exception occurs and is caught with catch!
Key Points
-
Sometimes, you need to build more than one value asynchronously, which is usually done with sequences or streams.
-
Sequences are lazy and cold, but blocking when you need to consume events. It’s better to use and suspend coroutines instead.
-
If you build streams using
Channels, you have coroutine support and suspendability, but they’re hot by default. -
Being cold means the data isn’t computed until you start observing. As opposed to being cold, being hot means the data is computed right away, with or without any observers.
-
As such, streams have two sides: the producer, or observable construct, and a consumer, or the observer construct.
-
The main limitations of streams are they’re blocking and use backpressure.
-
Blocking happens when a stream needs to produce or consume events.
-
Backpressure is when a stream produces or consumes events too quickly and one side must be slowed to balance the stream.
-
Backpressure is usually done through blocking the thread of a producer or a consumer.
-
A good stream avoids blocking, supports context switching while still allowing for backpressure.
-
The
FlowAPI is built upon coroutines, allowing for suspending. -
Because you can suspend a consumer or a producer, you get intrinsic backpressure support.
-
Additionally, you’re avoiding blocking, which is what a good stream should do.
-
To create a
Flow, simply callflowand provide a way to emit values to theFlowCollectors that decide to process the values. -
To attach a
FlowCollectorto aFlow, you have to callcollect, with a lambda in which you’ll consume each value. -
collectis a suspending function, so it has to be within a coroutine or another suspending function. -
You have access to a
FlowCollectorfrom withincollect, so you can emit values. -
Flows can be transformed and mutated by various operators likemap,flatMapConcat. -
You can apply manual backpressure using
debounce,onEachandonStart. -
Switching the context of a
Flowallows you to change the threads in which you consume each piece of data or perform each operator. -
To switch context, call
flowOn(context)after the operators you wish to switch the context of. -
The
Flowcollects the values always in the context of theCoroutineScopeit’s located in. So if you callcollecton the main thread, you also consume the values there. -
Flows don’t allow you to produce values concurrently. If you try to do that, an exception occurs. -
If you need to produce values from multiple threads, you can use
channelFlow. -
It’s better to use
flowOnto switch contexts of theFlowthan bury them in coroutines. -
Flows should be transparent when it comes to exceptions. -
To handle exceptions with
Flows, use thecatchoperator. -
catchwill intercept any uncaught exceptions from all the operators you called beforecatchitself.
Where to Go From Here?
This chapter introduced you to the fundamentals of the Kotlin Coroutines Flow API. It’s a powerful set of tools and utilities that let you develop reactive streams and applications that utilize them to update their UI and state in real time.
In the next chapter, you’ll dive deeper into two specific types of Flows - the SharedFlow and StateFlow. They live at the core part of any stateful streams of data, where sharing data or accessing the last cached value is a common practice.
So keep on reading to understand the full potential of the Flow API and its helper functions and types.