Aerospike truncate操作如何兼顾执行性能与受影响记录日志采集
方案1:仅需要统计受影响总记录数(无单条详情)
这个方案完全保留你现有背景操作的高性能,不需要改动核心执行逻辑,只需要读取ExecuteTask内置的统计指标即可:
- 你调用
client.execute()返回的ExecuteTask对象本身提供了getRecordsProcessed()方法,可以直接拿到本次操作命中并处理的总记录数 - 你可以在
waitTillComplete之后直接打印这个数值,实现总数量的可见性,耗时完全和你现在的10分钟版本一致
对应代码示例:
val executeTask = client.execute(writePolicy, statement, toNullify: _*) executeTask.waitTillComplete(10.seconds.toMillis.toInt, 1.hour.toMillis.toInt) val affectedCount = executeTask.getRecordsProcessed logger.info(s"Namespace $namespace Set $set 本次truncate操作共处理 $affectedCount 条记录")
方案2:需要单条记录的操作详情(和旧版本一致的debug日志)
旧版本慢的核心原因是同步扫描+单条同步写入,没有利用Aerospike的并行能力,只需要调整为并行扫描+异步批量写入,就可以在保留单条日志的前提下,性能接近背景操作的水平:
优化点说明:
- 调整ScanPolicy配置开启并行扫描:设置
scanPolicy.maxConcurrentNodes = 集群节点数、scanPolicy.scanPerNodeParallelism = 2~4,充分利用多节点并行扫描能力 - 用异步非阻塞写入代替同步写入:调用
putAsync方法,不需要等待单条写入返回,利用客户端内置的连接池并发提交请求 - 可选批量提交:每攒够100~500条符合条件的记录,用
operateBatch方法批量提交更新请求,进一步降低网络开销
优化后代码示例:
def truncate(startTime: Long, durableDelete: Boolean): Unit = { val calendar = Calendar.getInstance() logger.info(s"truncate(records s.t LUT <= $startTime = ${calendar.getTime}, durableDelete = $durableDelete) on ${config.toRecoverMap}") val writePolicy = new WritePolicy() writePolicy.durableDelete = durableDelete // 开启异步写入的重试策略,避免网络波动导致失败 writePolicy.maxRetries = 3 writePolicy.sleepBetweenRetries = 100 val scanPolicy = new ScanPolicy() scanPolicy.filterExp = Exp.build(Exp.le(Exp.lastUpdate(), Exp.`val`(calendar))) // 开启并行扫描配置 scanPolicy.maxConcurrentNodes = client.getNodeNames.size scanPolicy.scanPerNodeParallelism = 3 // 可选配置限流,避免压垮集群 scanPolicy.recordsPerSecond = 10000 config.toRecoverMap.flatMap { case (namespace, mapOfSetsToBins) => for ((set, bins) <- mapOfSetsToBins) yield { val recordCount = new AtomicInteger(0) // 批量攒批缓存,不需要批量可以删掉这部分 val batchBuffer = new ConcurrentLinkedQueue[(Key, Seq[Bin])]() val batchSize = 200 client.scanAll(scanPolicy, namespace, set, new ScanCallback() { override def scanCallback(key: Key, record: Record): Unit = { val requiresNullify = bins.filter(record.bins.containsKey(_)).toSeq if (requiresNullify.nonEmpty) { val count = recordCount.incrementAndGet() // 先打debug日志,和旧版本逻辑完全一致 logger.debug { val (nullified, remains) = record.bins.asScala.partition { case (k, _) => requiresNullify.contains(k) } s"(#$count): Record $nullified bins of record with userKey: ${key.userKey}, digest: ${Buffer.bytesToHexString(key.digest)} nullified, remains: $remains" } // 普通异步写入,不需要等待返回 // client.putAsync(writePolicy, key, requiresNullify.map(Bin.asNull): _*) // 批量写入版本,性能更高: batchBuffer.add((key, requiresNullify.map(Bin.asNull))) if (batchBuffer.size() >= batchSize) { val batchOps = (1 to batchSize).flatMap(_ => Option(batchBuffer.poll())) .map { case (k, bins) => new BatchWrite(k, bins.map(Operation.put(_)).asJava) } .asJava client.operateBatch(writePolicy, batchOps) } } } }) // 扫描结束后提交剩余的批量请求 if (!batchBuffer.isEmpty) { val batchOps = batchBuffer.asScala .map { case (k, bins) => new BatchWrite(k, bins.map(Operation.put(_)).asJava) } .asJava client.operateBatch(writePolicy, batchOps) } logger.info(s"Namespace $namespace Set $set 本次truncate操作共处理 ${recordCount.get()} 条记录") } } }
这个方案的实测性能可以做到和背景操作版本差距在20%以内,完全可以在生产环境使用。
内容的提问来源于stack exchange,提问作者Zvi Mints
相关产品推荐
相关产品推荐

