refactor: 删除edit_plan_generation中legacy渲染引擎分支,统一走unified

This commit is contained in:
2026-07-19 19:19:54 +08:00
committed by CI Bot
parent c93be457bc
commit d539124895
@@ -1,18 +1,13 @@
"""剪辑计划渲染任务 — 支持 Feature Flag 灰度.
"""剪辑计划渲染任务 — 使用 UnifiedRenderService 统一渲染引擎.
Celery 任务 worker.render_edit_plan:
1. 加载 EditPlan + EditPlanClips
2. 根据 Feature Flag 选择渲染引擎(legacy / unified)
2. 通过 RenderAdapter 调用 UnifiedRenderService 渲染
3. 下载各片段素材 + 渲染
4. 上传渲染结果到 OSS
5. 创建 GeneratedVideo 记录 + 查重
6. 更新 EditPlan / EditPlanClip 状态
7. 更新 GenerationTask 进度
渲染引擎灰度:
- 走 Feature Flag (render_engine) 控制
- legacy: VideoComposeService + FFmpeg filter_complex
- unified: UnifiedRenderService 图层架构
"""
from __future__ import annotations
@@ -66,34 +61,6 @@ def _get_repos():
# ── Celery Task ───────────────────────────────────────────────────────────────
def _resolve_render_engine(user_id: str) -> str:
"""根据 Feature Flag 决定使用哪个渲染引擎。
Returns:
"legacy" 或 "unified"
"""
try:
from video_processing.render_engine_resolver import get_render_engine_resolver
resolver = get_render_engine_resolver()
engine = resolver.get_engine(user_id=user_id)
# 灰度期间打印详细 flag 配置,便于排查
config = resolver.get_config_snapshot()
logger.info(
"edit_plan 引擎选择: user_id=%s engine=%s enabled=%s percentage=%s whitelist=%d default=%s",
user_id,
engine,
config.get("enabled"),
config.get("percentage"),
len(config.get("whitelist", [])),
config.get("default_engine"),
)
return engine
except Exception as exc:
logger.warning("获取渲染引擎配置失败,fallback 到 legacy: %s", exc, exc_info=True)
return "legacy"
def _mark_plan_failed(plan_repo, plan_id, gen_task_repo, generation_task_id, error_msg: str):
"""统一的计划失败标记工具。"""
plan = plan_repo.get(plan_id)
@@ -311,331 +278,13 @@ def _render_with_unified(
)
def _render_with_legacy(
plan,
clips,
rendered_clip_ids: list[str],
failed_clip_ids: list[str],
tmpdir_path: Path,
plan_id: str,
generation_task_id: str,
plan_repo,
clip_repo,
gen_task_repo,
db,
) -> dict:
"""旧引擎路径(VideoComposeService + FFmpeg filter_complex)。"""
import os
from apps.api.app.services.video_compose_service import VideoComposeService
compose_svc = VideoComposeService(db)
# 校验合成条件
validation = compose_svc.validate_compose(plan_id)
if not validation.valid:
error_msg = "; ".join(validation.errors)
logger.error("合成校验失败(legacy): %s — %s", plan_id, error_msg)
_mark_plan_failed(plan_repo, plan_id, gen_task_repo, generation_task_id, f"合成校验失败: {error_msg}")
return {"status": "error", "message": error_msg}
# 构建 FFmpeg 命令
output_dir = os.environ.get("VIDEO_OUTPUT_DIR", str(tmpdir_path))
output_path = Path(output_dir) / f"{plan_id}.mp4"
# 从 plan.config.export 读取输出分辨率,兼容 plan 自定义配置
plan_config = plan.config or {}
export_config = plan_config.get("export", {}) or {}
output_width = OUTPUT_WIDTH
output_height = OUTPUT_HEIGHT
resolution = export_config.get("resolution", "")
if resolution and "x" in resolution:
try:
w_str, h_str = resolution.lower().split("x", 1)
output_width = int(w_str)
output_height = int(h_str)
except (ValueError, TypeError):
pass
fps = export_config.get("fps", 25)
try:
fps = int(fps)
except (ValueError, TypeError):
fps = 25
compose_cmd = compose_svc.build_compose_command(
plan_id,
str(output_path),
output_width=output_width,
output_height=output_height,
fps=fps,
)
logger.info("执行 FFmpeg (legacy): plan_id=%s cmd=%s", plan_id, " ".join(compose_cmd.command)[:500])
# 开始渲染,更新进度
if generation_task_id:
try:
gen_task = gen_task_repo.get(generation_task_id)
if gen_task and gen_task.progress < 40.0:
gen_task.progress = 40.0
gen_task.append_log(
stage="render_start",
message="开始FFmpeg渲染(legacy)",
level="INFO",
progress=40.0,
)
gen_task_repo.update(gen_task)
except Exception:
pass
try:
from video_processing.ffmpeg_utils import run_ffmpeg
run_ffmpeg(compose_cmd.command, timeout=3600)
except Exception as e:
# 提取完整 stderr(如果是 CalledProcessError)
stderr_text = ""
if hasattr(e, "stderr"):
stderr_raw = e.stderr
if isinstance(stderr_raw, bytes):
stderr_text = stderr_raw.decode("utf-8", errors="replace")
elif isinstance(stderr_raw, str):
stderr_text = stderr_raw
# 完整命令(截断前2000字符,避免日志过大)
full_cmd = " ".join(compose_cmd.command)
cmd_preview = full_cmd[:2000] + ("..." if len(full_cmd) > 2000 else "")
# 拼接完整错误信息:命令 + 异常 + stderr最后1500字符
error_parts = [f"FFmpeg渲染失败(exit={getattr(e, 'returncode', 'unknown')})"]
error_parts.append("--- cmd ---")
error_parts.append(cmd_preview)
if stderr_text:
# 取最后1500字符,通常错误信息在末尾
stderr_preview = stderr_text[-1500:] if len(stderr_text) > 1500 else stderr_text
error_parts.append("--- stderr (last 1500 chars) ---")
error_parts.append(stderr_preview)
error_msg = "\n".join(error_parts)
logger.error("FFmpeg 执行失败(legacy): plan_id=%s\n%s", plan_id, error_msg)
_mark_plan_failed(plan_repo, plan_id, gen_task_repo, generation_task_id, error_msg)
return {"status": "error", "message": error_msg}
# 获取文件大小 + 实际时长
file_size = output_path.stat().st_size if output_path.exists() else 0
duration = compose_cmd.estimated_duration or 0.0
try:
from video_processing.ffmpeg_utils import probe_duration
actual_duration = probe_duration(str(output_path))
if actual_duration > 0:
duration = actual_duration
except Exception:
pass
# ── 标题/字幕叠加(legacy 引擎补齐) ────────────────────────────────
plan_config = plan.config or {}
title_cfg = plan_config.get("title", {}) or {}
subtitle_cfg = plan_config.get("subtitle", {}) or {}
title_text = title_cfg.get("text", "") or ""
subtitle_text = subtitle_cfg.get("text", "") or ""
title_enabled = title_cfg.get("enabled", True) and bool(title_text.strip())
subtitle_enabled = subtitle_cfg.get("enabled", True) and bool(subtitle_text.strip())
# ASR 自动字幕 legacy 暂不支持(需要额外 ASR 服务,统一用 unified 引擎)
has_subtitle_overlay = title_enabled or subtitle_enabled
if has_subtitle_overlay and output_path.exists() and duration > 0:
try:
from video_processing.ffmpeg_utils import run_ffmpeg
from video_processing.render_subtitles import generate_ass_subtitles
ass_path = tmpdir_path / f"subtitles_{plan_id}.ass"
generate_ass_subtitles(
ass_path,
video_width=output_width,
video_height=output_height,
video_duration=duration,
title_text=title_text,
title_config=title_cfg,
subtitle_text=subtitle_text,
subtitle_config=subtitle_cfg,
)
# 用 subtitles 滤镜叠加 ASS 字幕,音频直接 copy
subtitled_path = tmpdir_path / f"{plan_id}_subtitled.mp4"
# 处理 Windows 路径下的 ass 滤镜转义问题
ass_filter_path = str(ass_path).replace("\\", "/").replace(":", r"\:")
run_ffmpeg(
[
"ffmpeg",
"-y",
"-i",
str(output_path),
"-vf",
f"subtitles='{ass_filter_path}'",
"-c:a",
"copy",
str(subtitled_path),
],
timeout=1800,
)
if subtitled_path.exists() and subtitled_path.stat().st_size > 0:
output_path = subtitled_path
file_size = subtitled_path.stat().st_size
logger.info(
"legacy 标题/字幕叠加完成: plan_id=%s title=%s subtitle=%s",
plan_id,
title_enabled,
subtitle_enabled,
)
except Exception as sub_err:
logger.warning("legacy 标题/字幕叠加失败(不影响主流程): plan_id=%s err=%s", plan_id, sub_err)
# ── TTS 配音混音(legacy 引擎补齐) ────────────────────────────────
tts_cfg = plan_config.get("tts", {}) or {}
tts_enabled = tts_cfg.get("enabled", False) and bool(tts_cfg.get("text", "").strip())
if tts_enabled and output_path.exists() and duration > 0:
try:
from packages.domain.tts_config import TtsConfig
tts_config = TtsConfig.parse(tts_cfg)
if tts_config.enabled and tts_config.text.strip():
from apps.worker.services.tts_service_factory import get_tts_service
tts_service = get_tts_service()
voiceover_path = tmpdir_path / f"voiceover_{plan_id}.wav"
# 生成配音音频
audio_path = tts_service.synthesize(
text=tts_config.text,
voice_id=tts_config.voice_id,
speed=tts_config.speed,
pitch=tts_config.pitch,
output_path=voiceover_path,
)
if audio_path and audio_path.exists() and audio_path.stat().st_size > 0:
from video_processing.ffmpeg_utils import run_ffmpeg
mixed_path = tmpdir_path / f"{plan_id}_with_voiceover.mp4"
# 混音:配音音量按配置调整
voice_volume = max(0.0, min(1.0, tts_config.volume))
if tts_config.overlap_mode == "mix":
# 混音模式:原音 + 配音混合
filter_complex = (
f"[0:a]volume=1.0[a0];"
f"[1:a]volume={voice_volume:.2f}[a1];"
f"[a0][a1]amix=inputs=2:duration=first:dropout_transition=0[aout]"
)
else:
# replace 模式:配音替换原音
filter_complex = f"[1:a]volume={voice_volume:.2f}[aout]"
run_ffmpeg(
[
"ffmpeg",
"-y",
"-i",
str(output_path),
"-i",
str(audio_path),
"-filter_complex",
filter_complex,
"-map",
"0:v",
"-map",
"[aout]",
"-c:v",
"copy",
"-c:a",
"aac",
"-b:a",
"128k",
"-shortest",
str(mixed_path),
],
timeout=1800,
)
if mixed_path.exists() and mixed_path.stat().st_size > 0:
output_path = mixed_path
file_size = mixed_path.stat().st_size
logger.info(
"legacy TTS 配音混音完成: plan_id=%s voice_id=%s mode=%s",
plan_id,
tts_config.voice_id,
tts_config.overlap_mode,
)
except Exception as tts_err:
logger.warning("legacy TTS 配音混音失败(不影响主流程): plan_id=%s err=%s", plan_id, tts_err)
# 渲染完成,更新进度
if generation_task_id:
try:
gen_task = gen_task_repo.get(generation_task_id)
if gen_task and gen_task.progress < 80.0:
gen_task.progress = 80.0
gen_task.append_log(
stage="render_done",
message="FFmpeg渲染完成(legacy)",
level="INFO",
progress=80.0,
)
gen_task_repo.update(gen_task)
except Exception:
pass
# 上传到 OSS
storage_key = f"rendered/{plan_id}/output.mp4"
output_url = upload_to_oss(output_path, storage_key)
# 上传完成,更新进度
if generation_task_id:
try:
gen_task = gen_task_repo.get(generation_task_id)
if gen_task and gen_task.progress < 95.0:
gen_task.progress = 95.0
gen_task.append_log(
stage="upload_done",
message="OSS上传完成(legacy)",
level="INFO",
progress=95.0,
)
gen_task_repo.update(gen_task)
except Exception:
pass
return _finalize_render_success(
plan=plan,
plan_repo=plan_repo,
clip_repo=clip_repo,
gen_task_repo=gen_task_repo,
db=db,
plan_id=plan_id,
output_url=output_url or "",
storage_key=storage_key,
duration=duration,
file_size=file_size,
width=output_width,
height=output_height,
rendered_clip_ids=rendered_clip_ids,
failed_clip_ids=failed_clip_ids,
generation_task_id=generation_task_id,
output_path=output_path,
engine="legacy",
)
@celery_app.task(name="worker.render_edit_plan", bind=True, max_retries=2)
def render_edit_plan(self, plan_id: str) -> dict:
"""渲染剪辑计划
流程:
1. 加载 EditPlan + EditPlanClips
2. 根据 Feature Flag 选择渲染引擎(legacy / unified)
2. 通过 RenderAdapter 调用 UnifiedRenderService 渲染
3. 下载素材 + 渲染
4. 上传渲染结果到 OSS
5. 创建 GeneratedVideo 记录 + 查重
@@ -645,7 +294,6 @@ def render_edit_plan(self, plan_id: str) -> dict:
logger.info("开始渲染剪辑计划: plan_id=%s", plan_id)
generation_task_id = ""
engine = "legacy"
for repos in _get_repos():
plan_repo, clip_repo, gen_task_repo, db = repos
@@ -660,10 +308,8 @@ def render_edit_plan(self, plan_id: str) -> dict:
# 获取 generation_task_id(提前读取,确保 except 块可用)
generation_task_id = plan.config.get("generation_task_id", "")
# 2. 选择渲染引擎(Feature Flag 灰度控制)
# 2. 准备渲染(使用 unified 渲染引擎)
user_id = plan.created_by_user_id or ""
engine = _resolve_render_engine(user_id)
logger.info("剪辑计划渲染引擎: plan_id=%s engine=%s user_id=%s", plan_id, engine, user_id)
# 3. 加载片段列表(按 order 排序)
clips = clip_repo.list_by_plan(plan_id, skip=0, limit=10000)
@@ -681,9 +327,9 @@ def render_edit_plan(self, plan_id: str) -> dict:
gen_task.started_at = datetime.now(timezone.utc)
gen_task.append_log(
stage="render_start",
message=f"开始渲染,引擎 {engine},片段数 {len(clips)}",
message=f"开始渲染,片段数 {len(clips)}",
level="INFO",
engine=engine,
engine="unified",
clip_count=len(clips),
)
gen_task_repo.update(gen_task)
@@ -706,122 +352,19 @@ def render_edit_plan(self, plan_id: str) -> dict:
pass
return {"status": "cancelled", "plan_id": plan_id, "message": "任务已取消"}
# 4. 根据引擎选择渲染方式
if engine == "unified":
# ── unified 路径:RenderAdapter 统一处理(下载 + BGM + ASR + 渲染 + 上传)
result = _render_with_unified(
plan=plan,
clips=clips,
plan_id=plan_id,
generation_task_id=generation_task_id,
plan_repo=plan_repo,
clip_repo=clip_repo,
gen_task_repo=gen_task_repo,
db=db,
)
else:
# ── legacy 路径:原有的素材下载 + VideoComposeService
with tempfile.TemporaryDirectory(prefix="edit_plan_") as tmpdir:
tmpdir_path = Path(tmpdir)
asset_path_map: dict[str, Path] = {}
rendered_clip_ids: list[str] = []
failed_clip_ids: list[str] = []
# 4. 渲染(unified 引擎:RenderAdapter 统一处理下载 + BGM + ASR + 渲染 + 上传)
result = _render_with_unified(
plan=plan,
clips=clips,
plan_id=plan_id,
generation_task_id=generation_task_id,
plan_repo=plan_repo,
clip_repo=clip_repo,
gen_task_repo=gen_task_repo,
db=db,
)
# 预先批量查询所有素材的 storage_key
# 兼容存量数据:storage_key 为空时 fallback 到 file_url
from packages.adapters.sqlalchemy_impl.models import AssetModel
clip_asset_ids = [c.asset_id for c in clips if c.asset_id]
asset_storage_map: dict[str, str] = {}
if clip_asset_ids:
assets = db.query(AssetModel).filter(AssetModel.id.in_(clip_asset_ids)).all()
asset_storage_map = {
a.id: (a.storage_key or a.file_url or "") for a in assets if a.storage_key or a.file_url
}
for clip in clips:
if not clip.asset_id:
# 没有素材的片段跳过,标记为失败
clip.mark_failed()
clip_repo.update(clip)
failed_clip_ids.append(clip.id)
continue
if clip.asset_id in asset_path_map:
# 同一素材已下载(多个 clip 共享同一素材)
rendered_clip_ids.append(clip.id)
continue
storage_key = asset_storage_map.get(clip.asset_id)
if not storage_key:
logger.warning(
"片段素材无 storage_key,跳过: clip_id=%s asset_id=%s",
clip.id,
clip.asset_id,
)
clip.mark_failed()
clip_repo.update(clip)
failed_clip_ids.append(clip.id)
continue
# 下载素材
ext = Path(storage_key).suffix or ".mp4"
local_path = tmpdir_path / f"clip_{clip.order:04d}{ext}"
if download_asset(storage_key, local_path):
asset_path_map[clip.asset_id] = local_path
rendered_clip_ids.append(clip.id)
else:
clip.mark_failed()
clip_repo.update(clip)
failed_clip_ids.append(clip.id)
if not asset_path_map:
logger.error("所有片段素材下载失败: %s", plan_id)
plan.mark_failed()
plan_repo.update(plan)
if generation_task_id:
gen_task = gen_task_repo.get(generation_task_id)
if gen_task:
gen_task.status = "failed"
gen_task.error_message = "所有片段素材下载失败"
gen_task.completed_at = datetime.now(timezone.utc)
gen_task.append_log(
stage="download_failed",
message="所有片段素材下载失败",
level="ERROR",
)
gen_task_repo.update(gen_task)
return {"status": "error", "message": "所有片段素材下载失败"}
# 素材下载完成,记录日志
if generation_task_id:
gen_task = gen_task_repo.get(generation_task_id)
if gen_task:
gen_task.append_log(
stage="download_done",
message=f"素材下载完成,成功 {len(asset_path_map)} 个,失败 {len(failed_clip_ids)} 个",
level="INFO",
success_count=len(asset_path_map),
failed_count=len(failed_clip_ids),
)
gen_task.progress = 30.0
gen_task_repo.update(gen_task)
result = _render_with_legacy(
plan=plan,
clips=clips,
rendered_clip_ids=rendered_clip_ids,
failed_clip_ids=failed_clip_ids,
tmpdir_path=tmpdir_path,
plan_id=plan_id,
generation_task_id=generation_task_id,
plan_repo=plan_repo,
clip_repo=clip_repo,
gen_task_repo=gen_task_repo,
db=db,
)
result["engine"] = engine
result["engine"] = "unified"
return result
except Exception as exc: