apache / apache/pekko

Feature Request: Add Flow#concatAllDeferred operator.

Open
#1,652 5 comments 0 reactions 0 assignees View on GitHub
t:stream
Dominant language
Scala
Stars
1.6k
Forks
211
Avg merge
1d 6h
Merged PRs (30d)
89

Description

Motivation:
The original issue is https://github.com/apache/pekko/pull/1623 and https://github.com/apache/pekko/discussions/1566 ,
which do help find some problems, but with how the current interpreter and `concatAllLazy` are implemented, we can not fix the problem.

refs: https://projectreactor.io/docs/core/release/api/reactor/core/publisher/Flux.html#concat-java.lang.Iterable-

So a new operator is needed.

Modification:
I would like to add a new operator `concatAllDeferred` to support this usage.
```scala
def concatAllDeferred[U >: Out](those: Graph[SourceShape[U], _]*): Repr[U] =
concatLazy(Source.lazySource(
() => Source(those).flatMapConcat(ConstantFun.scalaIdentityFunction)))
```

Result:
```scala
package org.apache.pekko.stream.scaladsl

import org.apache.pekko.actor.ActorSystem
import org.apache.pekko.util.ByteString

import scala.concurrent.Await

object PekkoQuickstart extends App {
private implicit val system: ActorSystem = ActorSystem()

val s = Source
.repeat(())
.map(_ => ByteString('a' * 400000))
.take(1000000)
.prefixAndTail(50000)
.flatMapConcat { case (prefix, tail) => Source(prefix).concatLazy(tail) }

val r = Source.empty
.concatAllDeferred(List.tabulate(30000)(_ => s): _*)
.runWith(Sink.ignore)

Await.result(r, scala.concurrent.duration.Duration.Inf)
println(r.value)

// Source
// .repeat(s)
// .take(30000)
// .flatMapConcat(x => x)
// .runWith(Sink.ignore)
// .onComplete(println(_))

// Source.empty
// .concatAllLazy(List.tabulate(30000)(_ => Source.lazySource(() => s)): _*)
// .runWith(Sink.ignore).onComplete(println(_))
}

```

runs without problem

image

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.