From b04a8036559320a796052b40c4adfa70377d7779 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E4=B8=9A=E5=8A=A1=E6=9C=8D=E5=8A=A1=E5=99=A8=E8=BF=90?= =?UTF-8?q?=E7=BB=B4?= Date: Mon, 28 Sep 2026 13:30:31 +0800 Subject: [PATCH 1/3] =?UTF-8?q?feat(worker):=20=E5=85=A8GPU=E7=9B=B4?= =?UTF-8?q?=E8=BF=9E=E6=B8=B2=E6=9F=93=EF=BC=8C=E5=8F=96=E6=B6=88mezzanine?= =?UTF-8?q?=20CPU=E4=B8=AD=E9=97=B4=E7=BC=96=E7=A0=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - gpu_encoder 新增 render_inputs_to_output:多输入+filter_complex 直接交 P4000 一次出片 - 新增 gpu_direct_pipeline:原始素材/BGM 签名URL、TTS 上传、drawtext 标题/字幕、 filter_complex 完成 trim/scale/pad/concat/边缘crop/amix,末端 h264_nvenc 仅编码一次 - unified_render_service:render() 步骤4.8 接入直连,命中即出片; 不满足条件或 GPU 失败自动回退现有 mezzanine/CPU 链路;config gpu_direct_enabled 可灰度关闭 - render_adapter:_download_assets 额外返回 asset_storage_map, _do_render 给 clip 注入 _storage_key 供直连签名 消除 mezzanine 85s + 边缘裁剪 26s + 合成 2s + 中转 1s ≈ 114s CPU 开销 --- .../video_processing/gpu_direct_pipeline.py | 300 ++++++++++++++++++ .../worker/video_processing/render_adapter.py | 16 +- .../unified_render_service.py | 193 +++++++++++ packages/shared/gpu_encoder.py | 97 ++++++ 4 files changed, 603 insertions(+), 3 deletions(-) create mode 100644 apps/worker/video_processing/gpu_direct_pipeline.py diff --git a/apps/worker/video_processing/gpu_direct_pipeline.py b/apps/worker/video_processing/gpu_direct_pipeline.py new file mode 100644 index 000000000..35a8365d6 --- /dev/null +++ b/apps/worker/video_processing/gpu_direct_pipeline.py @@ -0,0 +1,300 @@ +"""全 GPU 直连渲染管线(P1)。 + +背景:旧链路 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 合成。 + +约束(P1): +- 仅覆盖智能剪辑主流场景:单一主视频轨、全硬切、无 PiP/overlay/水印/贴纸/片头片尾/绿幕。 + 不满足条件时调用方回退到现有 mezzanine/CPU 链路(功能不回归)。 +- 字幕先用 drawtext(P4000 装好中文字体后可再切 subtitles 滤镜烧 ASS)。 +""" + +from __future__ import annotations + +import logging +import uuid +from pathlib import Path +from typing import Any, Optional + +logger = logging.getLogger(__name__) + +# drawtext 默认字体名(fontconfig 解析);P4000 装好中文字体后可在 config 指定 +DEFAULT_DRAWTEXT_FONT = "sans" + + +def escape_drawtext_text(text: str) -> str: + """转义 drawtext text= 中的特殊字符(ffmpeg 过滤器语法)。""" + if not text: + return "" + # 顺序重要:先转义反斜杠本身 + s = text.replace("\\", "\\\\") + s = s.replace(":", "\\:") + s = s.replace("'", "\\'") + s = s.replace("%", "\\%") + s = s.replace(",", "\\,") + s = s.replace("[", "\\[").replace("]", "\\]") + s = s.replace(";", "\\;") + # 换行保留(drawtext 支持 %{...};真实换行需转成字面) + s = s.replace("\n", " ") + return s + + +def build_drawtext_filter( + *, + text: str, + start: float, + end: float, + font: str = DEFAULT_DRAWTEXT_FONT, + font_size: int = 0, + font_color: str = "white", + x_expr: str = "(w-text_w)/2", + y_expr: str = "h-th-60", + box: bool = False, + box_color: str = "black@0.5", + borderw: int = 0, + 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: + parts.append(f"fontsize={int(font_size)}") + parts.append(f"fontcolor={font_color}") + if box: + parts.append("box=1") + parts.append(f"boxcolor={box_color}") + if borderw and borderw > 0: + parts.append(f"borderw={int(borderw)}") + parts.append(f"bordercolor={border_color}") + parts.append(f"x={x_expr}") + parts.append(f"y={y_expr}") + if enable: + parts.append(f"enable='between(t,{start:.3f},{end:.3f})'") + return "drawtext=" + ":".join(parts) + + +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 + + storage = _storage() + key = f"{tmp_prefix.rstrip('/')}/{uuid.uuid4().hex}{local_audio.suffix or '.mp3'}" + content_type = "audio/mpeg" if local_audio.suffix.lower() in (".mp3", ".mpeg") else "audio/mp4" + storage.upload_file(local_audio, key, content_type=content_type) + url = storage.get_download_url(key, expires) + return url, key + + +def sign_asset_url(storage_key: str, *, expires: int = 3600) -> str: + """给原始素材 storage_key 签公网 GET URL(供 P4000 直接下载)。""" + from video_processing.oss_helpers import _storage # 类型: 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(供事后清理) + + +def build_direct_render( + *, + resolved_clips: list[Any], + output_width: int, + output_height: int, + output_fps: int, + tts_audio: Optional[Path] = None, + bgm_audio: Optional[Path] = None, + title_text: str = "", + subtitle_segments: Optional[list[Any]] = None, + font: str = DEFAULT_DRAWTEXT_FONT, + vcodec: str = "h264_nvenc", + preset: str = "p4", + video_bitrate: str = "", + cq: int = 23, + 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 各段 + + n = len(resolved_clips) + + # 1. 视频输入(原始素材签名 URL) + for i, clip in enumerate(resolved_clips): + sk = getattr(clip, "_storage_key", None) + if not sk: + raise ValueError(f"clip {getattr(clip,'clip_id',i)} missing _storage_key") + fname = f"v{i}.mp4" + 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) + if eff > 0: + if start > 0: + filters.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) + 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}") + + label = f"vc{i}" + fc.append(f"[{i}:v]{','.join(filters)}[{label}]") + pre_labels.append(label) + + # 3. concat(全硬切;原生音频丢弃,v=1:a=0) + concat_in = "".join(f"[{l}]" for l in pre_labels) + fc.append(f"{concat_in}concat=n={n}:v=1:a=0[vcat]") + cur = "vcat" + + # 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" + fc.append( + f"[{cur}]crop=w='{cw_expr}':h='{ch_expr}':x='(iw-{cw_expr})/2':y='(ih-{ch_expr})/2'," + f"scale={output_width}:{output_height}[vcrop]" + ) + cur = "vcrop" + + # 5. drawtext 字幕(标题整段 + ASR 逐句) + draw_filters: list[str] = [] + if title_text.strip(): + title_size = max(int(output_height * 0.05), 24) + draw_filters.append( + build_drawtext_filter( + text=title_text, + start=0.0, + end=max(total_duration, 0.1), + font=font, + font_size=title_size, + y_expr="h-th-40", + box=True, + ) + ) + sub_size = max(int(output_height * 0.045), 20) + for seg in subtitle_segments or []: + txt = getattr(seg, "text", "") or "" + if not txt.strip(): + continue + draw_filters.append( + build_drawtext_filter( + text=txt, + start=float(getattr(seg, "start", 0)), + end=float(getattr(seg, "end", 0)), + font=font, + font_size=sub_size, + y_expr="h-th-60", + borderw=2, + ) + ) + + if draw_filters: + prev = cur + 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]") + vfinal_label = "vfinal" + + # 6. 音频输入与混音 + audio_items: list[tuple[int, float]] = [] # (input_index, volume) + next_idx = n + if tts_audio and Path(tts_audio).exists(): + turl, tkey = upload_local_audio_and_sign(Path(tts_audio)) + tname = "tts" + (Path(tts_audio).suffix or ".mp3") + inputs[tname] = turl + oss_keys.append(tkey) + input_args.extend(["-i", tname]) + audio_items.append((next_idx, 1.0)) + next_idx += 1 + if bgm_audio and Path(bgm_audio).exists(): + burl, bkey = upload_local_audio_and_sign(Path(bgm_audio)) + bname = "bgm" + (Path(bgm_audio).suffix or ".mp3") + inputs[bname] = burl + oss_keys.append(bkey) + input_args.extend(["-i", bname]) + audio_items.append((next_idx, 0.35)) + 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}," + f"aformat=sample_fmts=fltp:channel_layouts=stereo[{alabel}]" + ) + mix_labels.append(alabel) + mix_in = "".join(f"[{l}]" for l in mix_labels) + fc.append( + f"{mix_in}amix=inputs={len(mix_labels)}:duration=first:dropout_transition=2," + f"aresample=44100[afinal]" + ) + maps.extend(["-map", "[afinal]", "-c:a", "aac", "-b:a", "128k"]) + + # 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"]) + + 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 e1f2924a5..bac7f36ad 100755 --- a/apps/worker/video_processing/render_adapter.py +++ b/apps/worker/video_processing/render_adapter.py @@ -189,7 +189,9 @@ class RenderAdapter: self._report_progress(progress_cb, 15.0, f"下载素材({len(ready_clips)} 个)") # 2. 下载素材 - asset_path_map, rendered_clip_ids, failed_clip_ids = self._download_assets(ready_clips, work_dir) + asset_path_map, rendered_clip_ids, failed_clip_ids, asset_storage_map = self._download_assets( + ready_clips, work_dir + ) if not asset_path_map: return RenderAdapterResult( success=False, @@ -206,6 +208,7 @@ class RenderAdapter: plan=plan, clips=ready_clips, asset_path_map=asset_path_map, + asset_storage_map=asset_storage_map, work_dir=work_dir, plan_id=plan_id, job_id=job_id, @@ -315,7 +318,7 @@ class RenderAdapter: def _download_assets( self, clips: list[EditPlanClip], work_dir: Path - ) -> tuple[dict[str, Path], list[str], list[str]]: + ) -> tuple[dict[str, Path], list[str], list[str], dict[str, str]]: """下载片段素材到本地。 先通过 asset_id 批量查询 assets 表获取 file_url(OSS存储路径), @@ -386,7 +389,7 @@ class RenderAdapter: failed_clip_ids.append(clip.id) logger.warning("素材下载失败: clip_id=%s asset_id=%s", clip.id, asset_id[:60]) - return asset_path_map, rendered_clip_ids, failed_clip_ids + return asset_path_map, rendered_clip_ids, failed_clip_ids, asset_storage_map def _prepare_bgm(self, plan, work_dir: Path, plan_id: str) -> str | None: """准备 BGM 音频文件(从 plan.config.bgm 读取配置)。 @@ -541,6 +544,7 @@ class RenderAdapter: rendered_clip_ids: list[str] | None = None, failed_clip_ids: list[str] | None = None, voiceover_audio_path: str | None = None, + asset_storage_map: dict[str, str] | None = None, ) -> RenderAdapterResult: """执行统一渲染核心流程(BGM + ASR + 渲染 + 缩略图 + 上传)。 @@ -590,6 +594,12 @@ class RenderAdapter: voiceover_audio_path=voiceover_audio_path, clip_has_text=clip_has_text, ) + # 注入每个视频段对应素材的 storage_key,供全 GPU 直连管线直接签名下载 + _storage_map = asset_storage_map or {} + for c in clips: + sk = _storage_map.get(getattr(c, "asset_id", "")) + if sk: + setattr(c, "_storage_key", sk) result = render_svc.render() # 4.5 渲染后校验输出完整性 diff --git a/apps/worker/video_processing/unified_render_service.py b/apps/worker/video_processing/unified_render_service.py index a92e3bbe8..a239e126f 100755 --- a/apps/worker/video_processing/unified_render_service.py +++ b/apps/worker/video_processing/unified_render_service.py @@ -341,6 +341,33 @@ class UnifiedRenderService: len(pip_sources), ) + # 4.8 全 GPU 直连管线(P1):命中主流场景则跳过 mezzanine/边缘裁剪 CPU 重编码 + output_path = self.work_dir / f"rendered_{self.plan.id}.mp4" + direct_result = self._try_gpu_direct( + layers=layers, + ass_path=ass_path, + video_duration=video_duration_final, + output_path=output_path, + ) + if direct_result is not None: + # 直连成功:直接探测并返回,跳过后续视频/音频 CPU 流程 + duration, file_size, width, height = self._probe_output(output_path) + logger.info( + "[unified-render] gpu-direct done: plan_id=%s total_ms=%d output_size=%d resolution=%dx%d", + self.plan.id, + int((time.time() - t_start) * 1000), + file_size, + width, + height, + ) + return RenderResult( + output_path=output_path, + duration=duration, + file_size=file_size, + width=width, + height=height, + ) + # 5. 视频主渲染 t_video_start = time.time() video_only_path = self.work_dir / f"rendered_{self.plan.id}_video.mp4" @@ -2180,6 +2207,172 @@ class UnifiedRenderService: # ── GPU NVENC 加速 ──────────────────────────────────────────────────── + # ── 全 GPU 直连渲染(P1)───────────────────────────────────────────── + + def _can_use_gpu_direct(self, layers: list[RenderLayer]) -> bool: + """判断是否命中直连支持的场景:单一主视频轨、全硬切、无复杂合成。""" + try: + cfg = self.plan.config or {} + # 特性开关(默认开启;可经 env/plan config 关闭灰度回退) + if not bool(cfg.get("gpu_direct_enabled", True)): + return False + + video_layers = [l for l in layers if l.role not in ("audio",)] + # 只允许一个视频层,且角色为主层 + if len(video_layers) != 1: + return False + role = video_layers[0].role + if role not in ("main", "broll"): + return False + + clips_v = [c for c in video_layers[0].clips if c.clip_type != "audio"] + if not clips_v: + return False + # 全硬切(第一个 clip 的转场忽略) + for c in clips_v[1:]: + te = c.transition_effect + if te not in (None, "", "cut"): + return False + # 无画中画 / 水印 / 贴纸 / 片头片尾 / 绿幕 / 倒放 / 调色 + if (cfg or {}).get("pip_config"): + return False + if (cfg or {}).get("intro_outro"): + return False + for c in clips_v: + cc = c.config or {} + if cc.get("watermark") or cc.get("stickers") or cc.get("chroma_key"): + return False + if ReverseConfig.from_dict(cc.get("reverse")).enabled: + return False + cg = ColorGradeConfig.from_dict(cc.get("color_grade")) + if cg.enabled and cg.has_effect(): + return False + if not getattr(c, "_storage_key", None): + return False + return True + except Exception: # noqa: BLE001 + logger.warning("[gpu-direct] eligibility check failed (fallback)", exc_info=True) + return False + + def _try_gpu_direct( + self, + *, + layers: list[RenderLayer], + ass_path: Path | None, + video_duration: float, + output_path: Path, + ) -> bool | None: + """尝试全 GPU 直连渲染。成功返回 True,不支持/失败返回 None(调用方走旧链路)。""" + if not self._can_use_gpu_direct(layers): + return None + if not self._gpu_encode_available(): + return None + + try: + from video_processing import gpu_direct_pipeline as gdp + + cfg = self.plan.config or {} + video_layer = next(l for l in layers if l.role not in ("audio",)) + video_clips = [c for c in video_layer.clips if c.clip_type != "audio"] + + # TTS:把 audio 层的 TTS 片段合并成一个文件给 P4000 + tts_merged: Path | None = None + audio_layer = next((l for l in layers if l.role == "audio"), None) + if audio_layer: + tts_clips = [ + c for c in audio_layer.clips if (c.config or {}).get("tts") and c.local_path.exists() + ] + if tts_clips: + tts_merged = self._concat_audio_clips(tts_clips, tag="tts_direct") + + # BGM 本地文件 + bgm_path = Path(self.bgm_path) if self.bgm_path else None + if bgm_path is not None and not bgm_path.exists(): + bgm_path = None + + # 字幕:标题 + ASR 时间轴 + title_cfg = cfg.get("title", {}) or cfg.get("title_config", {}) or {} + title_text = "" + if isinstance(title_cfg, dict) and title_cfg.get("enabled", True): + title_text = title_cfg.get("text", "") or "" + + subtitle_segments: list[Any] = [] + sub_cfg = cfg.get("subtitle", {}) or {} + if isinstance(sub_cfg, dict) and sub_cfg.get("enabled", True): + if sub_cfg.get("auto_generated") and self._asr_timeline_cache is not None: + subtitle_segments = list(self._asr_timeline_cache.segments) + + # 边缘裁剪比例(与 random_edge_crop 默认 2~5% 同口径,取固定 3%) + dedup = self._dedup_enabled() + edge_pct = 0.03 if dedup else 0.0 + + plan = gdp.build_direct_render( + resolved_clips=video_clips, + output_width=self.output_width, + output_height=self.output_height, + output_fps=self.output_fps, + tts_audio=tts_merged, + bgm_audio=bgm_path, + title_text=title_text, + subtitle_segments=subtitle_segments, + edge_crop_pct=edge_pct, + total_duration=video_duration, + ) + + client = get_gpu_encoder() + client.render_inputs_to_output(plan.inputs, plan.ffmpeg_args, output_path) + + # 清理本次上传的临时音频 + for key in plan.oss_keys: + try: + from video_processing.oss_helpers import _storage + + _storage().delete_file(key) if hasattr(_storage(), "delete_file") else None + except Exception: # noqa: BLE001 + pass + + logger.info("[gpu-direct] success: plan_id=%s clips=%d", self.plan.id, len(video_clips)) + return True + + except GpuEncodeError as e: + logger.warning("[gpu-direct] failed (fallback to legacy): %s", e) + try: + if output_path.exists(): + output_path.unlink() + except OSError: + pass + return None + except Exception: # noqa: BLE001 + logger.warning("[gpu-direct] unexpected error (fallback)", exc_info=True) + return None + + def _concat_audio_clips(self, clips: list[Any], *, tag: str) -> Path: + """把多个本地音频片段无间隙 concat 成一个 m4a(TTS 分段→单文件)。""" + out = self.work_dir / f"{tag}_{self.plan.id}.m4a" + listfile = self.work_dir / f"{tag}_{self.plan.id}.txt" + lines = [] + for c in clips: + ap = str(c.local_path).replace("'", "'\\''") + lines.append(f"file '{ap}'") + listfile.write_text("\n".join(lines), encoding="utf-8") + cmd = [ + FFMPEG_BIN, + "-y", + "-f", + "concat", + "-safe", + "0", + "-i", + str(listfile), + "-c:a", + "aac", + "-b:a", + "128k", + str(out), + ] + run_ffmpeg(cmd) + return out + def _gpu_encode_available(self) -> bool: """GPU 编码客户端是否已配置且健康(缓存健康状态,单任务内只探测一次)。""" if not getattr(self, "_gpu_health_ok", None): diff --git a/packages/shared/gpu_encoder.py b/packages/shared/gpu_encoder.py index 601da431b..77443d2c9 100644 --- a/packages/shared/gpu_encoder.py +++ b/packages/shared/gpu_encoder.py @@ -309,6 +309,103 @@ class GpuEncoderClient: except Exception as e: # noqa: BLE001 logger.warning("[gpu-encoder] failed to delete OSS mezzanine %s: %s", oss_key, e) + # ------------------------------------------------------------------ + # High-level: render arbitrary inputs → final output (all-GPU pipeline) + # ------------------------------------------------------------------ + def render_inputs_to_output( + self, + inputs: dict[str, str], + ffmpeg_args: list[str], + output_path: Path, + *, + timeout: Optional[int] = None, + ) -> dict[str, Any]: + """把多输入(原始素材/字幕/BGM)连同完整 filter_complex 交给 P4000 一次出片。 + + 与 encode_mezzanine_to_output 的区别:worker 侧不再生成/上传 mezzanine, + P4000 直接从 inputs 中的签名 URL 下载原始素材,filter_complex 内完成 + concat/scale/crop/drawtext/amix,末端 h264_nvenc 只编码一次。 + + 传输:成片仍走 relay 回传(P4000 PUT → worker GET),避免公网 OSS 往返。 + + Args: + inputs: {裸文件名: 可下载URL},key 即 ffmpeg_args 中引用的文件名 + ffmpeg_args: 完整 ffmpeg 参数(含 -i、-filter_complex、-map、NVENC 编码参数) + output_path: worker 本地成片落盘路径 + timeout: P4000 侧超时(秒) + """ + if not inputs: + raise GpuEncodeError("render_inputs_to_output: inputs is empty") + if not ffmpeg_args: + raise GpuEncodeError("render_inputs_to_output: ffmpeg_args is empty") + if not self.relay_base_url: + raise GpuEncodeError("gpu_encode_relay_base_url not configured") + + timeout = timeout or self.sync_timeout + t_total = time.time() + result_key: Optional[str] = None + + try: + secret = self._get_relay_secret() + + # 1. result relay URLs(成片 P4000 PUT → worker GET) + result_key = uuid.uuid4().hex + put_url = self._result_put_url(result_key, secret) + get_url = self._result_get_url(result_key, secret) + del_result_url = get_url + + # 2. pre-warm then call P4000 sync render + self._warm_up_if_needed() + body = { + "inputs": dict(inputs), + "ffmpeg_args": list(ffmpeg_args), + "output_url": put_url, + "timeout": int(timeout), + } + job = self._post_sync(body) + self._last_ok_ts = time.time() + logger.info( + "[gpu-encoder] P4000 direct done: job_id=%s rc=%s size=%s dur=%ss inputs=%d", + job.get("job_id"), + job.get("ffmpeg_rc"), + job.get("size"), + job.get("duration"), + len(inputs), + ) + + # 3. download result from relay to output_path + output_path.parent.mkdir(parents=True, exist_ok=True) + size = self._download_to_file(get_url, output_path) + + # 4. cleanup relay result + self._relay_delete(del_result_url) + + logger.info( + "[gpu-encoder] direct render ok → %s (%d bytes) total=%.2fs", + output_path.name, + size, + time.time() - t_total, + ) + return { + "job": job, + "output_size": size, + "output_path": str(output_path), + "transport": "direct", + } + + except GpuEncodeError: + raise + except Exception as e: # noqa: BLE001 + raise GpuEncodeError(f"unexpected: {e}") from e + finally: + # cleanup relay result (best-effort) + if result_key: + try: + secret = self._get_relay_secret() + self._relay_delete(self._relay_result_url(self.relay_internal_base_url, result_key, secret)) + except Exception as e: # noqa: BLE001 + logger.warning("[gpu-encoder] failed to delete relay result %s: %s", result_key, e) + # ------------------------------------------------------------------ # Internal helpers # ------------------------------------------------------------------ -- 2.54.0 From 1e23a3f094321958c9c6135550f1b11782323333 Mon Sep 17 00:00:00 2001 From: CI Bot Date: Mon, 28 Sep 2026 05:38:30 +0000 Subject: [PATCH 2/3] style: auto-format with black + isort + ruff + prettier [skip ci-format-check] --- apps/worker/video_processing/gpu_direct_pipeline.py | 3 +-- apps/worker/video_processing/render_adapter.py | 2 +- apps/worker/video_processing/unified_render_service.py | 4 +--- 3 files changed, 3 insertions(+), 6 deletions(-) diff --git a/apps/worker/video_processing/gpu_direct_pipeline.py b/apps/worker/video_processing/gpu_direct_pipeline.py index 35a8365d6..a7b70c2e1 100644 --- a/apps/worker/video_processing/gpu_direct_pipeline.py +++ b/apps/worker/video_processing/gpu_direct_pipeline.py @@ -283,8 +283,7 @@ def build_direct_render( mix_labels.append(alabel) mix_in = "".join(f"[{l}]" for l in mix_labels) fc.append( - f"{mix_in}amix=inputs={len(mix_labels)}:duration=first:dropout_transition=2," - f"aresample=44100[afinal]" + f"{mix_in}amix=inputs={len(mix_labels)}:duration=first:dropout_transition=2," f"aresample=44100[afinal]" ) maps.extend(["-map", "[afinal]", "-c:a", "aac", "-b:a", "128k"]) diff --git a/apps/worker/video_processing/render_adapter.py b/apps/worker/video_processing/render_adapter.py index bac7f36ad..2c11ee7ee 100755 --- a/apps/worker/video_processing/render_adapter.py +++ b/apps/worker/video_processing/render_adapter.py @@ -599,7 +599,7 @@ class RenderAdapter: for c in clips: sk = _storage_map.get(getattr(c, "asset_id", "")) if sk: - setattr(c, "_storage_key", sk) + c._storage_key = sk result = render_svc.render() # 4.5 渲染后校验输出完整性 diff --git a/apps/worker/video_processing/unified_render_service.py b/apps/worker/video_processing/unified_render_service.py index a239e126f..91ef61abe 100755 --- a/apps/worker/video_processing/unified_render_service.py +++ b/apps/worker/video_processing/unified_render_service.py @@ -2279,9 +2279,7 @@ class UnifiedRenderService: tts_merged: Path | None = None audio_layer = next((l for l in layers if l.role == "audio"), None) if audio_layer: - tts_clips = [ - c for c in audio_layer.clips if (c.config or {}).get("tts") and c.local_path.exists() - ] + tts_clips = [c for c in audio_layer.clips if (c.config or {}).get("tts") and c.local_path.exists()] if tts_clips: tts_merged = self._concat_audio_clips(tts_clips, tag="tts_direct") -- 2.54.0 From 6e8199581de9349617f5673d0147e57f14987b4c Mon Sep 17 00:00:00 2001 From: saas-backend-agent Date: Mon, 28 Sep 2026 14:02:57 +0800 Subject: [PATCH 3/3] =?UTF-8?q?fix(gpu-direct):=20=E4=BF=AE=E5=A4=8D?= =?UTF-8?q?=E5=8D=95=E6=B5=8B+style=20E741=EF=BC=8C=E9=BB=98=E8=AE=A4?= =?UTF-8?q?=E5=AD=97=E4=BD=93=E6=94=B9=E4=B8=BANoto=20Sans=20CJK=20SC?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - DEFAULT_DRAWTEXT_FONT 从 sans 改为 'Noto Sans CJK SC'(P4000已装中文字体) - _download_assets 返回 4-tuple 导致老单测解包失败,更新 test_render_adapter - 修复 E741 模糊变量名 l → _lbl/_lyr - ruff format 统一格式 --- .../video_processing/gpu_direct_pipeline.py | 15 ++++++--------- .../video_processing/unified_render_service.py | 10 +++++----- tests/unit/test_render_adapter.py | 14 +++++++------- 3 files changed, 18 insertions(+), 21 deletions(-) diff --git a/apps/worker/video_processing/gpu_direct_pipeline.py b/apps/worker/video_processing/gpu_direct_pipeline.py index a7b70c2e1..6f774da5e 100644 --- a/apps/worker/video_processing/gpu_direct_pipeline.py +++ b/apps/worker/video_processing/gpu_direct_pipeline.py @@ -22,7 +22,7 @@ from typing import Any, Optional logger = logging.getLogger(__name__) # drawtext 默认字体名(fontconfig 解析);P4000 装好中文字体后可在 config 指定 -DEFAULT_DRAWTEXT_FONT = "sans" +DEFAULT_DRAWTEXT_FONT = "Noto Sans CJK SC" def escape_drawtext_text(text: str) -> str: @@ -158,7 +158,7 @@ def build_direct_render( for i, clip in enumerate(resolved_clips): sk = getattr(clip, "_storage_key", None) if not sk: - raise ValueError(f"clip {getattr(clip,'clip_id',i)} missing _storage_key") + raise ValueError(f"clip {getattr(clip, 'clip_id', i)} missing _storage_key") fname = f"v{i}.mp4" inputs[fname] = sign_asset_url(sk) input_args.extend(["-i", fname]) @@ -192,7 +192,7 @@ def build_direct_render( pre_labels.append(label) # 3. concat(全硬切;原生音频丢弃,v=1:a=0) - concat_in = "".join(f"[{l}]" for l in pre_labels) + concat_in = "".join(f"[{l}]" for _lbl in pre_labels) fc.append(f"{concat_in}concat=n={n}:v=1:a=0[vcat]") cur = "vcat" @@ -277,14 +277,11 @@ def build_direct_render( for k, (idx, vol) in enumerate(audio_items): alabel = f"au{k}" fc.append( - f"[{idx}:a]aresample=44100,volume={vol:.2f}," - f"aformat=sample_fmts=fltp:channel_layouts=stereo[{alabel}]" + f"[{idx}:a]aresample=44100,volume={vol:.2f},aformat=sample_fmts=fltp:channel_layouts=stereo[{alabel}]" ) mix_labels.append(alabel) - mix_in = "".join(f"[{l}]" for l in mix_labels) - fc.append( - f"{mix_in}amix=inputs={len(mix_labels)}:duration=first:dropout_transition=2," f"aresample=44100[afinal]" - ) + mix_in = "".join(f"[{l}]" for _lbl in mix_labels) + fc.append(f"{mix_in}amix=inputs={len(mix_labels)}:duration=first:dropout_transition=2,aresample=44100[afinal]") maps.extend(["-map", "[afinal]", "-c:a", "aac", "-b:a", "128k"]) # 7. 组装 ffmpeg_args + NVENC 编码 diff --git a/apps/worker/video_processing/unified_render_service.py b/apps/worker/video_processing/unified_render_service.py index 91ef61abe..f052ee053 100755 --- a/apps/worker/video_processing/unified_render_service.py +++ b/apps/worker/video_processing/unified_render_service.py @@ -233,7 +233,7 @@ class UnifiedRenderService: return if abs(mt.brightness) > 1e-4 or abs(mt.contrast - 1.0) > 1e-4 or abs(mt.saturation - 1.0) > 1e-4: filters.append( - f"eq=brightness={mt.brightness:+.4f}:" f"contrast={mt.contrast:.4f}:saturation={mt.saturation:.4f}" + f"eq=brightness={mt.brightness:+.4f}:contrast={mt.contrast:.4f}:saturation={mt.saturation:.4f}" ) @staticmethod @@ -2217,7 +2217,7 @@ class UnifiedRenderService: if not bool(cfg.get("gpu_direct_enabled", True)): return False - video_layers = [l for l in layers if l.role not in ("audio",)] + video_layers = [_lyr for _lyr in layers if l.role not in ("audio",)] # 只允许一个视频层,且角色为主层 if len(video_layers) != 1: return False @@ -2272,12 +2272,12 @@ class UnifiedRenderService: from video_processing import gpu_direct_pipeline as gdp cfg = self.plan.config or {} - video_layer = next(l for l in layers if l.role not in ("audio",)) + video_layer = next(_lyr for _lyr in layers if l.role not in ("audio",)) video_clips = [c for c in video_layer.clips if c.clip_type != "audio"] # TTS:把 audio 层的 TTS 片段合并成一个文件给 P4000 tts_merged: Path | None = None - audio_layer = next((l for l in layers if l.role == "audio"), None) + audio_layer = next((_lyr for _lyr in layers if l.role == "audio"), None) if audio_layer: tts_clips = [c for c in audio_layer.clips if (c.config or {}).get("tts") and c.local_path.exists()] if tts_clips: @@ -2785,7 +2785,7 @@ class UnifiedRenderService: b = pixel_pert.get("color_b", 0) if r != 0 or g != 0 or b != 0: # color_balance 参数范围 -1.0 ~ 1.0,这里用 /100 转换 - filters.append(f"colorbalance=rs={r/100:.3f}:gs={g/100:.3f}:bs={b/100:.3f}") + filters.append(f"colorbalance=rs={r / 100:.3f}:gs={g / 100:.3f}:bs={b / 100:.3f}") @staticmethod def _clip_volume(clip: ResolvedClip) -> float: diff --git a/tests/unit/test_render_adapter.py b/tests/unit/test_render_adapter.py index e4d5b930b..502143479 100755 --- a/tests/unit/test_render_adapter.py +++ b/tests/unit/test_render_adapter.py @@ -666,7 +666,7 @@ class TestDownloadAssets: _make_clip("c2", order=1, asset_id="asset_002"), ] - asset_path_map, rendered_ids, failed_ids = adapter._download_assets(clips, tmp_path) + asset_path_map, rendered_ids, failed_ids, _storage_map = adapter._download_assets(clips, tmp_path) assert len(asset_path_map) == 2 assert "asset_001" in asset_path_map @@ -693,7 +693,7 @@ class TestDownloadAssets: _make_clip("c2", order=1, asset_id="asset_002"), ] - asset_path_map, rendered_ids, failed_ids = adapter._download_assets(clips, tmp_path) + asset_path_map, rendered_ids, failed_ids, _storage_map = adapter._download_assets(clips, tmp_path) assert len(asset_path_map) == 1 assert "asset_002" in asset_path_map @@ -714,7 +714,7 @@ class TestDownloadAssets: _make_clip("c1", order=0, asset_id="asset_001"), ] - asset_path_map, rendered_ids, failed_ids = adapter._download_assets(clips, tmp_path) + asset_path_map, rendered_ids, failed_ids, _storage_map = adapter._download_assets(clips, tmp_path) assert len(asset_path_map) == 0 assert len(rendered_ids) == 0 @@ -750,7 +750,7 @@ class TestDownloadAssets: _make_clip("c3", order=2, asset_id="asset_003"), ] - asset_path_map, rendered_ids, failed_ids = adapter._download_assets(clips, tmp_path) + asset_path_map, rendered_ids, failed_ids, _storage_map = adapter._download_assets(clips, tmp_path) assert len(asset_path_map) == 2 assert "c1" in rendered_ids @@ -771,7 +771,7 @@ class TestDownloadAssets: _make_clip("c2", order=1, asset_id="asset_shared"), ] - asset_path_map, rendered_ids, failed_ids = adapter._download_assets(clips, tmp_path) + asset_path_map, rendered_ids, failed_ids, _storage_map = adapter._download_assets(clips, tmp_path) assert len(asset_path_map) == 1 assert mock_download.call_count == 1 @@ -795,7 +795,7 @@ class TestDownloadAssets: adapter = RenderAdapter(mock_db) clips = [_make_clip("c1", order=0, asset_id="asset_fallback")] - asset_path_map, rendered_ids, failed_ids = adapter._download_assets(clips, tmp_path) + asset_path_map, rendered_ids, failed_ids, _storage_map = adapter._download_assets(clips, tmp_path) assert len(asset_path_map) == 1 assert "c1" in rendered_ids @@ -820,7 +820,7 @@ class TestDownloadAssets: adapter = RenderAdapter(mock_db) clips = [_make_clip("c1", order=0, asset_id="asset_no_key")] - asset_path_map, rendered_ids, failed_ids = adapter._download_assets(clips, tmp_path) + asset_path_map, rendered_ids, failed_ids, _storage_map = adapter._download_assets(clips, tmp_path) assert len(asset_path_map) == 0 assert "c1" in failed_ids -- 2.54.0