feat(cleanup): 解耦 request_candidates 与 provider_api_keys 生命周期

- 移除 request_candidates.key_id 对 provider_api_keys 的外键约束(含迁移脚本)
- 删除 Key 时不再级联删除候选记录,改为独立按保留天数定时清理
- 新增 request_candidates_retention_days / request_candidates_cleanup_batch_size 配置项
- batch_delete_task 增加 lock_timeout 及超时自动降批重试机制
- cleanup_key_references 提取阶段化清理流程,移除 RequestCandidate 联动删除
- 前端 CleanupPolicySection 新增候选记录保留天数和清理批次配置

Closes #227

Co-authored-by: Entropy-Xu <entropy.xu@cloudhabitatsh.com>
This commit is contained in:
fawney19
2026-03-14 01:31:39 +08:00
parent bdfe4adc98
commit 776dd2f8ea
14 changed files with 516 additions and 70 deletions

View File

@@ -2481,12 +2481,7 @@ class RequestCandidate(Base):
nullable=True,
index=True,
)
key_id = Column(
String(36),
ForeignKey("provider_api_keys.id", ondelete="CASCADE"),
nullable=True,
index=True,
)
key_id = Column(String(36), nullable=True, index=True, comment="Provider Key ID 快照")
# 状态信息
status = Column(
@@ -2533,7 +2528,6 @@ class RequestCandidate(Base):
api_key = relationship("ApiKey")
provider = relationship("Provider")
endpoint = relationship("ProviderEndpoint")
key = relationship("ProviderAPIKey")
# ==================== 统计数据模型 ====================

View File

@@ -27,13 +27,7 @@ class CandidateRecorder:
if getattr(row, "provider", None) is not None:
provider_name = getattr(row.provider, "name", None)
key_name = None
auth_type = None
priority = None
if getattr(row, "key", None) is not None:
key_name = getattr(row.key, "name", None)
auth_type = getattr(row.key, "auth_type", None)
priority = getattr(row.key, "priority", None)
key_name = getattr(row, "api_key_name", None)
result.append(
CandidateKey(
@@ -44,8 +38,6 @@ class CandidateRecorder:
endpoint_id=str(row.endpoint_id) if row.endpoint_id else None,
key_id=str(row.key_id) if row.key_id else None,
key_name=str(key_name) if key_name else None,
auth_type=str(auth_type) if auth_type else None,
priority=int(priority) if priority is not None else None,
is_cached=bool(getattr(row, "is_cached", False)),
status=str(getattr(row, "status", "") or "pending"),
skip_reason=getattr(row, "skip_reason", None),

View File

@@ -16,6 +16,7 @@ from concurrent.futures import Future
import redis.asyncio as aioredis
from sqlalchemy import delete as sa_delete
from sqlalchemy import text
from sqlalchemy.orm import Session
from src.clients.redis_client import get_redis_client
from src.core.logger import logger
@@ -36,9 +37,15 @@ _CLEANUP_BATCH_SIZE = 50
# 单个批次的数据库 statement 超时(秒)
_BATCH_STATEMENT_TIMEOUT_S = 30
# 单个批次的数据库锁等待超时(秒)
_BATCH_LOCK_TIMEOUT_S = 5
# 整个任务的最大执行时间(秒)
_TASK_TIMEOUT_S = 600
# 发生超时/锁等待时降批到的最小 Key 数
_MIN_RETRY_BATCH_SIZE = 1
# Redis key 前缀
_REDIS_KEY_PREFIX = "batch_delete_task"
@@ -50,6 +57,30 @@ def _task_key(task_id: str) -> str:
return f"{_REDIS_KEY_PREFIX}:{task_id}"
def _apply_statement_timeouts(db: Session) -> None:
db.execute(text(f"SET LOCAL statement_timeout = '{_BATCH_STATEMENT_TIMEOUT_S * 1000}'"))
db.execute(text(f"SET LOCAL lock_timeout = '{_BATCH_LOCK_TIMEOUT_S * 1000}'"))
def _is_retryable_batch_error(exc: Exception) -> bool:
messages = [str(exc)]
orig = getattr(exc, "orig", None)
if orig is not None:
messages.append(str(orig))
text_blob = " ".join(messages).lower()
return any(
marker in text_blob
for marker in (
"querycanceled",
"statement timeout",
"canceling statement due to statement timeout",
"lock timeout",
"canceling statement due to lock timeout",
)
)
class BatchDeleteTaskInfo:
"""任务状态数据对象(从 Redis 反序列化)。"""
@@ -177,6 +208,115 @@ async def get_batch_delete_task(task_id: str) -> BatchDeleteTaskInfo | None:
return await _load_task(task_id)
def _delete_key_batch(
db: Session,
provider_id: str,
batch: list[str],
) -> int:
from src.models.database import ProviderAPIKey
phase = "cleanup_key_references"
def _set_phase(stage_name: str, _batch_size: int) -> None:
nonlocal phase
phase = stage_name
try:
_apply_statement_timeouts(db)
cleanup_key_references(
db,
batch,
batch_size=len(batch),
stage_callback=_set_phase,
)
phase = "provider_api_keys"
result = db.execute(
sa_delete(ProviderAPIKey).where(
ProviderAPIKey.provider_id == provider_id,
ProviderAPIKey.id.in_(batch),
)
)
rowcount = getattr(result, "rowcount", 0) or 0
db.commit()
return int(rowcount)
except Exception as exc:
setattr(exc, "_aether_batch_phase", phase)
raise
def _delete_key_batch_with_retry(
db: Session,
provider_id: str,
batch: list[str],
*,
batch_idx: int,
total_batches: int,
start_offset: int,
attempt: int = 1,
) -> int:
try:
return _delete_key_batch(db, provider_id, batch)
except Exception as exc:
try:
db.rollback()
except Exception:
pass
is_retryable = _is_retryable_batch_error(exc)
phase = getattr(exc, "_aether_batch_phase", "unknown")
can_split = len(batch) > _MIN_RETRY_BATCH_SIZE
if is_retryable and can_split:
split_at = max(len(batch) // 2, _MIN_RETRY_BATCH_SIZE)
left = batch[:split_at]
right = batch[split_at:]
logger.warning(
"[BATCH_DELETE] batch {}/{} retrying after timeout/lock (keys {}-{}, size={}, attempt={}, phase={}): split into {} + {}",
batch_idx,
total_batches,
start_offset,
start_offset + len(batch),
len(batch),
attempt,
phase,
len(left),
len(right),
)
deleted = _delete_key_batch_with_retry(
db,
provider_id,
left,
batch_idx=batch_idx,
total_batches=total_batches,
start_offset=start_offset,
attempt=attempt + 1,
)
if right:
deleted += _delete_key_batch_with_retry(
db,
provider_id,
right,
batch_idx=batch_idx,
total_batches=total_batches,
start_offset=start_offset + len(left),
attempt=attempt + 1,
)
return deleted
logger.warning(
"[BATCH_DELETE] batch {}/{} failed (keys {}-{} size={} attempt={} phase={} retryable={}): {}",
batch_idx,
total_batches,
start_offset,
start_offset + len(batch),
len(batch),
attempt,
phase,
is_retryable,
exc,
)
return 0
def _sync_delete(
provider_id: str,
key_ids: list[str],
@@ -188,7 +328,6 @@ def _sync_delete(
每个批次独立事务,单批失败跳过并继续。
"""
from src.database import create_session
from src.models.database import ProviderAPIKey
db = create_session()
try:
@@ -209,33 +348,14 @@ def _sync_delete(
batch = key_ids[i : i + _CLEANUP_BATCH_SIZE]
batch_idx += 1
try:
# 设置 statement_timeout防止单条 SQL 无限等锁
timeout_ms = _BATCH_STATEMENT_TIMEOUT_S * 1000
db.execute(text(f"SET LOCAL statement_timeout = '{timeout_ms}'"))
cleanup_key_references(db, batch)
result = db.execute(
sa_delete(ProviderAPIKey).where(
ProviderAPIKey.provider_id == provider_id,
ProviderAPIKey.id.in_(batch),
)
)
rowcount = getattr(result, "rowcount", 0) or 0
affected += int(rowcount)
db.commit()
except Exception as exc:
logger.warning(
"[BATCH_DELETE] batch {}/{} failed (keys {}-{}): {}",
batch_idx,
total_batches,
i,
i + len(batch),
exc,
)
try:
db.rollback()
except Exception:
pass
affected += _delete_key_batch_with_retry(
db,
provider_id,
batch,
batch_idx=batch_idx,
total_batches=total_batches,
start_offset=i,
)
# 每个批次都上报一次进度
if progress_callback is not None:
progress_callback(affected)

View File

@@ -4,6 +4,8 @@ Provider Key 写操作后的副作用处理。
from __future__ import annotations
from collections.abc import Callable
from sqlalchemy import delete as sa_delete
from sqlalchemy import update as sa_update
from sqlalchemy.orm import Session
@@ -12,7 +14,6 @@ from src.core.logger import logger
from src.models.database import (
GeminiFileMapping,
ProviderAPIKey,
RequestCandidate,
Usage,
VideoTask,
)
@@ -22,11 +23,34 @@ from src.services.cache.provider_cache import ProviderCacheService
_SQLITE_BATCH_SIZE = 900
_DEFAULT_BATCH_SIZE = 2000
_CLEANUP_STAGES = (
(
"gemini_file_mappings",
lambda batch: sa_delete(GeminiFileMapping).where(GeminiFileMapping.key_id.in_(batch)),
),
(
"usage",
lambda batch: sa_update(Usage)
.where(Usage.provider_api_key_id.in_(batch))
.values(provider_api_key_id=None),
),
(
"video_tasks",
lambda batch: sa_update(VideoTask).where(VideoTask.key_id.in_(batch)).values(key_id=None),
),
)
def cleanup_key_references(db: Session, key_ids: list[str]) -> None:
def cleanup_key_references(
db: Session,
key_ids: list[str],
*,
batch_size: int | None = None,
stage_callback: Callable[[str, int], None] | None = None,
) -> None:
"""在删除 ProviderAPIKey 前,先显式处理关联表引用,降低级联删除/置空成本。
- request_candidates / gemini_file_mappings: 直接删除
- gemini_file_mappings: 直接删除
- usage / video_tasks: 先置空外键,保留快照与历史记录
PostgreSQL 下直接按 key_id/provider_api_key_id 批量处理;
@@ -34,16 +58,12 @@ def cleanup_key_references(db: Session, key_ids: list[str]) -> None:
"""
if not key_ids:
return
batch_size = _resolve_batch_size(db)
for batch in _iter_batches(key_ids, batch_size):
db.execute(sa_delete(RequestCandidate).where(RequestCandidate.key_id.in_(batch)))
db.execute(sa_delete(GeminiFileMapping).where(GeminiFileMapping.key_id.in_(batch)))
db.execute(
sa_update(Usage)
.where(Usage.provider_api_key_id.in_(batch))
.values(provider_api_key_id=None)
)
db.execute(sa_update(VideoTask).where(VideoTask.key_id.in_(batch)).values(key_id=None))
effective_batch_size = batch_size if batch_size is not None else _resolve_batch_size(db)
for batch in iter_key_batches(key_ids, effective_batch_size):
for stage_name, statement_factory in _CLEANUP_STAGES:
if stage_callback is not None:
stage_callback(stage_name, len(batch))
db.execute(statement_factory(batch))
def _resolve_batch_size(db: Session) -> int:
@@ -57,10 +77,13 @@ def _resolve_batch_size(db: Session) -> int:
return _DEFAULT_BATCH_SIZE
def _iter_batches(items: list[str], batch_size: int) -> list[list[str]]:
def iter_key_batches(items: list[str], batch_size: int) -> list[list[str]]:
"""将 key_ids 列表按 batch_size 拆分为子列表。"""
if not items:
return []
if batch_size <= 0:
return [items]
return [items[i : i + batch_size] for i in range(0, len(items), batch_size)]
return [list(items)]
return [list(items[i : i + batch_size]) for i in range(0, len(items), batch_size)]
async def run_update_key_side_effects(

View File

@@ -144,6 +144,14 @@ class SystemConfigService:
"value": 1000,
"description": "每批次清理的记录数,避免单次操作过大影响数据库性能",
},
"request_candidates_retention_days": {
"value": 30,
"description": "请求候选记录保留天数,超过此天数的 request_candidates 审计记录将被自动清理",
},
"request_candidates_cleanup_batch_size": {
"value": 5000,
"description": "请求候选记录每批次清理条数,使用独立批次控制大表删除压力",
},
"enable_provider_checkin": {
"value": True,
"description": "是否启用 Provider 自动签到任务",

View File

@@ -786,13 +786,28 @@ class MaintenanceScheduler:
return 0
retention_days = max(
SystemConfigService.get_config(db, "detail_log_retention_days", 7),
SystemConfigService.get_config(
db,
"request_candidates_retention_days",
SystemConfigService.get_config(db, "detail_log_retention_days", 7),
),
3,
)
batch_size = SystemConfigService.get_config(db, "cleanup_batch_size", 1000)
batch_size = max(
SystemConfigService.get_config(
db,
"request_candidates_cleanup_batch_size",
SystemConfigService.get_config(db, "cleanup_batch_size", 1000),
),
1,
)
cutoff_time = datetime.now(timezone.utc) - timedelta(days=retention_days)
logger.info(f"开始清理 {retention_days} 天前的请求候选记录...")
logger.info(
"开始清理 {} 天前的请求候选记录batch_size={}",
retention_days,
batch_size,
)
except Exception as e:
logger.exception(f"候选记录清理配置读取失败: {e}")
return 0
@@ -806,6 +821,7 @@ class MaintenanceScheduler:
records_to_delete = (
batch_db.query(RequestCandidate.id)
.filter(RequestCandidate.created_at < cutoff_time)
.order_by(RequestCandidate.created_at.asc(), RequestCandidate.id.asc())
.limit(batch_size)
.all()
)