diff --git a/apps/worker/video_processing/render_adapter.py b/apps/worker/video_processing/render_adapter.py index bf20d9cc1..04a90dabe 100644 --- a/apps/worker/video_processing/render_adapter.py +++ b/apps/worker/video_processing/render_adapter.py @@ -485,6 +485,7 @@ class RenderAdapter: rendered_clip_ids: list[str] | None = None, failed_clip_ids: list[str] | None = None, voiceover_audio_path: str | None = None, + is_preview: bool = False, ) -> RenderAdapterResult: """执行统一渲染核心流程(BGM + ASR + 渲染 + 缩略图 + 上传)。 @@ -503,8 +504,8 @@ class RenderAdapter: self._report_progress(progress_cb, 40.0, "执行视频渲染") - # 2. 初始化 ASR - asr_service = self._get_asr_service() + # 2. 初始化 ASR(预览模式跳过,节省启动开销) + asr_service = None if is_preview else self._get_asr_service() # 3. 读取输出分辨率 plan_config = plan.config or {} @@ -529,25 +530,27 @@ class RenderAdapter: bgm_path=bgm_path, asr_service=asr_service, voiceover_audio_path=voiceover_audio_path, + is_preview=is_preview, ) result = render_svc.render() - # 4.5 渲染后校验输出完整性 - - validation = validate_video_output(result.output_path) - if not validation.valid: - logger.error( - "[render-adapter] 渲染输出校验失败: plan_id=%s job_id=%s error=%s", - plan_id, - job_id, - validation.error_message, - ) - return RenderAdapterResult( - success=False, - error_message=f"渲染输出校验失败: {validation.error_message}", - error_detail=validation.error_message, - ) - + # 4.5 渲染后校验输出完整性(预览模式跳过,节省耗时) + if is_preview: + logger.info("[render-adapter] 预览模式:跳过输出校验") + else: + validation = validate_video_output(result.output_path) + if not validation.valid: + logger.error( + "[render-adapter] 渲染输出校验失败: plan_id=%s job_id=%s error=%s", + plan_id, + job_id, + validation.error_message, + ) + return RenderAdapterResult( + success=False, + error_message=f"渲染输出校验失败: {validation.error_message}", + error_detail=validation.error_message, + ) self._report_progress(progress_cb, 80.0, "上传渲染结果") # 5. 上传结果 @@ -556,19 +559,20 @@ class RenderAdapter: self._report_progress(progress_cb, 90.0, "生成封面缩略图") - # 6. 生成缩略图 + # 6. 生成缩略图(预览模式跳过,节省耗时) thumbnail_url = "" - try: - from video_processing.thumbnail_generator import generate_and_upload_thumbnail + if not is_preview: + try: + from video_processing.thumbnail_generator import generate_and_upload_thumbnail - thumb_storage_key = f"rendered/{plan_id}/thumbnail.jpg" - thumbnail_url = generate_and_upload_thumbnail(str(result.output_path), thumb_storage_key) - except Exception as thumb_err: - logger.warning( - "[render-adapter] 缩略图生成失败(不影响主流程): plan_id=%s error=%s", - plan_id, - thumb_err, - ) + thumb_storage_key = f"rendered/{plan_id}/thumbnail.jpg" + thumbnail_url = generate_and_upload_thumbnail(str(result.output_path), thumb_storage_key) + except Exception as thumb_err: + logger.warning( + "[render-adapter] 缩略图生成失败(不影响主流程): plan_id=%s error=%s", + plan_id, + thumb_err, + ) self._report_progress(progress_cb, 100.0, "渲染完成") @@ -614,6 +618,7 @@ class RenderAdapter: work_dir: Path | None = None, progress_cb: ProgressCallback | None = None, voiceover_audio_path: str | None = None, + is_preview: bool = False, ) -> RenderAdapterResult: """使用内存中的 plan/clips/asset_path_map 直接渲染。 @@ -672,6 +677,7 @@ class RenderAdapter: job_id=job_id, progress_cb=progress_cb, voiceover_audio_path=voiceover_audio_path, + is_preview=is_preview, ) except subprocess.CalledProcessError as exc: diff --git a/apps/worker/video_processing/unified_render_service.py b/apps/worker/video_processing/unified_render_service.py index c8f6229c8..57d8e5b2e 100755 --- a/apps/worker/video_processing/unified_render_service.py +++ b/apps/worker/video_processing/unified_render_service.py @@ -151,6 +151,7 @@ class UnifiedRenderService: asr_service: Any = None, # ASRService 实例,用于自动生成字幕 bgm_path: str | None = None, # BGM 本地文件路径 voiceover_audio_path: str | None = None, # 配音素材库音频本地路径 + is_preview: bool = False, # 预览模式:ultrafast 编码 + 跳过非必要步骤 ): self.plan = plan self.clips = clips @@ -163,6 +164,7 @@ class UnifiedRenderService: self.asr_service = asr_service self.bgm_path = bgm_path self.voiceover_audio_path = voiceover_audio_path + self.is_preview = is_preview self._transition_engine = TransitionEngine(default_duration=transition_duration) self._speed_engine = SpeedEngine() self._asr_timeline_cache: Any = None # ASR 字幕结果缓存,避免重复调用 @@ -1228,9 +1230,9 @@ class UnifiedRenderService: "-c:v", "libx264", "-crf", - "23", + "28" if self.is_preview else "23", "-preset", - "medium", + "ultrafast" if self.is_preview else "medium", "-pix_fmt", "yuv420p", "-movflags", @@ -1742,9 +1744,9 @@ class UnifiedRenderService: "-c:v", "libx264", "-crf", - "23", + "28" if self.is_preview else "23", "-preset", - "medium", + "ultrafast" if self.is_preview else "medium", "-pix_fmt", "yuv420p", "-movflags", @@ -1753,10 +1755,11 @@ class UnifiedRenderService: ] logger.info( - "执行渲染: plan_id=%s inputs=%d output=%s", + "执行渲染: plan_id=%s inputs=%d output=%s preview=%s", self.plan.id, input_args.count("-i"), output_path, + self.is_preview, ) try: run_ffmpeg(command) diff --git a/apps/worker/worker_app/tasks/generation.py b/apps/worker/worker_app/tasks/generation.py index d2c9ffb83..04ca03599 100644 --- a/apps/worker/worker_app/tasks/generation.py +++ b/apps/worker/worker_app/tasks/generation.py @@ -743,7 +743,8 @@ def _download_library_assets( len(asset_ids), ) - downloaded: list[Path] = [] + # 构建待下载列表 (index, asset, storage_key, local_file) + download_jobs: list[tuple[int, Any, str, Path]] = [] failed_assets: list[str] = [] for i, asset in enumerate(assets): storage_key = asset.file_url if asset.file_url else None @@ -772,52 +773,81 @@ def _download_library_assets( ext = Path(storage_key).suffix or ".mp4" local_file = temp_path / f"asset_{i:03d}_{asset.id}{ext}" - asset_start = time.monotonic() - download_ok = download_asset(storage_key, local_file) - asset_elapsed = time.monotonic() - asset_start + download_jobs.append((i, asset, storage_key, local_file)) - if download_ok: - file_size = local_file.stat().st_size if local_file.exists() else 0 - downloaded.append(local_file) - logger.info( - "[task_id=%s] Downloaded asset: %s -> %s (size=%d, time=%.1fs)", - task_id, - asset.name, - local_file, - file_size, - asset_elapsed, - ) - if gen_task: - gen_task.append_log( - "下载素材", - f"下载成功: {asset.name}", - asset_id=asset.id, - asset_name=asset.name, - success=True, - file_size=file_size, - duration=round(asset_elapsed, 2), + # 并行下载素材(线程池,IO 密集型) + downloaded: list[Path] = [] + if download_jobs: + from concurrent.futures import ThreadPoolExecutor, as_completed + + max_workers = min(len(download_jobs), 6) + logger.info( + "[task_id=%s] 并行下载素材: count=%d, workers=%d", + task_id, + len(download_jobs), + max_workers, + ) + + def _download_one(item: tuple) -> tuple[int, Any, Path, bool, float]: + idx, asset, skey, lfile = item + t0 = time.monotonic() + ok = download_asset(skey, lfile) + elapsed = time.monotonic() - t0 + return idx, asset, lfile, ok, elapsed + + with ThreadPoolExecutor(max_workers=max_workers) as executor: + futures = {executor.submit(_download_one, job): job for job in download_jobs} + # 按原始顺序收集结果,保证 downloaded 列表顺序稳定 + results_map: dict[int, tuple[Path, bool, float, Any]] = {} + for future in as_completed(futures): + idx, asset, lfile, ok, elapsed = future.result() + results_map[idx] = (lfile, ok, elapsed, asset) + + # 按原始顺序处理结果 + for idx in sorted(results_map.keys()): + lfile, ok, elapsed, asset = results_map[idx] + if ok: + file_size = lfile.stat().st_size if lfile.exists() else 0 + downloaded.append(lfile) + logger.info( + "[task_id=%s] Downloaded asset: %s -> %s (size=%d, time=%.1fs)", + task_id, + asset.name, + lfile, + file_size, + elapsed, ) - else: - failed_assets.append(f"{asset.name}({asset.id})") - logger.warning( - "[task_id=%s] Failed to download asset: %s (id=%s)", - task_id, - asset.name, - asset.id, - ) - if gen_task: - gen_task.append_log( - "下载素材", - f"下载失败: {asset.name}", - level="WARN", - asset_id=asset.id, - asset_name=asset.name, - success=False, - file_size=0, - duration=round(asset_elapsed, 2), + if gen_task: + gen_task.append_log( + "下载素材", + f"下载成功: {asset.name}", + asset_id=asset.id, + asset_name=asset.name, + success=True, + file_size=file_size, + duration=round(elapsed, 2), + ) + else: + failed_assets.append(f"{asset.name}({asset.id})") + logger.warning( + "[task_id=%s] Failed to download asset: %s (id=%s)", + task_id, + asset.name, + asset.id, ) - if strict: - raise RuntimeError(f"素材下载失败: asset_id={asset.id}, name={asset.name}") + if gen_task: + gen_task.append_log( + "下载素材", + f"下载失败: {asset.name}", + level="WARN", + asset_id=asset.id, + asset_name=asset.name, + success=False, + file_size=0, + duration=round(elapsed, 2), + ) + if strict: + raise RuntimeError(f"素材下载失败: asset_id={asset.id}, name={asset.name}") # 指定了 asset_ids 但全部下载失败 → 无论 strict 与否都报错 if asset_ids and not downloaded: @@ -1171,6 +1201,7 @@ def _render_video( job_id=task_id, work_dir=temp_path, voiceover_audio_path=voice_path, + is_preview=is_preview, ) finally: db.close() diff --git a/tests/unit/test_1280_preview_speedup.py b/tests/unit/test_1280_preview_speedup.py new file mode 100644 index 000000000..13cb516cd --- /dev/null +++ b/tests/unit/test_1280_preview_speedup.py @@ -0,0 +1,289 @@ +"""#1280 预览视频生成加速 — 单元测试。 + +验证点: +1. UnifiedRenderService.is_preview 参数正确传递 +2. 预览模式使用 ultrafast preset + crf 28 +3. RenderAdapter.render_from_memory 正确传递 is_preview +4. 预览模式跳过 ASR 初始化 +5. 预览模式跳过输出校验和缩略图 +6. generation.py 并行下载逻辑 +""" + +from __future__ import annotations + +import os +import tempfile +from pathlib import Path +from unittest.mock import MagicMock, patch + +import pytest + +# ── 1. UnifiedRenderService is_preview 参数 ── + + +class TestUnifiedRenderServicePreviewFlag: + """is_preview 参数正确传递和存储。""" + + def test_default_is_preview_false(self): + from video_processing.unified_render_service import UnifiedRenderService + + svc = UnifiedRenderService( + plan=MagicMock(id="test"), + clips=[], + asset_path_map={}, + work_dir=Path(tempfile.mkdtemp()), + ) + assert svc.is_preview is False + + def test_is_preview_true(self): + from video_processing.unified_render_service import UnifiedRenderService + + svc = UnifiedRenderService( + plan=MagicMock(id="test"), + clips=[], + asset_path_map={}, + work_dir=Path(tempfile.mkdtemp()), + is_preview=True, + ) + assert svc.is_preview is True + + def test_is_preview_false_explicit(self): + from video_processing.unified_render_service import UnifiedRenderService + + svc = UnifiedRenderService( + plan=MagicMock(id="test"), + clips=[], + asset_path_map={}, + work_dir=Path(tempfile.mkdtemp()), + is_preview=False, + ) + assert svc.is_preview is False + + +# ── 2. 预览模式 FFmpeg 参数 ── + + +class TestPreviewFFmpegPreset: + """预览模式使用 ultrafast preset + crf 28。""" + + def _make_clip(self): + from video_processing.unified_render_service import ResolvedClip + + return ResolvedClip( + clip_id="c1", + asset_id="a1", + local_path=Path("/tmp/fake.mp4"), + clip_type="main", + order=0, + start_time=0, + duration=10.0, + playback_speed=1.0, + transition_effect="cut", + transition_duration=0.0, + config={}, + ) + + @patch("video_processing.unified_render_service.run_ffmpeg") + def test_execute_ffmpeg_preview_uses_ultrafast(self, mock_run): + from video_processing.unified_render_service import ( + RenderLayer, + UnifiedRenderService, + ) + + plan = MagicMock() + plan.id = "test_plan" + plan.config = {"export": {"resolution": "854x480"}} + + clip = self._make_clip() + + svc = UnifiedRenderService( + plan=plan, + clips=[clip], + asset_path_map={"a1": Path("/tmp/fake.mp4")}, + work_dir=Path(tempfile.mkdtemp()), + output_width=854, + output_height=480, + is_preview=True, + ) + + layers = [RenderLayer(role="main", clips=[clip])] + filter_complex, input_args = svc._build_filter_complex(layers) + output_path = Path(tempfile.mkdtemp()) / "out.mp4" + svc._execute_ffmpeg(filter_complex, input_args, output_path) + + mock_run.assert_called_once() + cmd = mock_run.call_args[0][0] + + # Check preset is ultrafast + preset_idx = cmd.index("-preset") + assert cmd[preset_idx + 1] == "ultrafast", f"Expected ultrafast, got {cmd[preset_idx + 1]}" + + # Check crf is 28 + crf_idx = cmd.index("-crf") + assert cmd[crf_idx + 1] == "28", f"Expected crf 28, got {cmd[crf_idx + 1]}" + + @patch("video_processing.unified_render_service.run_ffmpeg") + def test_execute_ffmpeg_normal_uses_medium(self, mock_run): + from video_processing.unified_render_service import ( + RenderLayer, + UnifiedRenderService, + ) + + plan = MagicMock() + plan.id = "test_plan" + plan.config = {} + + clip = self._make_clip() + + svc = UnifiedRenderService( + plan=plan, + clips=[clip], + asset_path_map={"a1": Path("/tmp/fake.mp4")}, + work_dir=Path(tempfile.mkdtemp()), + output_width=1280, + output_height=720, + is_preview=False, + ) + + layers = [RenderLayer(role="main", clips=[clip])] + filter_complex, input_args = svc._build_filter_complex(layers) + output_path = Path(tempfile.mkdtemp()) / "out.mp4" + svc._execute_ffmpeg(filter_complex, input_args, output_path) + + mock_run.assert_called_once() + cmd = mock_run.call_args[0][0] + + preset_idx = cmd.index("-preset") + assert cmd[preset_idx + 1] == "medium" + + crf_idx = cmd.index("-crf") + assert cmd[crf_idx + 1] == "23" + + +# ── 3. RenderAdapter passes is_preview ── + + +class TestRenderAdapterPreviewPassthrough: + """RenderAdapter 正确传递 is_preview 参数。""" + + def test_render_from_memory_passes_is_preview(self): + from video_processing.render_adapter import RenderAdapter + + db = MagicMock() + adapter = RenderAdapter(db) + + plan = MagicMock() + plan.id = "test_plan" + plan.config = {"export": {"resolution": "854x480"}} + + clip = MagicMock() + clip.id = "c1" + + with patch.object(adapter, "_do_render") as mock_do_render: + mock_do_render.return_value = MagicMock( + success=True, + output_path=Path("/tmp/out.mp4"), + thumbnail_url="", + duration=5.0, + file_size=1000, + width=854, + height=480, + output_url="https://oss/test.mp4", + rendered_clip_ids=["c1"], + failed_clip_ids=[], + ) + + adapter.render_from_memory( + plan=plan, + clips=[clip], + asset_path_map={"a1": Path("/tmp/fake.mp4")}, + is_preview=True, + ) + + mock_do_render.assert_called_once() + _, kwargs = mock_do_render.call_args + assert kwargs.get("is_preview") is True + + +# ── 4. 预览模式跳过 ASR ── + + +class TestPreviewSkipsASR: + """预览模式跳过 ASR 初始化。""" + + def test_render_method_source_has_asr_skip(self): + """_do_render 在 is_preview=True 时不调用 _get_asr_service。""" + with open("apps/worker/video_processing/render_adapter.py") as f: + source = f.read() + + assert ( + "None if is_preview else self._get_asr_service()" in source + ), "Should skip ASR initialization in preview mode" + + +# ── 5. 并行下载逻辑 ── + + +class TestParallelDownload: + """generation.py 并行下载素材。""" + + def test_parallel_download_uses_thread_pool(self): + with open("apps/worker/worker_app/tasks/generation.py") as f: + source = f.read() + + assert "ThreadPoolExecutor" in source, "Should use ThreadPoolExecutor for parallel downloads" + assert "as_completed" in source, "Should use as_completed for result collection" + + def test_parallel_download_preserves_order(self): + with open("apps/worker/worker_app/tasks/generation.py") as f: + source = f.read() + + assert "sorted(results_map.keys())" in source, "Should sort results by original index" + + +# ── 6. generation.py _render_video passes is_preview ── + + +class TestRenderVideoPassesPreview: + """_render_video 正确传递 is_preview 到 render_from_memory。""" + + def test_render_video_passes_is_preview(self): + with open("apps/worker/worker_app/tasks/generation.py") as f: + source = f.read() + + assert "is_preview=is_preview" in source, "Should pass is_preview to render_from_memory" + + +# ── 7. Preview mode skips thumbnail and validation ── + + +class TestPreviewSkipsThumbnailAndValidation: + """预览模式跳过缩略图生成和输出校验。""" + + def test_render_adapter_skips_thumbnail_in_preview(self): + with open("apps/worker/video_processing/render_adapter.py") as f: + source = f.read() + + assert "if not is_preview:" in source, "Thumbnail should be conditional on is_preview" + + def test_render_adapter_skips_validation_in_preview(self): + with open("apps/worker/video_processing/render_adapter.py") as f: + source = f.read() + + assert "预览模式:跳过输出校验" in source, "Should skip validation in preview mode" + + +# ── 8. Pass-through rendering uses ultrafast in preview ── + + +class TestPassThroughPreviewPreset: + """直通渲染在预览模式也使用 ultrafast。""" + + def test_pass_through_has_preview_preset(self): + with open("apps/worker/video_processing/unified_render_service.py") as f: + source = f.read() + + # The pass_through method should also use ultrafast for preview + # Count occurrences of "ultrafast" - should be at least 2 (execute_ffmpeg + pass_through) + count = source.count('"ultrafast" if self.is_preview') + assert count >= 2, f"Expected at least 2 ultrafast preset usages, found {count}"