【问题标题】:rejected from slick.util.AsyncExecutor on "large" Future.sequence在“大”Future.sequence 上被 slick.util.AsyncExecutor 拒绝
【发布时间】:2021-07-18 15:11:08
【问题描述】:

我花了一整天的时间试图弄清楚如何解决这个问题。

目的是将多个字符串序列插入到表的单个列中。

我有这样的方法:

case class Column(strings: Seq[String])

def insertColumns(columns: Seq[Column]) = for {
_ <- Future.sequence(columns.map(col => insert(col)))
} yield()

private def insert(column: Column) =
  db.run((stringTable ++= rows)) //slick batch insert

这在一定程度上是有效的。 我测试了 2100 列的序列(每列有 100 个字符串),它工作正常。 但是一旦我将列数增加到 3100+,我就有这个错误

Task slick.basic.BasicBackend$DatabaseDef$$anon$3@293ce053 rejected from slick.util.AsyncExecutor$$anon$1$$anon$2@3e423930[Running, pool size = 10, active threads = 10, queued tasks = 1000, completed tasks = 8160]

我在几个地方读到过这样做会有所帮助的地方

case class Column(strings: Seq[String])

val f = Future.sequence(columns.map(col => insert(col)))

def insertColumns(columns: Seq[Column]) = for {
_ <- f
} yield()

private def insert(column: Column) =
  db.run((stringTable ++= rows)) //slick batch insert

没有。

我尝试了insert内部的几种更改组合

Future.sequence(
rows.grouped(500).toSeq.map(group => db.run(DBIO.seq(stringTable ++= group)))
)
Source(rows).buffer(500, OverflowStrategy.backpressure)
  .via(
    Slick.flow(row => stringTable += row)
  )
  .log("nr-of-inserted-rows")
  .runWith(Sink.ignore)
Source(rows)
.runWith(Slick.sink(1, row => stringTable += row))

我试过了:

  • 不要在我的配置中使用reWriteBatchedInserts=true
  • (dataColumnStringsTable ++= rows).transactionally 选项
  • 使用特定的执行上下文来启用单个线程:implicit val ec: ExecutionContext = ExecutionContext.fromExecutor(Executors.newFixedThreadPool(1)) 尝试顺序执行期货

除了修改订阅者以接收和阻止我的消息(字符串序列)并处理队列消息方面的背压之外,我没有任何其他想法。

我正在使用 slick(带有 alpakka-slick)3.3.3 / HikariCP 3.2.0 / Postgres 13.2

我的配置是这样的

slick {
  profile = "slick.jdbc.PostgresProfile$"
  db {
      connectionPool = "HikariCP"
      dataSourceClass = "slick.jdbc.DriverDataSource"
      properties = {
        driver = "org.postgresql.Driver"
        user = "postgres"
        password = "password"
        url = "jdbc:postgresql://"${slick.db.host}":5432/slick?reWriteBatchedInserts=true"
      }
      host = "localhost"
      numThreads = 10
      maxConnections = 100
      minConnections = 1
    }
}

感谢您的帮助。

【问题讨论】:

    标签: postgresql scala slick concurrent.futures slick-3.0


    【解决方案1】:

    您不应将Future.sequence 用于包含多个元素的集合。每个Future 都是在后台运行的计算。所以当你运行这个集合时,比如说,3000 columns:

    Future.sequence(columns.map(col => insert(col)))
    

    您一次有效地产生了 3000 个操作。结果,执行者可能会开始拒绝新任务。

    解决方案是使用 Akka Streams 处理输入集合。在您的情况下,这意味着从columns(而不是从rows)创建Source。这将确保执行器不会被过多的并行操作所淹没。我没用过alpakka-slick,但是看看docs,解决方案应该是这样的:

    Source(columns)
      .via(
        Slick.flow(column => stringTable ++= column.rows) 
      )
      // further processing here
    

    此外,如果“列”来自消息队列,您甚至可能不需要中间 Seq[Column]。您可能只需要定义从队列中读取的Column 中的Source,并使用 Slick 流对其进行处理。

    【讨论】:

    • 你说得对,我设法避免了一次产生 3000 次操作。我需要通过源进行批量插入,应该没问题。感谢您的提示!
    猜你喜欢
    • 2017-08-08
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-08-10
    • 2011-05-13
    • 2013-07-10
    • 2014-07-22
    相关资源
    最近更新 更多