You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.06 17:15:03