perf(会员): 外部生成认领按账号折叠并发上限与在飞计数查询

- 新增 ExternalGenerationOwnerConcurrencyCache:单次 claim 事务内按账号只查一次 concurrent_job_limit,running 计数按账号缓存
- 认领 / lease 耗尽终结后失效该账号 running 缓存,判定结果与逐候选行现查一致
- 回收豁免与不限档不再触发 running 计数查询;补缓存失效单测
This commit is contained in:
2026-10-03 17:23:19 +08:00
parent 413b3645c2
commit 1d09d7e061
@@ -1,5 +1,6 @@
use crate::*;
use std::cmp::Ordering;
use std::collections::HashMap;
use std::ops::RangeFrom;
const EXTERNAL_GENERATION_STATUS_PENDING: &str = "pending";
@@ -885,14 +886,23 @@ fn claim_external_generation_jobs_tx(
});
let mut claimed = Vec::new();
let mut concurrency = ExternalGenerationOwnerConcurrencyCache::default();
for mut row in candidates.into_iter().take(limit) {
if external_generation_job_has_exhausted_attempts(&row) {
let owner_user_id = row.owner_user_id.clone();
finalize_external_generation_job_after_lease_exhaustion(ctx, row, claim_time)?;
// 终结可能把 running 行改成 failed,账号计数缓存随之失效。
concurrency.invalidate_running(&owner_user_id);
continue;
}
// 会员并发上限:只约束 `pending → running` 的认领。达上限的任务留在 `pending` 天然排队;
// 账号自己那条已过期的 `running` 回收不受限,否则超限账号的卡死任务永远无法回收。
if !external_generation_owner_has_concurrency_capacity(ctx, &row, claim_time) {
if !external_generation_owner_has_concurrency_capacity(
ctx,
&mut concurrency,
&row,
claim_time,
) {
continue;
}
let next_attempt = row.attempt.saturating_add(1);
@@ -921,33 +931,77 @@ fn claim_external_generation_jobs_tx(
Some(worker_id.clone()),
claim_time,
);
// 本次认领把一行推进到 running,账号计数缓存失效,后续候选行重新现查。
concurrency.invalidate_running(&row.owner_user_id);
claimed.push(map_external_generation_job_row(row));
}
Ok(claimed)
}
/// 单次 claim 事务内按账号缓存有效并发上限与 running 行数。
///
/// 上限在一次事务内是常量;running 计数在账号的任一行被认领或终结后立即失效、下次重查,
/// 因此判定结果与逐候选行现查一致,只是把同一账号的重复查询折叠为一次。
#[derive(Default)]
struct ExternalGenerationOwnerConcurrencyCache {
limits: HashMap<String, u32>,
running: HashMap<String, u32>,
}
impl ExternalGenerationOwnerConcurrencyCache {
fn limit(&mut self, ctx: &ReducerContext, owner_user_id: &str) -> u32 {
if let Some(limit) = self.limits.get(owner_user_id) {
return *limit;
}
let limit = crate::effective_profile_concurrent_job_limit(ctx, owner_user_id);
self.limits.insert(owner_user_id.to_string(), limit);
limit
}
fn running_count(&mut self, ctx: &ReducerContext, owner_user_id: &str) -> u32 {
if let Some(count) = self.running.get(owner_user_id) {
return *count;
}
let count = count_running_external_generation_jobs_for_owner(ctx, owner_user_id);
self.running.insert(owner_user_id.to_string(), count);
count
}
/// 账号任一 `running` 行发生状态变化(认领 / 终结)后调用,保证下次读取重新现查。
fn invalidate_running(&mut self, owner_user_id: &str) {
self.running.remove(owner_user_id);
}
fn allows(&mut self, ctx: &ReducerContext, owner_user_id: &str, is_recycling: bool) -> bool {
let limit = self.limit(ctx, owner_user_id);
if is_unlimited_concurrency(limit) {
return true;
}
// 回收豁免不需要计数:自身过期 running 已计入在飞数,放行不会推高它。
let running = if is_recycling {
0
} else {
self.running_count(ctx, owner_user_id)
};
external_generation_claim_within_concurrency_limit(is_recycling, running, limit)
}
}
/// 认领前的并发上限判定。
///
/// 计数只算 `status = running`(含 lease 过期待回收的 `expired_running`):`pending`(含延时重试)
/// 不占名额,失败退回 `pending` 会自然释放。`128` 哨兵由 [`is_unlimited_concurrency`] 解释。
fn external_generation_owner_has_concurrency_capacity(
ctx: &ReducerContext,
cache: &mut ExternalGenerationOwnerConcurrencyCache,
row: &ExternalGenerationJob,
now: Timestamp,
) -> bool {
let limit = crate::effective_profile_concurrent_job_limit(ctx, &row.owner_user_id);
if is_unlimited_concurrency(limit) {
return true;
}
// 回收豁免:候选本身已是该账号过期的 `running`,它已计入在飞数;放行回收不会推高在飞任务数。
let is_recycling = row.status == EXTERNAL_GENERATION_STATUS_RUNNING
&& is_external_generation_job_claimable(row, now);
external_generation_claim_within_concurrency_limit(
is_recycling,
count_running_external_generation_jobs_for_owner(ctx, &row.owner_user_id),
limit,
)
cache.allows(ctx, &row.owner_user_id, is_recycling)
}
/// 并发上限判定(纯函数):`running_count >= limit` 时只放行对账号自身过期 `running` 的回收。
@@ -3936,6 +3990,18 @@ mod tests {
));
}
#[test]
fn owner_concurrency_cache_only_invalidates_running_count() {
let mut cache = ExternalGenerationOwnerConcurrencyCache::default();
cache.running.insert("owner-a".to_string(), 2);
cache.limits.insert("owner-a".to_string(), 3);
// 认领 / 终结后 running 计数必须重查,否则会把上一行的变更带到下一候选行。
cache.invalidate_running("owner-a");
assert!(cache.running.get("owner-a").is_none());
// 上限在一次事务内是常量,不受 running 计数失效影响。
assert_eq!(cache.limits.get("owner-a"), Some(&3));
}
fn micros(value: i64) -> Timestamp {
Timestamp::from_micros_since_unix_epoch(value)
}