How to batch process functions with Scala Future's?

Viewed 3525

I have a single function I want to call 100 times but I want to do it in a batched way so that only 2 functions are being run at any one time. This is due to the fact that the function may place a high load on the internet connection, so its better to batch process the functions in groups of 2's.

This is my attempt to do it with Scala Futures but does not seem to work. Is there any standard way using Scala Futures to batch process a list of tasks?

  def futureString(s:String): String = {
    Thread.sleep(2000)// + (Math.random()*1000).toInt)
    println(s"Completed $s")
    "end:" + s
  }

 def processList(list: List[String], blockSize: Int) = {
    var futuresProcessing = Set[Future[String]]()
    async {
      val itemIterator = list.iterator
      while (itemIterator.hasNext) {
        val item = itemIterator.next()
        println("Item is " + item)

        if (futuresProcessing.size >= blockSize) {
          await {
            val completed = Future.firstCompletedOf(futuresProcessing.toSeq)
            println("Size : " + futuresProcessing.size)
            completed
          }
        }

        val f = future { futureString(item) }
        f.onComplete{ case Success(sss) => { futuresProcessing = futuresProcessing - f } }
        futuresProcessing = futuresProcessing + f
      }
    }
  }

  val list: List[String] = (1 to 200).map(n => "" + n).toList
  processList(list, 2)

What I want is that I can batch process any batch size and the futureString may finish at a random amount of time. So lets say the batch size is 10, then at first 10 items start, when one item finishes, then a new item should be added to the batch to process.

I'm starting to think I should use actors.

Update: After a good long sleep and awake I got it working, but I think it would be better done with Actors. Also I think there's some race condition problem with the following code and the use of the futuresProcessing Set.

  import scala.concurrent._
  import scala.concurrent.duration._
  import ExecutionContext.Implicits.global
  import scala.async.Async.{async, await}
  import scala.collection.parallel.mutable
  import scala.util.{Success, Try}
  import scala.concurrent.Await

  def futureString(s:String): Future[String] = {
    future {
    Thread.sleep(2000 + (Math.random()*1000).toInt)
    println(s"Completed $s")
    "end:" + s
    }
  }

  def processList(list: List[String], blockSize: Int) = {
    val futuresProcessing = mutable.ParSet[Future[String]]()
    async {
      val itemIterator = list.iterator
      while (itemIterator.hasNext) {
        val item = itemIterator.next()
        println("Item is " + item)

        if (futuresProcessing.size >= blockSize) {
          await {
            val completed = Future.firstCompletedOf(futuresProcessing.toList)
            println("Size : " + futuresProcessing.size)
            completed
          }
        }

        val f = futureString(item)
        futuresProcessing += f
        f.onComplete{ case Success(sss) => { futuresProcessing -= f } }
      }
    }
  }

val list: List[String] = (1 to 200).map(n => "" + n).toList
processList(list, 4)
2 Answers
Related