244691d335
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (pull_request) Successful in 2s
CI/CD Pipeline / Check push changed paths (pull_request) Has been skipped
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 2s
AI Code Review / AI Code Review (pull_request) Failing after 1m46s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 2m47s
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 2m54s
CI/CD Pipeline / Validate - Style (pull_request) Failing after 2m29s
CI/CD Pipeline / Frontend Lint (pull_request) Has been skipped
CI/CD Pipeline / Frontend Unit Tests (pull_request) Has been skipped
CI/CD Pipeline / PR Build Web Image (pull_request) Has been skipped
CI/CD Pipeline / Integration Tests (pull_request) Successful in 2m6s
CI/CD Pipeline / Build Staging API Image (pull_request) Has been skipped
CI/CD Pipeline / Build Staging Web Image (pull_request) Has been skipped
CI/CD Pipeline / Build Staging Worker Image (pull_request) Has been skipped
CI/CD Pipeline / PR Build API Image (pull_request) Successful in 24s
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 22s
CI/CD Pipeline / Retag skipped Staging API Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging Web Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (pull_request) Has been skipped
CI/CD Pipeline / Staging E2E Tests (pull_request) Has been skipped
CI/CD Pipeline / Staging API Integration Tests (pull_request) Has been skipped
CI/CD Pipeline / ACR Image Cleanup (pull_request) Has been skipped
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Successful in 4m45s
CI/CD Pipeline / Unit Tests (pull_request) Successful in 8m59s
PR Automation / Auto Merge on CI Green + Approved (pull_request) Successful in 10m57s
ACR Cleanup / ACR Image Cleanup (pull_request_target) Successful in 7s
Preview Cleanup / Cleanup Preview Environment (pull_request) Successful in 1m37s
CI/CD Pipeline / Validate - Security (pull_request) Successful in 28m20s
CI/CD Pipeline / Build Production API Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production Web Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production Worker Image (pull_request) Has been skipped
CI/CD Pipeline / CI Gate (pull_request) Failing after 1s
CI/CD Pipeline / Deploy Production (pull_request) Has been skipped
CI/CD Pipeline / Canary Release to Production (pull_request) Has been skipped
CI/CD Pipeline / Production Browser E2E (pull_request) Has been skipped
175 lines
5.5 KiB
Python
175 lines
5.5 KiB
Python
#!/usr/bin/env python3
|
||
"""存量指纹重建脚本 — 为已有视频生成 video_fingerprint_chunks 分片数据。
|
||
|
||
功能:
|
||
- 查询 generated_videos 中 video_fingerprint IS NOT NULL 但尚无分片数据的视频
|
||
- 从 OSS 下载视频 → 用新的分片算法重新计算指纹 → 写入分片表
|
||
- 支持 --dry-run(只打印不写入)和 --batch-size(默认 50)
|
||
- 幂等:已存在分片数据的视频跳过
|
||
|
||
用法:
|
||
# 预览(不写入)
|
||
python rebuild_fingerprint_chunks.py --dry-run
|
||
|
||
# 执行重建
|
||
python rebuild_fingerprint_chunks.py --batch-size 50
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import argparse
|
||
import logging
|
||
import os
|
||
import sys
|
||
import tempfile
|
||
|
||
# 确保可以 import worker_app 和 packages
|
||
sys.path.insert(0, os.path.join(os.path.dirname(__file__), "..", "..", "..", "worker"))
|
||
sys.path.insert(0, os.path.join(os.path.dirname(__file__), "..", "..", ".."))
|
||
|
||
logging.basicConfig(
|
||
level=logging.INFO,
|
||
format="%(asctime)s [%(levelname)s] %(name)s: %(message)s",
|
||
)
|
||
logger = logging.getLogger("rebuild_fingerprint_chunks")
|
||
|
||
|
||
def find_videos_needing_rebuild(session, batch_size: int) -> list[dict]:
|
||
"""查询需要重建分片指纹的视频。"""
|
||
from sqlalchemy import and_
|
||
|
||
from packages.adapters.sqlalchemy_impl.models import GeneratedVideoModel, VideoFingerprintChunkModel
|
||
|
||
# 有 video_fingerprint 的视频
|
||
has_fingerprint = GeneratedVideoModel.video_fingerprint.isnot(None)
|
||
has_fingerprint = and_(has_fingerprint, GeneratedVideoModel.video_fingerprint != "")
|
||
|
||
# 排除已有分片数据的视频
|
||
subq = session.query(VideoFingerprintChunkModel.video_id).distinct().subquery()
|
||
no_chunks = ~GeneratedVideoModel.id.in_(subq)
|
||
|
||
videos = (
|
||
session.query(GeneratedVideoModel)
|
||
.filter(and_(has_fingerprint, no_chunks))
|
||
.order_by(GeneratedVideoModel.generated_at.desc())
|
||
.limit(batch_size)
|
||
.all()
|
||
)
|
||
|
||
return [
|
||
{
|
||
"id": v.id,
|
||
"project_id": v.project_id,
|
||
"user_id": v.user_id or "",
|
||
"duration": v.duration,
|
||
}
|
||
for v in videos
|
||
]
|
||
|
||
|
||
def rebuild_one(video_info: dict, dry_run: bool = False) -> int:
|
||
"""重建单个视频的分片数据。返回写入的 chunk 数量。"""
|
||
from video_processing.dedup import VideoDeduplicator, _save_fingerprint_chunks
|
||
from worker_app.db import SessionLocal
|
||
|
||
from packages.adapters.sqlalchemy_impl.models import VideoFingerprintChunkModel
|
||
from packages.shared.storage import get_storage_service
|
||
|
||
video_id = video_info["id"]
|
||
project_id = video_info["project_id"]
|
||
user_id = video_info["user_id"]
|
||
|
||
if dry_run:
|
||
logger.info("[DRY-RUN] Would rebuild video %s (project=%s)", video_id, project_id)
|
||
return 0
|
||
|
||
session = SessionLocal()
|
||
temp_dir = tempfile.mkdtemp()
|
||
|
||
try:
|
||
# 再次检查幂等性
|
||
existing_count = (
|
||
session.query(VideoFingerprintChunkModel).filter(VideoFingerprintChunkModel.video_id == video_id).count()
|
||
)
|
||
if existing_count > 0:
|
||
logger.info("Video %s already has %d chunks, skipping", video_id, existing_count)
|
||
return 0
|
||
|
||
# 下载视频
|
||
storage_service = get_storage_service()
|
||
local_path = os.path.join(temp_dir, f"{video_id}.mp4")
|
||
storage_key = f"projects/{project_id}/generated/{video_id}/{video_id}.mp4"
|
||
storage_service.download_file(storage_key, local_path)
|
||
|
||
# 重新计算指纹
|
||
deduplicator = VideoDeduplicator()
|
||
fingerprint = deduplicator.compute_fingerprint(local_path)
|
||
|
||
# 写入分片表
|
||
_save_fingerprint_chunks(fingerprint, video_id, project_id, user_id, session)
|
||
session.commit()
|
||
|
||
chunk_count = len(fingerprint.chunks)
|
||
logger.info("Rebuilt %d chunks for video %s", chunk_count, video_id)
|
||
return chunk_count
|
||
|
||
except Exception as e:
|
||
logger.error("Failed to rebuild video %s: %s", video_id, e)
|
||
session.rollback()
|
||
return -1
|
||
finally:
|
||
session.close()
|
||
import shutil
|
||
|
||
shutil.rmtree(temp_dir, ignore_errors=True)
|
||
|
||
|
||
def main():
|
||
parser = argparse.ArgumentParser(description="存量指纹重建脚本")
|
||
parser.add_argument("--dry-run", action="store_true", help="只打印不写入")
|
||
parser.add_argument("--batch-size", type=int, default=50, help="每批处理数量(默认 50)")
|
||
parser.add_argument("--total-limit", type=int, default=0, help="总处理数量限制(0=不限制)")
|
||
args = parser.parse_args()
|
||
|
||
from worker_app.db import SessionLocal
|
||
|
||
session = SessionLocal()
|
||
|
||
try:
|
||
videos = find_videos_needing_rebuild(session, args.batch_size)
|
||
logger.info("Found %d videos needing rebuild", len(videos))
|
||
|
||
if args.dry_run:
|
||
for v in videos:
|
||
logger.info("[DRY-RUN] Video %s | project=%s | duration=%.1fs", v["id"], v["project_id"], v["duration"])
|
||
return
|
||
|
||
total_chunks = 0
|
||
processed = 0
|
||
failed = 0
|
||
|
||
for v in videos:
|
||
if args.total_limit > 0 and processed >= args.total_limit:
|
||
break
|
||
|
||
result = rebuild_one(v, dry_run=False)
|
||
if result < 0:
|
||
failed += 1
|
||
else:
|
||
total_chunks += result
|
||
processed += 1
|
||
|
||
logger.info(
|
||
"Rebuild complete: processed=%d, chunks=%d, failed=%d",
|
||
processed,
|
||
total_chunks,
|
||
failed,
|
||
)
|
||
|
||
finally:
|
||
session.close()
|
||
|
||
|
||
if __name__ == "__main__":
|
||
main()
|