fix(worker): GPU直连管线音频bug修复+边缘裁剪入GPU省5s #2089

Merged
auto-approve-bot merged 1 commits from fix/gpu-pipeline-audio-and-edgecrop into develop 2026-09-28 20:17:55 +08:00
4 changed files with 145 additions and 93 deletions
@@ -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)
@@ -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(
@@ -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 分段→单文件)。"""
+10 -2
View File
@@ -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