Akka HTTP - max-open-requests and substreams?

Viewed 76

I'm writing an app using Scala 2.13 with Akka HTTP 10.2.4 and Akka Stream 2.6.15. I'm trying to query a web service in a parallel manner, like so:

package com.example

import akka.actor.typed.scaladsl.ActorContext
import akka.http.scaladsl.Http
import akka.http.scaladsl.client.RequestBuilding.Get
import akka.http.scaladsl.model.HttpResponse
import akka.http.scaladsl.unmarshalling.Unmarshal
import akka.stream.scaladsl.{Flow, JsonFraming, Sink, Source}
import spray.json.DefaultJsonProtocol
import spray.json.DefaultJsonProtocol.jsonFormat2

import scala.util.Try

case class ClientStockPortfolio(id: Long, symbol: String)
case class StockTicker(symbol: String, price: Double)

trait SprayFormat extends DefaultJsonProtocol {
  implicit val stockTickerFormat = jsonFormat2(StockTicker)
}

class StockTrader(context: ActorContext[_]) extends SprayFormat {


  implicit val system = context.system.classicSystem
  val httpPool = Http().superPool()[Seq[ClientStockPortfolio]]

  def collectPrices() = {
    val src = Source(Seq(
      ClientStockPortfolio(1, "GOOG"),
      ClientStockPortfolio(2, "AMZN"),
      ClientStockPortfolio(3, "MSFT")
      )
    )

    val graph = src
      .groupBy(8, _.id % 8)
      .via(createPost)
      .via(httpPool)
      .via(decodeTicker)
      .mergeSubstreamsWithParallelism(8)
      .to(
        Sink.fold(0.0) { (totalPrice, ticker) =>
          insertIntoDatabase(ticker)

          totalPrice + ticker.price
        }
      )

    graph.run()
  }

  def createPost = Flow[ClientStockPortfolio]
    .grouped(10)
    .map { port =>
      (
        Get(uri = s"http://wherever/?symbols=${port.map(_.symbol).mkString(",")}"),
        port
      )
    }

  def decodeTicker = Flow[(Try[HttpResponse], Seq[ClientStockPortfolio])]
    .flatMapConcat { x =>
      x._1.get.entity.dataBytes
        .via(JsonFraming.objectScanner(Int.MaxValue))
        .mapAsync(4)(bytes => Unmarshal(bytes).to[StockTicker])
        .mapConcat { ticker =>
          lookupPreviousPrices(ticker)
        }
    }

  def lookupPreviousPrices(ticker: StockTicker): List[StockTicker] = ???
  def insertIntoDatabase(ticker: StockTicker) = ???
}

I have two questions. First, will the groupBy call that splits the stream into substreams run them in parallel like I want? And second, when I call this code, I run into the max-open-requests error, since I haven't increased the setting from the default. But even if I am running in parallel, I'm only running 8 threads - how is the Http().superPool() getting backed up with 32 requests?

0 Answers
Related