# Can Flow be used for an uncertain producer?

**URL:** <https://discuss.kotlinlang.org/t/can-flow-be-used-for-an-uncertain-producer/18511>\
**Category:** Language Design\
**Created:** [July 21, 2020, 11:53pm UTC](https://discuss.kotlinlang.org/t/can-flow-be-used-for-an-uncertain-producer/18511 "2020-07-21T23:53:29Z")\
**Posts on this page:** 7\
**Page:** 1

<div class="post-metadata">

**Author:** ![socialguy](https://sea1.discourse-cdn.com/flex019/user_avatar/discuss.kotlinlang.org/socialguy/32/7748_2.png) [@socialguy](https://discuss.kotlinlang.org/u/socialguy)\
**Post date:** [July 21, 2020, 11:53pm UTC](https://discuss.kotlinlang.org/t/can-flow-be-used-for-an-uncertain-producer/18511/1 "2020-07-21T23:53:29Z")

</div>

I’ve a use case where a producer emits elements at random intervals. I want the consumer to wait for a certain time, then give up (timeout) and close the channel. Currently, I’m doing this using a buffered channel, where the consumer loops over the channel within a `withTimeout` block.

Having heard that flows are rad, I’m wondering if this is possible using a flow. The problem is, once the consumer gets hold of the flow reference (immediately), there’s no way for the producer to emit more values into it (like `channel.send`), so the consumer always times out receiving nothing. It’s almost as if I need a deferred emit that can keep appending to the flow (like a queue).

---

<div class="post-metadata">

**Author:** ![nickallendev](https://avatars.discourse-cdn.com/v4/letter/n/cdc98d/32.png) [@nickallendev](https://discuss.kotlinlang.org/u/nickallendev)\
**Post date:** [July 22, 2020, 12:48am UTC](https://discuss.kotlinlang.org/t/can-flow-be-used-for-an-uncertain-producer/18511/2 "2020-07-22T00:48:46Z")

</div>

Flows and channels can be converted back and forth. Sounds like you may be looking for [ReceiveChannel.consumeAsFlow](https://kotlin.github.io/kotlinx.coroutines/kotlinx-coroutines-core/kotlinx.coroutines.flow/consume-as-flow.html). There’s also [BroadcastChannel.asFlow()](https://kotlin.github.io/kotlinx.coroutines/kotlinx-coroutines-core/kotlinx.coroutines.flow/as-flow.html).

If you are interfacing with some sort of callback/listener API (you call `channel.send` from some listener/callback implementation), then I recommend [callbackFlow](https://kotlin.github.io/kotlinx.coroutines/kotlinx-coroutines-core/kotlinx.coroutines.flow/callback-flow.html).

---

<div class="post-metadata">

**Author:** ![socialguy](https://sea1.discourse-cdn.com/flex019/user_avatar/discuss.kotlinlang.org/socialguy/32/7748_2.png) [@socialguy](https://discuss.kotlinlang.org/u/socialguy)\
**Post date:** [July 22, 2020, 5:53am UTC](https://discuss.kotlinlang.org/t/can-flow-be-used-for-an-uncertain-producer/18511/3 "2020-07-22T05:53:40Z")

</div>

@nickallendev `callbackFlow` looked promising since I’m working with a gRPC client callback, but I need to output to two different channels, which doesn’t seem to be possible using the `callbackFlow`.

My use case is gRPC bidirectional streaming. The client receives server response, and may generate more messages for the server (request channel). It also processes the response it received, and puts it on another channel (response channel).

I’m open to flow options, but it seems like sticking with the channels is the simplest option for now.

---

<div class="post-metadata">

**Author:** ![dalewking](https://sea1.discourse-cdn.com/flex019/user_avatar/discuss.kotlinlang.org/dalewking/32/1496_2.png) [@dalewking](https://discuss.kotlinlang.org/u/dalewking)\
**Post date:** [July 22, 2020, 3:37pm UTC](https://discuss.kotlinlang.org/t/can-flow-be-used-for-an-uncertain-producer/18511/4 "2020-07-22T15:37:27Z")

</div>

There is an example of doing this with Flow right here: [https://kotlinlang.org/docs/reference/coroutines/flow.html#flow-cancellation-basics](https://kotlinlang.org/docs/reference/coroutines/flow.html#flow-cancellation-basics)

It isn’t that you cancel the flow you are cancelling the coroutine where the flow collection is happening.

---

<div class="post-metadata">

**Author:** ![socialguy](https://sea1.discourse-cdn.com/flex019/user_avatar/discuss.kotlinlang.org/socialguy/32/7748_2.png) [@socialguy](https://discuss.kotlinlang.org/u/socialguy)\
**Post date:** [July 28, 2020, 12:40am UTC](https://discuss.kotlinlang.org/t/can-flow-be-used-for-an-uncertain-producer/18511/5 "2020-07-28T00:40:43Z")

</div>

I’m not sure why you mentioned cancellation. I think what I need is shared flow, which seems to be a [WIP](https://github.com/Kotlin/kotlinx.coroutines/issues/2047). There’s also an issue with the server never actually calling `onComplete` or `onError`, only `onNext`. Based on a custom value received in `onNext`, I can call `onError` myself, but there’s no way for me know when the flow is complete.

---

<div class="post-metadata">

**Author:** ![dalewking](https://sea1.discourse-cdn.com/flex019/user_avatar/discuss.kotlinlang.org/dalewking/32/1496_2.png) [@dalewking](https://discuss.kotlinlang.org/u/dalewking)\
**Post date:** [July 28, 2020, 12:58am UTC](https://discuss.kotlinlang.org/t/can-flow-be-used-for-an-uncertain-producer/18511/6 "2020-07-28T00:58:11Z")

</div>

> [@socialguy](#):
>
> I’m not sure why you mentioned cancellation

You were talking about timeout and that was an example of a timeout.

---

<div class="post-metadata">

**Author:** ![nickallendev](https://avatars.discourse-cdn.com/v4/letter/n/cdc98d/32.png) [@nickallendev](https://discuss.kotlinlang.org/u/nickallendev)\
**Post date:** [July 28, 2020, 6:39am UTC](https://discuss.kotlinlang.org/t/can-flow-be-used-for-an-uncertain-producer/18511/7 "2020-07-28T06:39:27Z")

</div>

> [@socialguy](#):
>
> I think what I need is shared flow, which seems to be a [WIP](https://github.com/Kotlin/kotlinx.coroutines/issues/2047).

A naive implementation is not exactly difficult. This is not tested, just typed out in the forum, but basic concept should be clear.

```auto
fun <T> Flow<T>.mySharedFlowIn(scope: CoroutineScope) {
    val channel = BroadcastChannel(Channel.BUFFERED)
    scope.launch {
        try {
            collect { channel.send(it) }
            channel.close()
        } catch (ex: Exception) {
            channel.close(ex)
        }
    }
    return channel.asFlow()
}

```

> [@socialguy](#):
>
> Based on a custom value received in `onNext` , I can call `onError` myself, but there’s no way for me know when the flow is complete.

You know a flow is complete when `collect` finishes. `Flow<T>.onComplete` is also an option.
