Skip to content
Open
Original file line number Diff line number Diff line change
Expand Up @@ -141,6 +141,35 @@ class PackageDao : SimpleMongoDao<TPackage>() {
this.updateFirst(query, update)
}

/**
* 向指定 package 的 historyVersion 追加一批版本名(去重语义,等价 `$addToSet.$each`)。
*
* 屏蔽 Spring Data MongoDB `AddToSetBuilder.each(Object...)` 的 vararg 陷阱:
* Kotlin 侧必须用 spread(`*`) 展开成独立元素,否则整个 Collection 会被当作单个数组元素塞进
* `$each`,生产端 historyVersion 会被写入错误结构。此处已在 DAO 内一次性处理,业务方无需感知。
*
* @return true 表示 mongo 侧发生了实际写入(modifiedCount > 0)
*/
fun appendHistoryVersions(packageId: String, names: Collection<String>): Boolean {
if (names.isEmpty()) return false
val query = Query(Criteria.where(ID).isEqualTo(packageId))
// toTypedArray() 产生独立数组副本,spread 展开为 vararg 独立元素,两个问题一次解决
val update = Update().addToSet(TPackage::historyVersion.name).each(*names.toTypedArray())
return this.updateFirst(query, update).modifiedCount > 0
}

/**
* 从指定 package 的 historyVersion 移除一批版本名(等价 `$pullAll`)。
*
* @return true 表示 mongo 侧发生了实际写入(modifiedCount > 0)
*/
fun removeHistoryVersions(packageId: String, names: Collection<String>): Boolean {
if (names.isEmpty()) return false
val query = Query(Criteria.where(ID).isEqualTo(packageId))
val update = Update().pullAll(TPackage::historyVersion.name, names.toTypedArray())
return this.updateFirst(query, update).modifiedCount > 0
}

