From a33dcf65e2afe39a226c0c23fdf54f07c699e418 Mon Sep 17 00:00:00 2001 From: xiaoxia Date: Sun, 19 Jul 2026 19:19:54 +0800 Subject: [PATCH] =?UTF-8?q?refactor:=20=E5=88=A0=E9=99=A4compose=5Fvideo?= =?UTF-8?q?=E4=B8=ADlegacy=E6=B8=B2=E6=9F=93=E5=BC=95=E6=93=8E=E5=88=86?= =?UTF-8?q?=E6=94=AF=EF=BC=8C=E7=BB=9F=E4=B8=80=E8=B5=B0unified?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- apps/worker/worker_app/tasks/compose_video.py | 95 +------------------ 1 file changed, 4 insertions(+), 91 deletions(-) diff --git a/apps/worker/worker_app/tasks/compose_video.py b/apps/worker/worker_app/tasks/compose_video.py index c9307d2da..f4b706395 100755 --- a/apps/worker/worker_app/tasks/compose_video.py +++ b/apps/worker/worker_app/tasks/compose_video.py @@ -1,6 +1,6 @@ """视频合成 Celery 任务 — Phase 8 任务 2.10. -使用 JobService 管理任务生命周期,集成 VideoComposeService 执行合成。 +使用 JobService 管理任务生命周期,通过 RenderAdapter 调用 UnifiedRenderService 执行合成。 """ from __future__ import annotations @@ -35,9 +35,7 @@ def _get_job_service(): def compose_video(self, job_id: str, **kwargs): """视频合成任务。 - 根据 RENDER_ENGINE 配置选择渲染引擎: - - legacy: 旧 VideoComposeService(filter_complex 模式) - - unified: 新 UnifiedRenderService(图层架构) + 使用 UnifiedRenderService(图层架构)进行渲染。 Args: job_id: JobService 中的任务 ID @@ -56,30 +54,8 @@ def compose_video(self, job_id: str, **kwargs): job_service.fail_job(job_id, "Missing plan_id in job payload") return {"status": "error", "message": "Missing plan_id"} - # 判断使用哪个渲染引擎 - # 优先级:Redis Feature Flag(白名单 > 百分比) > 环境变量默认 - from video_processing.render_engine_resolver import get_render_engine_resolver - - resolver = get_render_engine_resolver() - user_id = job.created_by_user_id or None - engine = resolver.get_engine(user_id=user_id) - # 灰度期间打印详细 flag 配置,便于排查 - config = resolver.get_config_snapshot() - logger.info( - "compose_video 引擎选择: job_id=%s engine=%s user_id=%s enabled=%s percentage=%s whitelist=%d default=%s", - job_id, - engine, - user_id, - config.get("enabled"), - config.get("percentage"), - len(config.get("whitelist", [])), - config.get("default_engine"), - ) - - if engine == "unified": - return _compose_with_unified_engine(self, job_service, job, plan_id, db) - else: - return _compose_with_legacy_engine(self, job_service, job, plan_id, db) + # 使用 unified 渲染引擎 + return _compose_with_unified_engine(self, job_service, job, plan_id, db) except self.retry_exc as exc: logger.warning("视频合成重试中: job_id=%s, exc=%s", job_id, exc) @@ -95,69 +71,6 @@ def compose_video(self, job_id: str, **kwargs): db.close() -def _compose_with_legacy_engine(task, job_service, job, plan_id: str, db) -> dict: - """旧引擎渲染路径(VideoComposeService)。""" - job_id = job.id - - # 标记为 running - job_service.update_progress(job_id, progress=10.0, current_stage="初始化合成环境") - - # 延迟导入 VideoComposeService - from apps.api.app.services.video_compose_service import VideoComposeService - - compose_svc = VideoComposeService(db) - - # 校验合成条件 - job_service.update_progress(job_id, progress=20.0, current_stage="校验合成条件") - validation = compose_svc.validate_compose(plan_id) - if not validation.valid: - error_msg = "; ".join(validation.errors) - job_service.fail_job(job_id, f"合成校验失败: {error_msg}") - return {"status": "error", "message": error_msg} - - # 构建合成命令 - job_service.update_progress(job_id, progress=30.0, current_stage="构建 FFmpeg 命令") - _output_dir = os.environ.get("VIDEO_OUTPUT_DIR", os.path.join(tempfile.gettempdir(), "video_output")) - output_path = os.path.join(_output_dir, f"{job_id}.mp4") - compose_cmd = compose_svc.build_compose_command(plan_id, output_path) - - # 执行 FFmpeg - job_service.update_progress(job_id, progress=50.0, current_stage="正在执行视频合成") - logger.info("Executing FFmpeg for job %s, plan %s", job_id, plan_id) - - try: - from video_processing.ffmpeg_utils import run_ffmpeg - - run_ffmpeg(compose_cmd.command, timeout=3600) - except Exception as e: - error_msg = f"FFmpeg 执行失败: {str(e)[:500]}" - job_service.fail_job(job_id, error_msg) - raise - - # 上传结果 - job_service.update_progress(job_id, progress=80.0, current_stage="上传合成结果") - storage_key = f"rendered/{plan_id}/{job_id}.mp4" - - from worker_app.tasks.edit_plan_generation import _upload_to_oss - - output_url = _upload_to_oss(Path(output_path), storage_key) - - # 更新 Job 状态为完成 - result_data = { - "plan_id": plan_id, - "output_path": output_path, - "storage_key": storage_key, - "output_url": output_url or "", - "estimated_duration": compose_cmd.estimated_duration, - "clip_count": len(compose_cmd.clip_chains), - "engine": "legacy", - } - job_service.complete_job(job_id, result=result_data) - - logger.info("视频合成完成(legacy): job_id=%s, plan_id=%s", job_id, plan_id) - return {"status": "completed", "job_id": job_id, "result": result_data} - - def _compose_with_unified_engine(task, job_service, job, plan_id: str, db) -> dict: """新引擎渲染路径(UnifiedRenderService + RenderAdapter)。""" job_id = job.id