From 100c8e01aa0f698bd321f2376311416e0bd7bc00 Mon Sep 17 00:00:00 2001 From: CI Bot Date: Mon, 13 Jul 2026 09:00:42 +0800 Subject: [PATCH] =?UTF-8?q?feat(worker):=20render=5Fedit=5Fplan=20?= =?UTF-8?q?=E6=8E=A5=E5=85=A5=20Feature=20Flag=20=E7=81=B0=E5=BA=A6?= =?UTF-8?q?=E6=8E=A7=E5=88=B6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 新增 _resolve_render_engine() 根据 Redis Feature Flag 选择引擎 - 新增 _render_with_unified() — 统一渲染引擎路径(UnifiedRenderService) - 新增 _render_with_legacy() — 旧引擎路径(VideoComposeService + FFmpeg) - 抽公共 _mark_plan_failed() / _finalize_render_success() 减少重复 - 默认走 legacy,灰度白名单/百分比命中才切 unified - 结果返回增加 engine 字段,便于灰度观测 - 修复测试 StubEditPlan 缺 created_by_user_id 属性 --- .../worker_app/tasks/edit_plan_generation.py | 401 +++++++++++++----- tests/unit/test_edit_plan_worker_failure.py | 2 + 2 files changed, 304 insertions(+), 99 deletions(-) mode change 100644 => 100755 apps/worker/worker_app/tasks/edit_plan_generation.py mode change 100644 => 100755 tests/unit/test_edit_plan_worker_failure.py diff --git a/apps/worker/worker_app/tasks/edit_plan_generation.py b/apps/worker/worker_app/tasks/edit_plan_generation.py old mode 100644 new mode 100755 index a89b9e3e9..a0011daf2 --- a/apps/worker/worker_app/tasks/edit_plan_generation.py +++ b/apps/worker/worker_app/tasks/edit_plan_generation.py @@ -1,13 +1,18 @@ -"""剪辑计划渲染任务 — Phase 8 任务 2.05. +"""剪辑计划渲染任务 — 支持 Feature Flag 灰度. Celery 任务 worker.render_edit_plan: 1. 加载 EditPlan + EditPlanClips - 2. 下载各片段素材 - 3. 使用 UnifiedRenderService 按时间线+图层渲染 + 2. 根据 Feature Flag 选择渲染引擎(legacy / unified) + 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 @@ -63,14 +68,268 @@ 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() + return resolver.get_engine(user_id=user_id) + except Exception as exc: + logger.warning("获取渲染引擎配置失败,fallback 到 legacy: %s", exc) + 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) + if plan and plan.status.value == "rendering": + plan.mark_failed() + plan_repo.update(plan) + if generation_task_id: + gen_task = gen_task_repo.get(generation_task_id) + if gen_task and gen_task.status.value != "failed": + gen_task.status = "failed" + gen_task.error_message = error_msg + gen_task.completed_at = datetime.now(timezone.utc) + gen_task_repo.update(gen_task) + + +def _finalize_render_success( + plan, + plan_repo, + clip_repo, + gen_task_repo, + db, + plan_id: str, + output_url: str, + storage_key: str, + duration: float, + file_size: int, + width: int, + height: int, + rendered_clip_ids: list[str], + failed_clip_ids: list[str], + generation_task_id: str, + output_path: Path, + engine: str, +) -> dict: + """渲染成功后的统一收尾:查重 + 更新状态 + 返回结果。""" + # 创建 GeneratedVideo 记录 + 查重 + project_id = plan.project_id or "" + batch_id = plan.config.get("batch_id", "") + mode = plan.config.get("mode", "edit_plan") + if generation_task_id and project_id: + try: + create_video_record_and_dedup( + generation_task_id=generation_task_id, + project_id=project_id, + batch_id=batch_id, + file_url=output_url or "", + file_size=file_size, + duration=duration, + video_path=str(output_path), + mode=mode, + session=db, + width=width, + height=height, + fps=OUTPUT_FPS, + ) + except Exception as dedup_err: + logger.warning("查重失败(不影响渲染结果): %s", dedup_err) + + # 更新片段状态为 rendered + for clip_id in rendered_clip_ids: + clip = clip_repo.get(clip_id) + if clip and clip.status.value == "ready": + clip.mark_rendered() + clip_repo.update(clip) + + # 更新 EditPlan 状态为 completed + plan.config["rendered_url"] = output_url or "" + plan.config["rendered_storage_key"] = storage_key + plan.mark_completed() + plan_repo.update(plan) + + # 更新 GenerationTask 状态为 completed + if generation_task_id: + gen_task = gen_task_repo.get(generation_task_id) + if gen_task: + gen_task.status = "completed" + gen_task.progress = 100.0 + gen_task.result_count = len(rendered_clip_ids) + gen_task.completed_at = datetime.now(timezone.utc) + gen_task_repo.update(gen_task) + + logger.info( + "剪辑计划渲染完成: plan_id=%s engine=%s rendered=%d failed=%d duration=%.1fs", + plan_id, + engine, + len(rendered_clip_ids), + len(failed_clip_ids), + duration, + ) + + return { + "status": "completed", + "plan_id": plan_id, + "rendered_count": len(rendered_clip_ids), + "failed_count": len(failed_clip_ids), + "output_url": output_url, + "duration": duration, + } + + +def _render_with_unified( + plan, + clips, + asset_path_map: dict[str, Path], + tmpdir_path: Path, + rendered_clip_ids: list[str], + plan_id: str, + generation_task_id: str, + plan_repo, + clip_repo, + gen_task_repo, + db, +) -> dict: + """统一渲染引擎路径(UnifiedRenderService 图层架构)。""" + render_service = UnifiedRenderService( + plan=plan, + clips=clips, + asset_path_map=asset_path_map, + work_dir=tmpdir_path, + output_width=OUTPUT_WIDTH, + output_height=OUTPUT_HEIGHT, + output_fps=int(OUTPUT_FPS), + ) + + try: + render_result = render_service.render() + except Exception as render_err: + logger.error("渲染失败(unified): %s — %s", plan_id, render_err) + _mark_plan_failed(plan_repo, plan_id, gen_task_repo, generation_task_id, f"渲染失败: {render_err}") + return {"status": "error", "message": f"渲染失败: {render_err}"} + + output_path = render_result.output_path + + # 上传到 OSS + storage_key = f"rendered/{plan_id}/output.mp4" + output_url = upload_to_oss(output_path, storage_key) + + failed_clip_ids: list[str] = [] + 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=render_result.duration, + file_size=render_result.file_size, + width=render_result.width, + height=render_result.height, + rendered_clip_ids=rendered_clip_ids, + failed_clip_ids=failed_clip_ids, + generation_task_id=generation_task_id, + output_path=output_path, + engine="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 + import subprocess + + 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" + compose_cmd = compose_svc.build_compose_command(plan_id, str(output_path)) + + logger.info("执行 FFmpeg (legacy): plan_id=%s", plan_id) + try: + subprocess.run( + compose_cmd.command, + check=True, + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + text=True, + timeout=3600, + ) + except subprocess.CalledProcessError as e: + error_msg = f"FFmpeg 执行失败: {e.stderr[:500]}" + logger.error("FFmpeg 执行失败(legacy): %s — %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 + + # 上传到 OSS + storage_key = f"rendered/{plan_id}/output.mp4" + output_url = upload_to_oss(output_path, storage_key) + + 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. 下载各片段素材到临时目录,构建 asset_path_map - 3. 使用 UnifiedRenderService 按时间线+图层渲染 + 2. 根据 Feature Flag 选择渲染引擎(legacy / unified) + 3. 下载素材 + 渲染 4. 上传渲染结果到 OSS 5. 创建 GeneratedVideo 记录 + 查重 6. 更新 EditPlan → completed, EditPlanClips → rendered @@ -79,6 +338,7 @@ 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 @@ -93,7 +353,12 @@ def render_edit_plan(self, plan_id: str) -> dict: # 获取 generation_task_id(提前读取,确保 except 块可用) generation_task_id = plan.config.get("generation_task_id", "") - # 2. 加载片段列表(按 order 排序) + # 2. 选择渲染引擎(Feature Flag 灰度控制) + 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) if not clips: logger.warning("剪辑计划没有片段: %s", plan_id) @@ -174,100 +439,38 @@ def render_edit_plan(self, plan_id: str) -> dict: gen_task_repo.update(gen_task) return {"status": "error", "message": "所有片段素材下载失败"} - # 4. 使用 UnifiedRenderService 渲染 - render_service = UnifiedRenderService( - plan=plan, - clips=clips, - asset_path_map=asset_path_map, - work_dir=tmpdir_path, - output_width=OUTPUT_WIDTH, - output_height=OUTPUT_HEIGHT, - output_fps=int(OUTPUT_FPS), - ) + # 4. 根据引擎选择渲染方式 + if engine == "unified": + result = _render_with_unified( + plan=plan, + clips=clips, + asset_path_map=asset_path_map, + tmpdir_path=tmpdir_path, + rendered_clip_ids=rendered_clip_ids, + 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: + 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, + ) - try: - render_result = render_service.render() - except Exception as render_err: - logger.error("渲染失败: %s — %s", plan_id, render_err) - 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 = f"渲染失败: {render_err}" - gen_task.completed_at = datetime.now(timezone.utc) - gen_task_repo.update(gen_task) - return {"status": "error", "message": f"渲染失败: {render_err}"} - - output_path = render_result.output_path - - # 5. 上传到 OSS - storage_key = f"rendered/{plan_id}/output.mp4" - output_url = upload_to_oss(output_path, storage_key) - - # 6. 创建 GeneratedVideo 记录 + 查重 - project_id = plan.project_id or "" - batch_id = plan.config.get("batch_id", "") - mode = plan.config.get("mode", "edit_plan") - if generation_task_id and project_id: - try: - create_video_record_and_dedup( - generation_task_id=generation_task_id, - project_id=project_id, - batch_id=batch_id, - file_url=output_url or "", - file_size=render_result.file_size, - duration=render_result.duration, - video_path=str(output_path), - mode=mode, - session=db, - width=render_result.width, - height=render_result.height, - fps=OUTPUT_FPS, - ) - except Exception as dedup_err: - logger.warning("查重失败(不影响渲染结果): %s", dedup_err) - - # 7. 更新片段状态为 rendered - for clip_id in rendered_clip_ids: - clip = clip_repo.get(clip_id) - if clip and clip.status.value == "ready": - clip.mark_rendered() - clip_repo.update(clip) - - # 8. 更新 EditPlan 状态为 completed - plan.config["rendered_url"] = output_url or "" - plan.config["rendered_storage_key"] = storage_key - plan.mark_completed() - plan_repo.update(plan) - - # 9. 更新 GenerationTask 状态为 completed - if generation_task_id: - gen_task = gen_task_repo.get(generation_task_id) - if gen_task: - gen_task.status = "completed" - gen_task.progress = 100.0 - gen_task.result_count = len(rendered_clip_ids) - gen_task.completed_at = datetime.now(timezone.utc) - gen_task_repo.update(gen_task) - - logger.info( - "剪辑计划渲染完成: plan_id=%s rendered=%d failed=%d duration=%.1fs", - plan_id, - len(rendered_clip_ids), - len(failed_clip_ids), - render_result.duration, - ) - - return { - "status": "completed", - "plan_id": plan_id, - "rendered_count": len(rendered_clip_ids), - "failed_count": len(failed_clip_ids), - "output_url": output_url, - "duration": render_result.duration, - } + result["engine"] = engine + return result except Exception as exc: logger.exception("渲染剪辑计划异常: %s", plan_id) diff --git a/tests/unit/test_edit_plan_worker_failure.py b/tests/unit/test_edit_plan_worker_failure.py old mode 100644 new mode 100755 index 137e28767..aad74f916 --- a/tests/unit/test_edit_plan_worker_failure.py +++ b/tests/unit/test_edit_plan_worker_failure.py @@ -63,6 +63,8 @@ class StubEditPlan: template_id: str = "tmpl-001" status: Any = None config: dict = field(default_factory=dict) + project_id: str = "" + created_by_user_id: str = "user-001" def mark_failed(self): self.status = _StubStatus("failed") -- 2.54.0