【问题标题】:How do I limit write operations to 1k records/sec?如何将写操作限制为 1k 记录/秒?
【发布时间】:2018-11-20 23:00:05
【问题描述】:

目前,我能够以 500 的批量写入数据库。但是由于内存不足错误和子聚合器与数据库叶节点之间的延迟同步,有时我会遇到叶节点内存错误。唯一的解决方案是,如果我将写入操作限制为每秒 1k 条记录,我可以摆脱错误。

dataStream
  .map(line => readJsonFromString(line))
  .grouped(memsqlBatchSize)
  .foreach { recordSet =>
          val dbRecords = recordSet.map(m => (m, Events.transform(m)))
          dbRecords.map { record =>
            try {
              Events.setValues(eventInsert, record._2)
              eventInsert.addBatch
            } catch {
              case e: Exception =>
                logger.error(s"error adding batch: ${e.getMessage}")
                val error_event = Events.jm.writeValueAsString(mapAsJavaMap(record._1.asInstanceOf[Map[String, Object]]))
                logger.error(s"event: $error_event")
            }
          }

          // Bulk Commit Records
          try {
            eventInsert.executeBatch
          } catch {
            case e: java.sql.BatchUpdateException =>
              val updates = e.getUpdateCounts
              logger.error(s"failed commit: ${updates.toString}")
              updates.zipWithIndex.filter { case (v, i) => v == Statement.EXECUTE_FAILED }.foreach { case (v, i) =>
                val error = Events.jm.writeValueAsString(mapAsJavaMap(dbRecords(i)._1.asInstanceOf[Map[String, Object]]))
                logger.error(s"insert error: $error")
                logger.error(e.getMessage)
              }
          }
          finally {
            connection.commit
            eventInsert.clearBatch
            logger.debug(s"committed: ${dbRecords.length.toString}")
          }
        }

1k 记录的原因是,我尝试写入的一些数据可能包含大量 json 记录,如果批量大小为 500,则可能导致每秒 30k 记录。有什么方法可以确保无论记录数量多少,一次只写入 1000 条记录?

【问题讨论】:

  • 在我看来,给定的代码并没有告诉我们有关您提到的问题的任何信息。看起来dbRecords.map 应该是dbRecords.foreach
  • @SarveshKumarSingh 想出了一个解决方案。我想要么我引入 1.0 秒的延迟,要么使用 Thread.sleep(1000) 以便我可以限制它。可能是我之前没有理解这个要求。根据 memsql 的架构,所有加载到 memsql 中的操作都首先进入行存储,然后 memsql 将合并到列存储中。但是 Memsql 首先将数据收集到本地内存中,然后刷新到叶子。这导致每次都出现叶子错误。减少批量大小并引入延迟帮助我每秒写入 120,000 条记录。使用此基准执行测试。

标签: multithreading scala performance thread-synchronization singlestore


【解决方案1】:

我认为 Thead.sleep 不是处理这种情况的好主意。一般来说,我们不建议在 Scala 中这样做,而且我们也不想在任何情况下阻塞线程。

一个建议是使用任何流技术,例如 Akka.Stream、Monix.Observable。这些库之间有一些优点和缺点,我不想在上面花太多的篇幅。但是当消费者比生产者慢时,它们确实支持背压来控制生产速率。例如,在您的情况下,您的消费者正在编写数据库,而您的生产者可能正在读取一些 json 文件并进行一些聚合。

以下代码说明了这个想法,您需要根据需要进行修改:

val sourceJson = Source(dataStream.map(line => readJsonFromString(line)))
val sinkDB = Sink(Events.jm.writeValueAsString) // you will need to figure out how to generate the Sink
val flowThrottle = Flow[String]
  .throttle(1, 1.second, 1, ThrottleMode.shaping)

val runnable = sourceJson.via[flowThrottle].toMat(sinkDB)(Keep.right)
val result = runnable.run()

【讨论】:

  • 但我不明白您代码中的可运行语句。我的要求是,如果我的线程 = 4,batchSize = 500 并且我想以某种方式限制生产者端的写入限制为 1k 记录/秒,我如何用你的代码来保证这一点?
【解决方案2】:

代码块已经被一个线程调用并且有多个线程并行运行。我可以在这个 Scala 代码中使用Thread.sleep(1000)delay(1.0)。但是如果我使用delay(),它将使用一个可能必须在函数外部调用的承诺。看起来Thread.sleep() 是最好的选择以及1000 的批量大小。执行测试后,我可以毫无问题地对 120,000 条记录/线程/秒进行基准测试。

根据 memsql 的体系结构,所有加载到 memsql 中的内容都是先在行存储中完成到本地内存中,然后 memsql 将在最后离开时合并到列存储中。每次我推送更多数量的数据导致瓶颈时,都会导致叶子错误。减少批大小并引入 Thread.sleep() 帮助我每秒写入 120,000 条记录。使用此基准执行测试。

【讨论】:

    猜你喜欢
    • 2021-11-07
    • 1970-01-01
    • 1970-01-01
    • 2012-07-29
    • 1970-01-01
    • 2021-12-22
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多