【发布时间】: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