Files
xiaoxia-saas/apps/worker/worker_app/tasks/atom_clips.py
T
CI Bot 01018a23f4
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (pull_request) Successful in 1s
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 1s
CI/CD Pipeline / Check push changed paths (pull_request) Has been skipped
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 / 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 / PR Build Web Image (pull_request) Successful in 1m2s
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 1m3s
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 1m49s
CI/CD Pipeline / Integration Tests (pull_request) Successful in 2m9s
CI/CD Pipeline / PR Build API Image (pull_request) Successful in 2m30s
CI/CD Pipeline / Frontend Unit Tests (pull_request) Successful in 2m36s
CI/CD Pipeline / Validate - Style (pull_request) Successful in 2m42s
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Successful in 3m0s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 3m15s
CI/CD Pipeline / Validate - Security (pull_request) Successful in 4m51s
CI/CD Pipeline / Frontend Lint (pull_request) Successful in 5m5s
AI Code Review / AI Code Review (pull_request) Successful in 6m46s
CI/CD Pipeline / Unit Tests (pull_request) Failing after 12m49s
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 2s
CI/CD Pipeline / Deploy Production (pull_request) Has been skipped
CI/CD Pipeline / Production Browser E2E (pull_request) Has been skipped
CI/CD Pipeline / Canary Release to Production (pull_request) Has been skipped
PR Automation / Auto Merge on CI Green + Approved (pull_request) Successful in 9m46s
style: auto-format with black + isort + ruff + prettier [skip ci-format-check]
2026-09-25 03:58:32 +00:00

140 lines
5.6 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 未就绪时选片有内存兜底)。
P2 增强:切片完成后自动链式触发 AI 标签任务(每个 clip 一个 tag_atom_clip 任务)。
"""
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
from packages.shared.mediakit_client import get_mediakit_client
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)
# #2035:metadata 中没有 scene_change_points 时,按需调用 MediaKit 检测
# (templates_editor 路由会主动写 metadata,ingest 流程此前未触发检测导致切点无法对齐)
if not scene_points:
try:
mk = get_mediakit_client()
video_url = getattr(asset, "file_url", "") or ""
if mk.is_available and video_url:
timestamps = mk.detect_scene_changes(video_url)
if timestamps:
scene_points = timestamps
# 持久化到 metadata,避免下次重复检测
new_meta = dict(asset.metadata or {})
new_meta["scene_change_points"] = list(timestamps)
asset.metadata = new_meta
asset_repo.update(asset)
db.commit()
logger.info(
"[atom_clips] asset_id=%s 自动检测到 %d 个场景切换点并写回metadata",
asset_id,
len(timestamps),
)
except Exception as detect_err: # noqa: BLE001
logger.warning(
"[atom_clips] asset_id=%s scene_change自动检测失败,降级为均匀切片: %s",
asset_id,
detect_err,
)
db.rollback() # 回滚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),
)
# P2 增强:链式触发 AI 标签任务(每个 clip 一个异步任务)
_dispatch_tagging_tasks(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()
def _dispatch_tagging_tasks(clips: list) -> None:
"""为每个新建片段发送 AI 标签异步任务.
失败不阻断(标签任务是锦上添花,不影响核心流程)。
"""
try:
for clip in clips:
celery_app.send_task(
"worker.tag_atom_clip",
args=[clip.id],
)
logger.info(
"[atom_clips] 已发送 %d 个 AI 标签任务",
len(clips),
)
except Exception as e:
logger.warning(
"[atom_clips] 发送 AI 标签任务失败(不影响切片结果): %s",
e,
)