companion object {
private const val HISTORY_VERSION = "historyVersion"
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,9 @@ import com.tencent.bkrepo.common.metadata.model.TPackageVersion
import com.tencent.bkrepo.common.metadata.util.PackageQueryHelper
import com.tencent.bkrepo.common.mongo.dao.simple.SimpleMongoDao
import com.tencent.bkrepo.repository.pojo.metadata.MetadataModel
import org.bson.Document
import org.springframework.context.annotation.Conditional
import org.springframework.data.domain.Sort
import org.springframework.data.mongodb.core.FindAndModifyOptions
import org.springframework.data.mongodb.core.query.Query
import org.springframework.data.mongodb.core.query.Update
Expand All @@ -56,6 +58,47 @@ class PackageVersionDao : SimpleMongoDao<TPackageVersion>() {
return this.find(PackageQueryHelper.versionListQuery(packageId))
}

/**
* 按 `_id` ASC 顺序分页返回指定 package 下的版本名,用 [lastId] 作为游标:
* - 首次传 `null`,从最小 `_id` 开始
* - 后续传上一批最后一个元素的 `_id`,实现「游标分页」,避免 `skip` 在大数据集下的性能退化
*
* 只投影 `_id` 与 `name` 两个字段,单批 payload 极小,适合 10w+ 版本包的分批修复场景。
* 返回 `Pair<id, name>` 列表,调用方可从最后一个元素取 id 作为下一批的游标。
*/
fun pageVersionNamesAfterId(packageId: String, lastId: String?, batchSize: Int): List<Pair<String, String>> {
val criteria = where(TPackageVersion::packageId).isEqualTo(packageId)
if (lastId != null) {
criteria.and("_id").gt(org.bson.types.ObjectId(lastId))
}
val query = Query(criteria)
.with(Sort.by(Sort.Direction.ASC, "_id"))
.limit(batchSize)
query.fields().include("_id").include(TPackageVersion::name.name)
// 采用 Document 接收 projection 结果,避免反序列化为 TPackageVersion 时因缺少必填字段报错
return this.find(query, Document::class.java).map {
it.getObjectId("_id").toHexString() to it.getString(TPackageVersion::name.name)
}
}

/**
* 反查指定 name 集合中哪些在 `package_version` 里实际存在。
* 走 `(packageId, name)` 复合索引,一次查询即可,用于分批探测 `historyVersion` 中的脏数据:
* `batch - 返回值 = 该批的脏数据`。
*
* 只投影 `name` 字段,返回结果集大小上限为 `names.size`。
*/
fun findExistingNames(packageId: String, names: Collection<String>): Set<String> {
if (names.isEmpty()) return emptySet()
val criteria = where(TPackageVersion::packageId).isEqualTo(packageId)
.and(TPackageVersion::name.name).`in`(names)
val query = Query(criteria)
query.fields().include(TPackageVersion::name.name)
// 同 pageVersionNamesAfterId:仅取 name 字段,反序列化到 Document 避免实体必填字段初始化报错
return this.find(query, Document::class.java)
.mapTo(HashSet(names.size)) { it.getString(TPackageVersion::name.name) }
}

fun findByTag(packageId: String, tag: String): TPackageVersion? {
return this.findOne(PackageQueryHelper.versionQuery(packageId, tag = tag))
}
Expand Down
Original file line number Diff line number Diff line change
@@ -1,5 +1,7 @@
package com.tencent.bkrepo.common.metadata.service.packages

import com.tencent.bkrepo.repository.pojo.packages.PackageMetadataRepairResult

interface PackageRepairService {

/**
Expand All @@ -11,4 +13,22 @@ interface PackageRepairService {
* 修正包的版本数
*/
fun repairVersionCount()

/**
* 按范围修复 Package 元数据字段(latest、historyVersion)。
*
* 以 package_version 集合为权威数据源:
* 1. latest 重算为 ordinal DESC 排序的第一个版本;
* 2. historyVersion 全量覆盖为当前所有版本名的集合。
*
* @param projectId 项目 ID,必填
* @param repoName 仓库名,必填
* @param packageKey 包唯一标识;为空时修复该仓库下所有 package
* @return 修复结果统计
*/
fun repairPackageMetadata(
projectId: String,
repoName: String,
packageKey: String? = null
): PackageMetadataRepairResult
}
Original file line number Diff line number Diff line change
Expand Up @@ -2,21 +2,23 @@ package com.tencent.bkrepo.common.metadata.service.packages.impl

import com.tencent.bkrepo.common.api.pojo.Page
import com.tencent.bkrepo.common.api.util.HumanReadable
import com.tencent.bkrepo.common.artifact.exception.PackageNotFoundException
import com.tencent.bkrepo.common.metadata.condition.SyncCondition
import com.tencent.bkrepo.common.mongo.dao.util.Pages
import com.tencent.bkrepo.common.metadata.dao.packages.PackageDao
import com.tencent.bkrepo.common.metadata.dao.packages.PackageVersionDao
import com.tencent.bkrepo.common.metadata.model.TPackage
import com.tencent.bkrepo.repository.pojo.packages.VersionListOption
import com.tencent.bkrepo.common.metadata.service.packages.PackageRepairService
import com.tencent.bkrepo.common.metadata.service.packages.PackageService
import com.tencent.bkrepo.common.metadata.util.PackageQueryHelper
import com.tencent.bkrepo.repository.pojo.packages.PackageMetadataRepairResult
import org.slf4j.Logger
import org.slf4j.LoggerFactory
import org.springframework.context.annotation.Conditional
import org.springframework.data.domain.Sort
import org.springframework.data.mongodb.core.query.Query
import org.springframework.data.mongodb.core.query.Update
import org.springframework.data.mongodb.core.query.isEqualTo
import org.springframework.data.mongodb.core.query.where
import org.springframework.stereotype.Service
import java.time.Duration
import java.time.LocalDateTime
Expand All @@ -25,7 +27,6 @@ import kotlin.system.measureNanoTime
@Service
@Conditional(SyncCondition::class)
class PackageRepairServiceImpl(
private val packageService: PackageService,
private val packageDao: PackageDao,
private val packageVersionDao: PackageVersionDao
) : PackageRepairService {
Expand Down Expand Up @@ -57,7 +58,7 @@ class PackageRepairServiceImpl(
val repoName = it.repoName
val key = it.key
try {
// 添加包管理
// 只修 historyVersion,不动 latest,保持全库入口的原始语义
doRepairPackageHistoryVersion(it)
logger.info("Success to repair history version for [$key] in repo [$projectId/$repoName].")
successCount += 1
Expand Down Expand Up @@ -120,6 +121,194 @@ class PackageRepairServiceImpl(
}
}

override fun repairPackageMetadata(
projectId: String,
repoName: String,
packageKey: String?
): PackageMetadataRepairResult {
logger.info(
"Start repair package metadata. projectId=[$projectId], repoName=[$repoName], " +
"packageKey=[${packageKey ?: "<ALL>"}]"
)
val startTime = LocalDateTime.now()
val failedItems = mutableListOf<PackageMetadataRepairResult.FailedItem>()
var total = 0
var updated = 0
var skipped = 0

val packageIterator: Sequence<TPackage> = if (!packageKey.isNullOrBlank()) {
val pkg = packageDao.findByKey(projectId, repoName, packageKey)
?: throw PackageNotFoundException("$projectId/$repoName/$packageKey")
Comment thread
zzdjx marked this conversation as resolved.
sequenceOf(pkg)
} else {
iterateRepoPackages(projectId, repoName)
}

packageIterator.forEach { pkg ->
total += 1
try {
val result = doRepairPackageMetadata(pkg)
if (result) {
updated += 1
} else {
skipped += 1
}
} catch (e: Exception) {
logger.error(
"Failed to repair metadata for package [${pkg.key}] " +
"in repo [$projectId/$repoName]: ${e.message}",
e
)
failedItems.add(
PackageMetadataRepairResult.FailedItem(
packageKey = pkg.key,
reason = e.message
)
)
}
}

val duration = Duration.between(startTime, LocalDateTime.now()).seconds
logger.info(
"Finish repair package metadata. projectId=[$projectId], repoName=[$repoName], " +
"total=$total, updated=$updated, skipped=$skipped, failed=${failedItems.size}, " +
"elapse=${duration}s."
)
return PackageMetadataRepairResult(
total = total,
updated = updated,
skipped = skipped,
failed = failedItems.size,
failedItems = failedItems
)
}

/**
* 以 package_version 为权威源修复单个 package 的 latest 与 historyVersion。
*
* @return true 表示实际发生 DB 更新,false 表示元数据已一致
*/
private fun doRepairPackageMetadata(tPackage: TPackage): Boolean {
val packageId = tPackage.id ?: return false
val query = PackageQueryHelper.packageQuery(tPackage.projectId, tPackage.repoName, tPackage.key)
val latestMutated = repairLatestField(tPackage, packageId, query)
val historyMutated = doRepairPackageHistoryVersion(tPackage)
return latestMutated || historyMutated
}

/** 若 latest 与 [PackageVersionDao.findLatest](ordinal DESC)不一致,则 `$set` 覆盖。 */
private fun repairLatestField(tPackage: TPackage, packageId: String, query: Query): Boolean {
val actualLatest: String? = packageVersionDao.findLatest(packageId)?.name
if (tPackage.latest == actualLatest) return false
val r = packageDao.updateFirst(query, Update().set(TPackage::latest.name, actualLatest))
logger.info(
"Repair latest for [${tPackage.key}] in repo " +
"[${tPackage.projectId}/${tPackage.repoName}]: [${tPackage.latest}] -> [$actualLatest]"
)
return r.modifiedCount > 0L
}

/**
* 分批修复 historyVersion:补差 [addMissingVersionNames] + 删差 [removeOrphanVersionNames],
* 等价于以 package_version 为源全量覆盖。单批 payload 恒为 O([REPAIR_VERSION_BATCH_SIZE]),
* 与包版本总数解耦,可安全处理 10w+ 版本的大包。
*
* @return true 表示触发过 DB 更新
*/
private fun doRepairPackageHistoryVersion(tPackage: TPackage): Boolean {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

这个会破换原有方法的实现吧

val packageId = tPackage.id ?: return false
val originalHistory = tPackage.historyVersion

val (addMutated, addedCount) = addMissingVersionNames(packageId, originalHistory)
val (removeMutated, removedCount) = removeOrphanVersionNames(packageId, originalHistory)

if (addedCount > 0 || removedCount > 0) {
logger.info(
"Repair historyVersion for [${tPackage.key}] in repo " +
"[${tPackage.projectId}/${tPackage.repoName}]: " +
"originalSize=${originalHistory.size}, added=$addedCount, removed=$removedCount"
)
}
return addMutated || removeMutated
}

/**
* 补差:按 `_id` 游标分页读 package_version,将 [originalHistory] 中缺失的 name 分批 `$addToSet`。
* 用应用侧 [HashSet] 跨批去重仅为让计数精确,DB 侧 `$addToSet` 天然幂等。
*/
private fun addMissingVersionNames(
packageId: String,
originalHistory: Set<String>
): Pair<Boolean, Int> {
var mutated = false
var addedCount = 0
val addBuffer = ArrayList<String>(REPAIR_VERSION_BATCH_SIZE)
val addedNames = HashSet<String>()
var lastId: String? = null

while (true) {
val batch = packageVersionDao.pageVersionNamesAfterId(packageId, lastId, REPAIR_VERSION_BATCH_SIZE)
if (batch.isEmpty()) break
batch.forEach { (_, name) ->
if (name !in originalHistory && addedNames.add(name)) addBuffer.add(name)
}
lastId = batch.last().first
if (addBuffer.isNotEmpty()) {
// 传入快照拷贝而非复用 buffer,避免下游持有引用后被 clear 清空
if (packageDao.appendHistoryVersions(packageId, ArrayList(addBuffer))) mutated = true
addedCount += addBuffer.size
addBuffer.clear()
}
if (batch.size < REPAIR_VERSION_BATCH_SIZE) break
}
return mutated to addedCount
}

/**
* 删差:将 [originalHistory] 按 [REPAIR_VERSION_BATCH_SIZE] 分块反查存在性,
* 不存在的分批 `$pullAll`。仅清理修复前已存在的 name,与补差新增的天然不相交。
*/
private fun removeOrphanVersionNames(
packageId: String,
originalHistory: Set<String>
): Pair<Boolean, Int> {
var mutated = false
var removedCount = 0
originalHistory.chunked(REPAIR_VERSION_BATCH_SIZE).forEach { chunk ->
val existing = packageVersionDao.findExistingNames(packageId, chunk)
val orphans = chunk.filter { it !in existing }
if (orphans.isNotEmpty()) {
if (packageDao.removeHistoryVersions(packageId, orphans)) mutated = true
removedCount += orphans.size
}
}
return mutated to removedCount
}

/**
* 按仓库范围分页遍历 package,返回 Sequence 以支持惰性消费。
*/
private fun iterateRepoPackages(projectId: String, repoName: String): Sequence<TPackage> = sequence {
var page = 1
while (true) {
val pageResult = queryPackageByRepo(projectId, repoName, page)
if (pageResult.records.isEmpty()) break
yieldAll(pageResult.records)
if (pageResult.records.size < REPAIR_PAGE_SIZE) break
page += 1
}
}

private fun queryPackageByRepo(projectId: String, repoName: String, page: Int): Page<TPackage> {
val criteria = where(TPackage::projectId).isEqualTo(projectId)
.and(TPackage::repoName.name).isEqualTo(repoName)
val query = Query(criteria).with(Sort.by(Sort.Direction.ASC, TPackage::key.name))
val totalRecords = packageDao.count(query)
val pageRequest = Pages.ofRequest(page, REPAIR_PAGE_SIZE)
val records = packageDao.find(query.with(pageRequest))
return Pages.ofResponse(pageRequest, totalRecords, records)
}

private fun updateVersionCount(tPackage: TPackage): Boolean? {
val actualCount = packageVersionDao.countVersion(tPackage.id!!)
return if (actualCount == tPackage.versions) {
Expand All @@ -135,15 +324,6 @@ class PackageRepairServiceImpl(
}
}

private fun doRepairPackageHistoryVersion(tPackage: TPackage) {
with(tPackage) {
val allVersion = packageService.listAllVersion(projectId, repoName, key, VersionListOption())
.map { it.name }
historyVersion = historyVersion.toMutableSet().apply { addAll(allVersion) }
packageDao.save(this)
}
}

private fun queryPackage(page: Int): Page<TPackage> {
val query = Query().with(
Sort.by(Sort.Direction.DESC, TPackage::projectId.name, TPackage::repoName.name, TPackage::key.name)
Expand All @@ -156,5 +336,12 @@ class PackageRepairServiceImpl(

companion object {
private val logger: Logger = LoggerFactory.getLogger(PackageRepairServiceImpl::class.java)
private const val REPAIR_PAGE_SIZE = 500

/**
* 修复 historyVersion 的分批大小,同时用于分页读取 package_version 与分块反查脏数据。
* 取 1000:mongo `$addToSet.$each` / `$in` / `$pullAll` 在此量级下 payload 与性能最平衡。
*/
private const val REPAIR_VERSION_BATCH_SIZE = 1_000
}
}
Loading
Loading