diff --git a/apps/worker/video_processing/gpu_direct_pipeline.py b/apps/worker/video_processing/gpu_direct_pipeline.py index 5203f7eff..ea4e3bc20 100644 --- a/apps/worker/video_processing/gpu_direct_pipeline.py +++ b/apps/worker/video_processing/gpu_direct_pipeline.py @@ -3,8 +3,8 @@ 背景:旧链路 worker 先用 CPU libx264 把 filter_complex 输出成 mezzanine(1080p 约 85s), 上传后再由 P4000 NVENC 编码,渲染后还要单独跑一次随机边缘裁剪重编码(约 26s)。 本管线取消 mezzanine:把原始素材签名 URL 作为多输入直接交给 P4000,filter_complex 内 -一步完成 trim/scale/pad/concat/边缘裁剪/drawtext 字幕,末端 h264_nvenc 只编码一次; -TTS/BGM 音频也在同一命令里 amix 合成。 +一步完成 trim/scale/pad/concat/边缘随机裁剪/drawtext 字幕,末端 h264_nvenc 只编码一次; +原素材音轨 concat + TTS/BGM 混音也在同一命令里完成。 约束(P1): - 仅覆盖智能剪辑主流场景:单一主视频轨、全硬切、无 PiP/overlay/水印/贴纸/片头片尾/绿幕。 @@ -15,21 +15,21 @@ TTS/BGM 音频也在同一命令里 amix 合成。 from __future__ import annotations import logging +import random import uuid from pathlib import Path from typing import Any, Optional logger = logging.getLogger(__name__) -# drawtext 默认字体名(fontconfig 解析);P4000 装好中文字体后可在 config 指定 DEFAULT_DRAWTEXT_FONT = "Noto Sans CJK SC" +EDGE_CROP_MIN_PCT = 0.02 +EDGE_CROP_MAX_PCT = 0.05 def escape_drawtext_text(text: str) -> str: - """转义 drawtext text= 中的特殊字符(ffmpeg 过滤器语法)。""" if not text: return "" - # 顺序重要:先转义反斜杠本身 s = text.replace("\\", "\\\\") s = s.replace(":", "\\:") s = s.replace("'", "\\'") @@ -37,7 +37,6 @@ def escape_drawtext_text(text: str) -> str: s = s.replace(",", "\\,") s = s.replace("[", "\\[").replace("]", "\\]") s = s.replace(";", "\\;") - # 换行保留(drawtext 支持 %{...};真实换行需转成字面) s = s.replace("\n", " ") return s @@ -58,10 +57,6 @@ def build_drawtext_filter( border_color: str = "black", enable: bool = True, ) -> str: - """构造单个 drawtext 滤镜字符串(不含输入/输出标签)。 - - [start, end] 秒的显示窗口通过 enable='between(t,...)' 控制。 - """ txt = escape_drawtext_text(text) parts = [f"font={font}", f"text='{txt}'"] if font_size and font_size > 0: @@ -80,18 +75,29 @@ def build_drawtext_filter( return "drawtext=" + ":".join(parts) +def _build_atempo_chain(speed: float) -> str: + if abs(speed - 1.0) < 1e-6: + return "" + stages: list[float] = [] + remaining = speed + while remaining > 2.0: + stages.append(2.0) + remaining /= 2.0 + while remaining < 0.5: + stages.append(0.5) + remaining /= 0.5 + if abs(remaining - 1.0) >= 1e-6: + stages.append(remaining) + return ",".join(f"atempo={s:.5f}" for s in stages) + + def upload_local_audio_and_sign( local_audio: Path, *, tmp_prefix: str = "tmp/gpu-direct-audio/", expires: int = 3600, ) -> tuple[str, str]: - """把本地音频(TTS/BGM)上传 OSS tmp 目录并签公网 GET URL。 - - Returns: - (signed_get_url, oss_key) - """ - from video_processing.oss_helpers import _storage # 类型: ignore + from video_processing.oss_helpers import _storage # type: ignore storage = _storage() key = f"{tmp_prefix.rstrip('/')}/{uuid.uuid4().hex}{local_audio.suffix or '.mp3'}" @@ -102,23 +108,17 @@ def upload_local_audio_and_sign( def sign_asset_url(storage_key: str, *, expires: int = 3600) -> str: - """给原始素材 storage_key 签公网 GET URL(供 P4000 直接下载)。""" - from video_processing.oss_helpers import _storage # 类型: ignore + from video_processing.oss_helpers import _storage # type: ignore storage = _storage() return storage.get_download_url(storage_key, expires) -# ── 直连渲染编排器 ──────────────────────────────────────────────────────── - - class DirectRenderPlan: - """一次直连渲染的产物:inputs(裸文件名→URL)与完整 ffmpeg_args。""" - def __init__(self, inputs: dict[str, str], ffmpeg_args: list[str], oss_keys: list[str]): self.inputs = inputs self.ffmpeg_args = ffmpeg_args - self.oss_keys = oss_keys # 本次上传的临时音频 key(供事后清理) + self.oss_keys = oss_keys def build_direct_render( @@ -139,22 +139,28 @@ def build_direct_render( edge_crop_pct: float = 0.0, total_duration: float = 0.0, ) -> DirectRenderPlan: - """构造 P4000 直连渲染所需的 inputs 与 ffmpeg_args。 - - 视频:每段 trim/setpts/scale/pad/fps → concat(全硬切)→ 可选边缘 crop+scale → drawtext。 - 音频:concat 时丢弃原生音轨(只映射 [vfinal]),TTS/BGM 上传签名后 amix 混音。 - """ if not resolved_clips: raise ValueError("build_direct_render: no resolved clips") inputs: dict[str, str] = {} oss_keys: list[str] = [] input_args: list[str] = [] - fc: list[str] = [] # filter_complex 各段 - + fc: list[str] = [] n = len(resolved_clips) - # 1. 视频输入(原始素材签名 URL) + clip_starts: list[float] = [] + clip_effs: list[float] = [] + clip_speeds: list[float] = [] + for clip in resolved_clips: + start = float(getattr(clip, "start_time", 0) or 0) + eff = float(getattr(clip, "duration", 0) or 0) + if eff <= 0: + eff = float(getattr(clip, "actual_duration", 0) or 0) + speed = float(getattr(clip, "playback_speed", 1.0) or 1.0) + clip_starts.append(start) + clip_effs.append(eff) + clip_speeds.append(speed) + for i, clip in enumerate(resolved_clips): sk = (getattr(clip, "config", None) or {}).get("_storage_key") if not sk: @@ -163,52 +169,69 @@ def build_direct_render( inputs[fname] = sign_asset_url(sk) input_args.extend(["-i", fname]) - # 2. 每个视频段预处理 pre_labels: list[str] = [] - for i, clip in enumerate(resolved_clips): - filters: list[str] = [] - start = float(getattr(clip, "start_time", 0) or 0) - eff = float(getattr(clip, "duration", 0) or 0) - if eff <= 0: - eff = float(getattr(clip, "actual_duration", 0) or 0) + for i in range(n): + vf: list[str] = [] + start, eff, speed = clip_starts[i], clip_effs[i], clip_speeds[i] if eff > 0: if start > 0: - filters.append(f"trim=start={start:.3f}:duration={eff:.3f}") + vf.append(f"trim=start={start:.3f}:duration={eff:.3f}") else: - filters.append(f"trim=duration={eff:.3f}") - filters.append("setpts=PTS-STARTPTS") - - speed = float(getattr(clip, "playback_speed", 1.0) or 1.0) + vf.append(f"trim=duration={eff:.3f}") + vf.append("setpts=PTS-STARTPTS") if abs(speed - 1.0) >= 1e-6: - filters.append(f"setpts=PTS/{speed:.4f}") - - filters.append(f"scale={output_width}:{output_height}:force_original_aspect_ratio=decrease") - filters.append(f"pad={output_width}:{output_height}:trunc((ow-iw)/2):trunc((oh-ih)/2):black") - filters.append("setpts=PTS-STARTPTS") - filters.append(f"fps={output_fps}") - + vf.append(f"setpts=PTS/{speed:.4f}") + vf.append(f"scale={output_width}:{output_height}:force_original_aspect_ratio=decrease") + vf.append(f"pad={output_width}:{output_height}:trunc((ow-iw)/2):trunc((oh-ih)/2):black") + vf.append("setpts=PTS-STARTPTS") + vf.append(f"fps={output_fps}") label = f"vc{i}" - fc.append(f"[{i}:v]{','.join(filters)}[{label}]") + fc.append(f"[{i}:v]{','.join(vf)}[{label}]") pre_labels.append(label) - # 3. concat(全硬切;原生音频丢弃,v=1:a=0) - concat_in = "".join(f"[{_lbl}]" for _lbl in pre_labels) - fc.append(f"{concat_in}concat=n={n}:v=1:a=0[vcat]") - cur = "vcat" + audio_pre_labels: list[str] = [] + for i in range(n): + af: list[str] = [] + start, eff, speed = clip_starts[i], clip_effs[i], clip_speeds[i] + if eff > 0: + if start > 0: + af.append(f"atrim=start={start:.3f}:duration={eff:.3f}") + else: + af.append(f"atrim=duration={eff:.3f}") + af.append("asetpts=PTS-STARTPTS") + if abs(speed - 1.0) >= 1e-6: + atempo = _build_atempo_chain(speed) + if atempo: + af.append(atempo) + af.append("aresample=44100") + af.append("aformat=sample_fmts=fltp:channel_layouts=stereo") + alabel = f"ac{i}" + fc.append(f"[{i}:a]{','.join(af)}[{alabel}]") + audio_pre_labels.append(alabel) + + concat_in = "".join(f"[{v}][{a}]" for v, a in zip(pre_labels, audio_pre_labels, strict=True)) + fc.append(f"{concat_in}concat=n={n}:v=1:a=1[vcat][acat]") + cur_v = "vcat" + cur_a = "acat" - # 4. 边缘裁剪降重(合并进同一条,不再单独重编码) if edge_crop_pct and edge_crop_pct > 0: - p = float(edge_crop_pct) - keep = 1.0 - 2.0 * p - cw_expr = f"trunc(iw*{keep:.4f}/2)*2" - ch_expr = f"trunc(ih*{keep:.4f}/2)*2" + _r = random.Random() + p_min = EDGE_CROP_MIN_PCT + p_max = EDGE_CROP_MAX_PCT + crop_top = p_min + _r.random() * (p_max - p_min) + crop_bottom = p_min + _r.random() * (p_max - p_min) + crop_left = p_min + _r.random() * (p_max - p_min) + crop_right = p_min + _r.random() * (p_max - p_min) + w_expr = f"trunc(iw*(1-{crop_left:.4f}-{crop_right:.4f})/2)*2" + h_expr = f"trunc(ih*(1-{crop_top:.4f}-{crop_bottom:.4f})/2)*2" + x_expr = f"trunc(iw*{crop_left:.4f}/2)*2" + y_expr = f"trunc(ih*{crop_top:.4f}/2)*2" fc.append( - f"[{cur}]crop=w='{cw_expr}':h='{ch_expr}':x='(iw-{cw_expr})/2':y='(ih-{ch_expr})/2'," + f"[{cur_v}]crop=w='{w_expr}':h='{h_expr}':x='{x_expr}':y='{y_expr}'," f"scale={output_width}:{output_height}[vcrop]" ) - cur = "vcrop" + cur_v = "vcrop" - # 5. drawtext 字幕(标题整段 + ASR 逐句) draw_filters: list[str] = [] if title_text.strip(): title_size = max(int(output_height * 0.05), 24) @@ -241,18 +264,17 @@ def build_direct_render( ) if draw_filters: - prev = cur + prev = cur_v for idx, df in enumerate(draw_filters): out_l = "vfinal" if idx == len(draw_filters) - 1 else f"vd{idx}" fc.append(f"[{prev}]{df}[{out_l}]") prev = out_l vfinal_label = prev else: - fc.append(f"[{cur}]format=yuv420p[vfinal]") + fc.append(f"[{cur_v}]format=yuv420p[vfinal]") vfinal_label = "vfinal" - # 6. 音频输入与混音 - audio_items: list[tuple[int, float]] = [] # (input_index, volume) + mix_labels: list[str] = [cur_a] next_idx = n if tts_audio and Path(tts_audio).exists(): turl, tkey = upload_local_audio_and_sign(Path(tts_audio)) @@ -260,7 +282,11 @@ def build_direct_render( inputs[tname] = turl oss_keys.append(tkey) input_args.extend(["-i", tname]) - audio_items.append((next_idx, 1.0)) + alabel = "au_tts" + fc.append( + f"[{next_idx}:a]aresample=44100,volume=1.00,aformat=sample_fmts=fltp:channel_layouts=stereo[{alabel}]" + ) + mix_labels.append(alabel) next_idx += 1 if bgm_audio and Path(bgm_audio).exists(): burl, bkey = upload_local_audio_and_sign(Path(bgm_audio)) @@ -268,29 +294,35 @@ def build_direct_render( inputs[bname] = burl oss_keys.append(bkey) input_args.extend(["-i", bname]) - audio_items.append((next_idx, 0.35)) + alabel = "au_bgm" + fc.append( + f"[{next_idx}:a]aresample=44100,volume=0.35,aformat=sample_fmts=fltp:channel_layouts=stereo[{alabel}]" + ) + mix_labels.append(alabel) next_idx += 1 maps: list[str] = ["-map", f"[{vfinal_label}]"] - if audio_items: - mix_labels: list[str] = [] - for k, (idx, vol) in enumerate(audio_items): - alabel = f"au{k}" - fc.append( - f"[{idx}:a]aresample=44100,volume={vol:.2f},aformat=sample_fmts=fltp:channel_layouts=stereo[{alabel}]" - ) - mix_labels.append(alabel) - mix_in = "".join(f"[{_lbl}]" for _lbl in mix_labels) - fc.append(f"{mix_in}amix=inputs={len(mix_labels)}:duration=first:dropout_transition=2,aresample=44100[afinal]") + if mix_labels: + mix_in = "".join(f"[{lb}]" for lb in mix_labels) + n_mix = len(mix_labels) + mix_parts = [ + f"amix=inputs={n_mix}:duration=longest:dropout_transition=2:normalize=0", + "aresample=44100", + ] + if total_duration and total_duration > 0: + mix_parts.append(f"atrim=0:{total_duration:.3f}") + mix_parts.append("asetpts=PTS-STARTPTS") + fc.append(f"{mix_in}{','.join(mix_parts)}[afinal]") maps.extend(["-map", "[afinal]", "-c:a", "aac", "-b:a", "128k"]) + else: + logger.info("[gpu-direct] no audio tracks; output silent video") - # 7. 组装 ffmpeg_args + NVENC 编码 ffmpeg_args = ["-y", *input_args, "-filter_complex", ";".join(fc), *maps] ffmpeg_args.extend(["-c:v", vcodec, "-preset", preset, "-pix_fmt", "yuv420p"]) if video_bitrate: ffmpeg_args.extend(["-b:v", video_bitrate]) else: ffmpeg_args.extend(["-cq", str(cq)]) - ffmpeg_args.extend(["-movflags", "+faststart", "-f", "mp4", "pipe:1"]) + ffmpeg_args.extend(["-movflags", "+faststart", "-shortest", "-f", "mp4", "pipe:1"]) return DirectRenderPlan(inputs=inputs, ffmpeg_args=ffmpeg_args, oss_keys=oss_keys) diff --git a/apps/worker/video_processing/render_adapter.py b/apps/worker/video_processing/render_adapter.py index 512344952..68ca88619 100755 --- a/apps/worker/video_processing/render_adapter.py +++ b/apps/worker/video_processing/render_adapter.py @@ -86,6 +86,7 @@ class RenderAdapterResult: None # 封面候选帧 [{"image_url": "...", "frame_time": 5.0, "storage_key": "..."}] ) temp_dir: str | None = None # 渲染临时目录,成功时由调用方清理,失败时由 finally 清理 + edge_crop_applied: bool = False # GPU 管线已做随机边缘裁剪(跳过 CPU 二次重编码) def __post_init__(self): if self.rendered_clip_ids is None: @@ -605,6 +606,9 @@ class RenderAdapter: c.config["_storage_key"] = sk result = render_svc.render() + # 4.4 透传 GPU 直连路径的 edge_crop 状态(供外层跳过 CPU 二次裁剪) + edge_crop_applied_flag = bool(getattr(result, "edge_crop_applied", False)) + # 4.5 渲染后校验输出完整性 validation = validate_video_output(result.output_path) if not validation.valid: @@ -712,6 +716,7 @@ class RenderAdapter: rendered_clip_ids=final_rendered_ids, failed_clip_ids=final_failed_ids, cover_candidates=cover_candidates, + edge_crop_applied=edge_crop_applied_flag, ) def render_from_memory( diff --git a/apps/worker/video_processing/unified_render_service.py b/apps/worker/video_processing/unified_render_service.py index c5d4efb16..4001acf17 100755 --- a/apps/worker/video_processing/unified_render_service.py +++ b/apps/worker/video_processing/unified_render_service.py @@ -112,6 +112,7 @@ class RenderResult: file_size: int width: int height: int + edge_crop_applied: bool = False # True = GPU管线已做随机边缘裁剪 # ── clip_type → layer role 映射 ────────────────────────────────────────────── @@ -349,7 +350,8 @@ class UnifiedRenderService: video_duration=video_duration_final, output_path=output_path, ) - if direct_result is not None: + if direct_result is not None and direct_result[0]: + _direct_edge_crop = bool(direct_result[1]) # 直连成功:直接探测并返回,跳过后续视频/音频 CPU 流程 duration, file_size, width, height = self._probe_output(output_path) logger.info( @@ -360,12 +362,14 @@ class UnifiedRenderService: width, height, ) + direct_edge_cropped = _direct_edge_crop # GPU直连时若dedup=True已在GPU内做随机边缘裁剪 return RenderResult( output_path=output_path, duration=duration, file_size=file_size, width=width, height=height, + edge_crop_applied=direct_edge_cropped, ) # 5. 视频主渲染 @@ -2261,12 +2265,12 @@ class UnifiedRenderService: ass_path: Path | None, video_duration: float, output_path: Path, - ) -> bool | None: - """尝试全 GPU 直连渲染。成功返回 True,不支持/失败返回 None(调用方走旧链路)。""" + ) -> tuple[bool, bool] | tuple[None, bool]: + """尝试全 GPU 直连渲染。成功返回 (True, edge_crop_applied),不支持/失败返回 (None, False)。""" if not self._can_use_gpu_direct(layers): - return None + return (None, False) if not self._gpu_encode_available(): - return None + return (None, False) try: from video_processing import gpu_direct_pipeline as gdp @@ -2329,8 +2333,11 @@ class UnifiedRenderService: except Exception: # noqa: BLE001 pass - logger.info("[gpu-direct] success: plan_id=%s clips=%d", self.plan.id, len(video_clips)) - return True + did_edge_crop = bool(edge_pct) + logger.info( + "[gpu-direct] success: plan_id=%s clips=%d edge_crop=%s", self.plan.id, len(video_clips), did_edge_crop + ) + return (True, did_edge_crop) except GpuEncodeError as e: logger.warning("[gpu-direct] failed (fallback to legacy): %s", e) @@ -2339,10 +2346,10 @@ class UnifiedRenderService: output_path.unlink() except OSError: pass - return None + return (None, False) except Exception: # noqa: BLE001 logger.warning("[gpu-direct] unexpected error (fallback)", exc_info=True) - return None + return (None, False) def _concat_audio_clips(self, clips: list[Any], *, tag: str) -> Path: """把多个本地音频片段无间隙 concat 成一个 m4a(TTS 分段→单文件)。""" diff --git a/apps/worker/worker_app/tasks/generation.py b/apps/worker/worker_app/tasks/generation.py index e0da33384..5e2d8978d 100644 --- a/apps/worker/worker_app/tasks/generation.py +++ b/apps/worker/worker_app/tasks/generation.py @@ -694,11 +694,11 @@ def _render_from_edit_plan( task_id: str, source_edit_plan_id: str, task_info: dict, -) -> tuple[Path, float, list[dict] | None, str | None, str | None, str]: +) -> tuple[Path, float, list[dict] | None, str | None, str | None, str, bool]: """从 EditPlan 数据库记录直接渲染(不再内存重建clips)。 Returns: - (output_path, render_duration, cover_candidates, voiceover_path, temp_dir, thumbnail_url) + (output_path, render_duration, cover_candidates, voiceover_path, temp_dir, thumbnail_url, edge_crop_applied) """ from video_processing.render_adapter import RenderAdapter from worker_app.db import SessionLocal @@ -747,6 +747,7 @@ def _render_from_edit_plan( voiceover_path, render_temp_dir, result.thumbnail_url or "", + bool(getattr(result, "edge_crop_applied", False)), ) finally: db.close() @@ -900,6 +901,7 @@ def generate_video(self, task_id: str) -> dict: voiceover_tmp_path, render_temp_dir, thumbnail_url, + _gpu_edge_crop_done, ) = _render_from_edit_plan( task_id=task_id, source_edit_plan_id=current_plan_id, @@ -936,6 +938,12 @@ def generate_video(self, task_id: str) -> dict: if gen_task and render_attempt == 0: gen_task.append_log("降重", "已关闭边缘裁剪与微变换(确定性渲染)") _flush_logs(task_id, gen_task) + elif _gpu_edge_crop_done: + # GPU 直连管线已经在 filter_complex 中做了随机边缘裁剪,跳过 CPU 二次重编码 + if gen_task and render_attempt == 0: + gen_task.append_log("边缘裁剪", "已在 GPU 直连管线内完成随机边缘裁剪") + _flush_logs(task_id, gen_task) + logger.info("[task_id=%s] GPU直连已完成边缘裁剪,跳过CPU二次重编码", task_id) else: from video_processing.ffmpeg_utils import random_edge_crop