diff --git a/apps/worker/worker_app/tasks/edit_plan_generation.py b/apps/worker/worker_app/tasks/edit_plan_generation.py index 80b26b9c2..1ef351d58 100755 --- a/apps/worker/worker_app/tasks/edit_plan_generation.py +++ b/apps/worker/worker_app/tasks/edit_plan_generation.py @@ -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: