在以下示例中,并行評估(列印)具有不同鑒別器 ( "a","b"和"c") 的專案:
package org.example
import cats.effect.std.Random
import cats.effect.{ExitCode, IO, IOApp, Temporal}
import cats.syntax.all._
import cats.{Applicative, Monad}
import fs2._
import scala.concurrent.duration._
object GitterQuestion extends IOApp {
override def run(args: List[String]): IO[ExitCode] =
Random.scalaUtilRandom[IO].flatMap { implicit random =>
val flat = Stream(
("a", 1),
("a", 2),
("a", 3),
("b", 1),
("b", 2),
("b", 3),
("c", 1),
("c", 2),
("c", 3)
).covary[IO]
val a = flat.filter(_._1 === "a").through(rndDelay)
val b = flat.filter(_._1 === "b").through(rndDelay)
val c = flat.filter(_._1 === "c").through(rndDelay)
val nested = Stream(a, b, c)
nested.parJoin(100).printlns.compile.drain.as(ExitCode.Success)
}
def rndDelay[F[_]: Monad: Random: Temporal, A]: Pipe[F, A, A] =
in =>
in.evalMap { v =>
(Random[F].nextDouble.map(_.seconds) >>= Temporal[F].sleep) >> Applicative[F].pure(v)
}
}
運行此程式的結果將類似于:
(c,1)
(a,1)
(c,2)
(a,2)
(c,3)
(b,1)
(a,3)
(b,2)
(b,3)
請注意,具有相同鑒別器的專案之間沒有重新排序 - 它們是按順序處理的。(a, 2)之前永遠不會列印(a, 1)。
在我的實際場景中,鑒別器值不是提前知道的,而且可能有很多,但我想有相同的行為,我該怎么做?
uj5u.com熱心網友回復:
我認為您需要為此推出自己的groupBy功能。我認為您必須Queue為每個鑒別器創建一個。然后對于每個Queue發出一個Stream從那個中拉出元素的內部Queue。
這是我想到的未經測驗且可能是幼稚的實作:
import cats.effect.std.Queue
val nested =
(flat.map(Some(_)) Stream(None))
.evalScan(Map.empty[String, Queue[IO, Option[(String, Int)]]] -> Option.empty[Stream[IO, (String, Int)]]){
case ((map, _), t @ Some((key, value))) =>
if (map.contains(key))
map(key).offer(t).as(map -> None)
else {
for {
q <- Queue.unbounded[IO, Option[(String, Int)]]
_ <- q.offer(t)
r = (map (key -> q)) -> Some(Stream.fromQueueNoneTerminated(q))
} yield r
}
case ((map, _), None) =>
// None means the flat stream is finished
map.values.toList.traverse(_.offer(None))
.as(Map.empty -> None)
}
.map(_._2).unNone
val parallelism: Int = ???
nested
.map(_.through(rndDelay))
// produce and consume in parallel in order to
// avoid deadlocks in case of bounded parJoin
.prefetchN(parallelism)
.parJoin(parallelism)
.printlns
.compile
.drain
.as(ExitCode.Success)
uj5u.com熱心網友回復:
我相信這broadcastThrough可以滿足您的需求。
(但一定要仔細檢查Scaladoc)
IO為了簡單起見,我直接使用,但它應該很容易適應抽象F[_]
def discriminateProcessing[A, B](stream: Stream[IO, A])(discriminators: List[A => Boolean])(pipe: Pipe[IO, A, B]): Stream[IO, B] = {
val allPipes: List[Pipe[IO, A, B]] = discriminators.map { p =>
s => s.filter(p).through(pipe)
}
stream.broadcastThrough(allPipes : _*)
}
這將像這樣使用:
val result = discriminateProcessing(stream = flat)(discriminators = List(
_._1 === "a",
_._1 === "b",
_._1 === "c"
)) { s =>
s.evalMap { v =>
random.nextDouble.map(_.seconds).flatMap(IO.sleep).as(v)
}
}
您可以在此處看到運行的代碼。
轉載請註明出處,本文鏈接:https://www.uj5u.com/qukuanlian/358577.html
