【问题标题】:Scalaz-stream chunking UP to NScalaz 流分块 UP 到 N
【发布时间】:2015-08-28 04:40:59
【问题描述】:

给定一个这样的队列:

val queue: Queue[Int] = async.boundedQueue[Int](1000)

我想从 UP 到 100 的块中取出这个队列并将其流式传输到下游 Sink。

queue.dequeue.chunk(100).to(downstreamConsumer) 

有点工作,但如果我说 101 条消息,它不会清空队列。将剩下 1 条消息,除非再推入 99 条消息。我想尽可能多地从队列中取出最多 100 条消息,以我的下游进程可以处理的速度尽可能快。

是否有现有的组合器可用?

【问题讨论】:

    标签: scala scalaz scalaz-stream


    【解决方案1】:

    为此,您可能需要在从队列中出列时监控队列的大小。然后,如果大小达到 0,您将不再等待任何元素。事实上,您可以根据队列的大小来实现elastic 的批处理大小。 IE。 :

    val q = async.unboundedQueue[String]
    
    val deq:Process[Task,(String,Int)] = q.dequeue zip q.size
    val elasticChunk: Process1[(String,Int), Vector[String]] = ???
    val downstreamConsumer : Sink[Task,Vector[String]] = ???
    
    deq.pipe(elasticChunk) to downstreamConsumer
    

    【讨论】:

    • 你将如何实现 elasticChunk?
    • 我实际上使用方便的 q.dequeueBatch 方法解决了这个问题。不知道它存在。
    【解决方案2】:

    我实际上以不同于我预期的方式解决了这个问题。

    scalaz-stream 队列现在包含一个 dequeueBatch 方法,该方法允许将队列中的所有值(最多 N 个或块)出列。

    https://github.com/scalaz/scalaz-stream/issues/338

    【讨论】:

      猜你喜欢
      • 2015-02-15
      • 2021-08-06
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2017-09-23
      • 1970-01-01
      • 2020-05-23
      相关资源
      最近更新 更多