Files
xiaoxia-saas/apps/worker/worker_app/tasks/atom_clips.py
T
xiaoxia f1621ace9f
CI/CD Pipeline / Check if frontend-only change (push) Has been skipped
CI/CD Pipeline / PR Build API Image (push) Has been skipped
CI/CD Pipeline / PR Build Web Image (push) Has been skipped
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (push) Successful in 2s
CI/CD Pipeline / Check push changed paths (push) Successful in 5s
CI/CD Pipeline / Frontend Lint (push) Has been skipped
CI/CD Pipeline / PR Build Worker Image (push) Has been skipped
CI/CD Pipeline / Build Staging Worker Image (push) Successful in 2m17s
CI/CD Pipeline / Build Staging Web Image (push) Successful in 2m26s
CI/CD Pipeline / Integration Tests (push) Successful in 3m48s
CI/CD Pipeline / Frontend Unit Tests (push) Successful in 4m13s
CI/CD Pipeline / Build Staging API Image (push) Successful in 4m44s
CI/CD Pipeline / Retag skipped Staging API Image (push) Has been skipped
CI/CD Pipeline / Retag skipped Staging Web Image (push) Has been skipped
CI/CD Pipeline / Retag skipped Staging Worker Image (push) Has been skipped
CI/CD Pipeline / Validate - Python (mypy + alembic) (push) Successful in 5m19s
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (push) Successful in 59s
CI/CD Pipeline / Validate - Style (push) Successful in 7m51s
CI/CD Pipeline / ACR Image Cleanup (push) Successful in 2m54s
CI/CD Pipeline / Staging API Integration Tests (push) Successful in 3m44s
CI/CD Pipeline / Staging E2E Tests (push) Failing after 3m58s
CI/CD Pipeline / Unit Tests (push) Successful in 10m42s
CI/CD Pipeline / Validate - Security (push) Successful in 12m25s
CI/CD Pipeline / Build Production API Image (push) Has been skipped
CI/CD Pipeline / Build Production Web Image (push) Has been skipped
CI/CD Pipeline / Build Production Worker Image (push) Has been skipped
CI/CD Pipeline / CI Gate (push) Has been skipped
CI/CD Pipeline / Deploy Production (push) Has been skipped
CI/CD Pipeline / Canary Release to Production (push) Has been skipped
CI/CD Pipeline / Production Browser E2E (push) Has been skipped
feat(#1970): 素材原子化切片 P1 - 数据层/切片逻辑/原子片段级选片 (#1974)
Co-authored-by: xiaoxia <dev@xiaoxiajianji.com>
Co-committed-by: xiaoxia <dev@xiaoxiajianji.com>
2026-09-18 03:57:07 +08:00

82 lines
3.1 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""素材原子切片 Celery 任务 — #1970 智能剪辑流程重构 P1.
素材入库预处理完成(ingest 置 READY)后异步触发:
根据素材时长和已缓存的 scdet 切换点计算原子片段并落库。
失败不阻断素材入库主流程(atom_clips 未就绪时选片有内存兜底)。
"""
from __future__ import annotations
from celery.utils.log import get_task_logger
from worker_app.celery_app import celery_app
from worker_app.db import SessionLocal
from packages.adapters.sqlalchemy_impl.asset_atom_clip_repository import (
SQLAlchemyAssetAtomClipRepository,
)
from packages.adapters.sqlalchemy_impl.asset_repository import SQLAlchemyAssetRepository
from packages.domain.atom_clip_service import compute_atom_clips
from packages.domain.plan_generator_utils import extract_scene_points_from_metadata
logger = get_task_logger(__name__)
@celery_app.task(name="worker.generate_atom_clips")
def generate_atom_clips(asset_id: str) -> dict:
"""为单条视频素材生成原子片段。
Returns:
任务结果 dict:status / asset_id / clips_count。
"""
db = SessionLocal()
try:
asset_repo = SQLAlchemyAssetRepository(db)
atom_repo = SQLAlchemyAssetAtomClipRepository(db)
asset = asset_repo.find_by_id(asset_id)
if asset is None:
return {"status": "skipped", "reason": "asset not found", "asset_id": asset_id}
# 仅视频素材切片
if asset.mime_type and not asset.mime_type.startswith("video/"):
return {"status": "skipped", "reason": "not a video", "asset_id": asset_id}
if not asset.duration or asset.duration <= 0:
return {"status": "skipped", "reason": "invalid duration", "asset_id": asset_id}
# 已生成过则幂等跳过(重新切片需先显式删除)
existing = atom_repo.count_by_asset(asset_id)
if existing > 0:
return {
"status": "skipped",
"reason": "already generated",
"asset_id": asset_id,
"clips_count": existing,
}
scene_points = extract_scene_points_from_metadata(asset.metadata)
# P1 阶段继承素材的标签 ID;片段级语义标签是 P2 功能
tags = list(getattr(asset, "tag_ids", []) or [])
clips = compute_atom_clips(
asset_id=asset_id,
duration=float(asset.duration),
scene_change_points=scene_points,
tags=tags,
)
if not clips:
return {"status": "skipped", "reason": "no clips computed", "asset_id": asset_id}
atom_repo.batch_create(clips)
logger.info(
"[atom_clips] asset_id=%s 生成 %d 个原子片段",
asset_id,
len(clips),
)
return {"status": "completed", "asset_id": asset_id, "clips_count": len(clips)}
except Exception as exc: # noqa: BLE001 - 后台任务兜底,失败不阻断主流程
db.rollback()
logger.exception("[atom_clips] asset_id=%s 生成失败: %s", asset_id, exc)
return {"status": "failed", "asset_id": asset_id, "error": str(exc)}
finally:
db.close()