Compare commits

...

2 Commits

Author SHA1 Message Date
CI Bot 1ee120653e feat: 接入渲染灰度观测指标到compose_video和stream_copy
- compose_video 入口接入 render_task_metrics 上下文
  自动记录:引擎选择、任务耗时、成功/失败、错误类型、活跃任务数
- unified_render_service 接入 stream_copy 埋点
  三种状态:hit(命中直通)/ miss(条件不满足跳过)/ fallback(失败回退)
  miss 原因归一化为8类:编码/分辨率/帧率/字幕/像素格式/裁剪/多片段/其他
2026-07-13 14:39:15 +08:00
CI Bot f481bcd396 feat: 新增渲染引擎灰度观测指标模块
新增 render_metrics.py,提供 Prometheus 指标埋点:
- render_task_total: 任务总数(按引擎/任务类型/状态)
- render_task_duration_seconds: 耗时直方图(12个bucket,1s~2h)
- render_error_total: 错误统计(按引擎/任务类型/错误类型)
- render_stream_copy_total: stream_copy命中率(hit/miss/fallback)
- render_active_tasks: 活跃任务数(Gauge)

配套 9 个单元测试,覆盖成功/失败/超时/活跃数增减/stream_copy/错误分类等场景
2026-07-13 14:37:19 +08:00
4 changed files with 441 additions and 4 deletions
+214
View File
@@ -0,0 +1,214 @@
"""渲染引擎灰度观测指标。
提供 Prometheus 指标埋点,用于灰度发布期间观测新旧引擎的:
- 任务成功率
- 耗时分布
- 错误类型分布
- stream_copy 命中率
使用方式(任务入口):
with render_task_metrics(engine="unified", task_type="compose_video"):
# 执行渲染任务
result = do_render()
stream_copy 埋点:
record_stream_copy(result="hit", reason="成功")
record_stream_copy(result="miss", reason="编码不匹配")
record_stream_copy(result="fallback", reason="失败回退")
Worker 进程启动时需调用 start_metrics_server(port) 暴露 /metrics 端点。
"""
from __future__ import annotations
import logging
import threading
import time
from contextlib import contextmanager
from typing import Iterator, Optional
from prometheus_client import REGISTRY, Counter, Gauge, Histogram, start_http_server
logger = logging.getLogger(__name__)
# ── 指标定义 ──────────────────────────────────────────────────────────────────
# 渲染任务总数(按引擎、任务类型、状态区分)
RENDER_TASK_TOTAL = Counter(
"render_task_total",
"Total number of render tasks",
["engine", "task_type", "status"],
registry=REGISTRY,
)
# 渲染任务耗时直方图(按引擎、任务类型区分)
# buckets 覆盖从秒级到小时级,适配短视频到长视频的渲染场景
RENDER_TASK_DURATION_SECONDS = Histogram(
"render_task_duration_seconds",
"Render task duration in seconds",
["engine", "task_type"],
buckets=(1, 5, 10, 30, 60, 120, 300, 600, 900, 1800, 3600, 7200),
registry=REGISTRY,
)
# 渲染错误统计(按引擎、任务类型、错误类型区分)
RENDER_ERROR_TOTAL = Counter(
"render_error_total",
"Total number of render errors",
["engine", "task_type", "error_type"],
registry=REGISTRY,
)
# stream_copy 命中统计(仅新引擎有意义)
RENDER_STREAM_COPY_TOTAL = Counter(
"render_stream_copy_total",
"Total number of stream copy attempts and results",
["result", "reason"],
registry=REGISTRY,
)
# 当前正在执行的渲染任务数
RENDER_ACTIVE_TASKS = Gauge(
"render_active_tasks",
"Number of render tasks currently in progress",
["engine", "task_type"],
registry=REGISTRY,
)
# ── Metrics Server ───────────────────────────────────────────────────────────
_metrics_server_started = False
_metrics_server_lock = threading.Lock()
def start_metrics_server(port: int = 9101) -> None:
"""启动 Prometheus metrics HTTP 服务器。
在 worker 进程启动时调用一次即可,多进程环境下每个 worker 进程
会启动自己的 metrics server(需配置不同端口或使用进程号偏移)。
Args:
port: metrics 服务端口,默认 9101
"""
global _metrics_server_started
with _metrics_server_lock:
if _metrics_server_started:
return
try:
start_http_server(port)
_metrics_server_started = True
logger.info("Render metrics server started on port %d", port)
except OSError as e:
# 端口已占用可能是多 worker 进程场景,记录警告不阻断
logger.warning("Failed to start metrics server on port %d: %s", port, e)
# ── 任务级埋点 ───────────────────────────────────────────────────────────────
@contextmanager
def render_task_metrics(engine: str, task_type: str) -> Iterator[None]:
"""渲染任务指标上下文管理器。
自动记录:任务开始(活跃数+1)、任务结束(耗时 + 状态 + 活跃数-1)。
Args:
engine: 渲染引擎类型,"legacy""unified"
task_type: 任务类型,"compose_video" / "edit_plan" / "generate_video"
Usage:
with render_task_metrics(engine="unified", task_type="compose_video"):
result = do_render()
# 正常退出 = success
# 抛异常 = failure(会记录 error_type
"""
RENDER_ACTIVE_TASKS.labels(engine=engine, task_type=task_type).inc()
start_time = time.perf_counter()
status = "success"
try:
yield
except Exception as e:
status = "failure"
error_type = _classify_error(e)
RENDER_ERROR_TOTAL.labels(
engine=engine,
task_type=task_type,
error_type=error_type,
).inc()
raise
finally:
duration = time.perf_counter() - start_time
RENDER_TASK_TOTAL.labels(
engine=engine,
task_type=task_type,
status=status,
).inc()
RENDER_TASK_DURATION_SECONDS.labels(
engine=engine,
task_type=task_type,
).observe(duration)
RENDER_ACTIVE_TASKS.labels(engine=engine, task_type=task_type).dec()
def record_stream_copy(result: str, reason: str) -> None:
"""记录 stream_copy 命中/跳过/回退情况。
Args:
result: "hit"(命中直通) / "miss"(条件不满足跳过) / "fallback"(失败回退)
reason: 具体原因,如"编码不匹配""分辨率不同""成功""ffmpeg失败"
"""
RENDER_STREAM_COPY_TOTAL.labels(result=result, reason=reason).inc()
# ── 辅助函数 ─────────────────────────────────────────────────────────────────
def _classify_error(exc: BaseException) -> str:
"""将异常分类为标准错误类型。
用于 error_type 标签,控制指标基数不要爆炸。
"""
import subprocess
if isinstance(exc, subprocess.TimeoutExpired):
return "timeout"
if isinstance(exc, subprocess.CalledProcessError):
return "ffmpeg_error"
if isinstance(exc, ValueError):
return "validation"
if isinstance(exc, (OSError, IOError)):
return "io_error"
exc_name = type(exc).__name__
# 常见的 OSS / 网络相关异常
if "oss" in exc_name.lower() or "storage" in exc_name.lower():
return "oss_error"
if "timeout" in exc_name.lower():
return "timeout"
return "unknown"
def classify_stream_copy_miss_reason(reason: str) -> str:
"""将 stream_copy 未命中原因归一化为标准分类。
控制 reason 标签基数,避免爆炸。
"""
if "编码" in reason or "codec" in reason.lower():
return "编码不匹配"
if "分辨率" in reason or "width" in reason.lower() or "height" in reason.lower():
return "分辨率不匹配"
if "帧率" in reason or "fps" in reason.lower():
return "帧率不匹配"
if "字幕" in reason or "ass" in reason.lower():
return "有字幕叠加"
if "像素格式" in reason or "pix_fmt" in reason.lower():
return "像素格式不匹配"
if "trim" in reason or "时长" in reason:
return "裁剪不支持"
if "多clip" in reason or "多片段" in reason or "转场" in reason:
return "多片段不支持"
return "其他"
@@ -728,6 +728,10 @@ class UnifiedRenderService:
self.plan.id,
reason,
)
# 灰度观测:stream_copy 未命中
from video_processing.render_metrics import classify_stream_copy_miss_reason, record_stream_copy
miss_reason = classify_stream_copy_miss_reason(reason)
record_stream_copy("miss", miss_reason)
return False
# 构建 copy 命令
@@ -784,9 +788,15 @@ class UnifiedRenderService:
self.plan.id,
output_path.stat().st_size,
)
# 灰度观测:stream_copy 命中成功
from video_processing.render_metrics import record_stream_copy
record_stream_copy("hit", "成功")
return True
else:
logger.warning("[unified-render] stream_copy 输出为空: plan_id=%s", self.plan.id)
# 灰度观测:stream_copy 失败回退(输出为空)
from video_processing.render_metrics import record_stream_copy
record_stream_copy("fallback", "输出为空")
return False
except (subprocess.CalledProcessError, subprocess.TimeoutExpired) as e:
logger.warning(
@@ -794,6 +804,9 @@ class UnifiedRenderService:
self.plan.id,
str(e)[:200],
)
# 灰度观测:stream_copy 失败回退(ffmpeg错误)
from video_processing.render_metrics import record_stream_copy
record_stream_copy("fallback", "ffmpeg失败")
# 清理可能的损坏输出文件
if output_path.exists():
try:
+19 -4
View File
@@ -63,15 +63,30 @@ def compose_video(self, job_id: str, **kwargs):
# 判断使用哪个渲染引擎
# 优先级:Redis Feature Flag(白名单 > 百分比) > 环境变量默认
from video_processing.render_engine_resolver import get_render_engine_resolver
from video_processing.render_metrics import render_task_metrics
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)
# 灰度观测指标埋点
with render_task_metrics(engine=engine, task_type="compose_video"):
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)
except self.retry_exc as exc:
logger.warning("视频合成重试中: job_id=%s, exc=%s", job_id, exc)
+195
View File
@@ -0,0 +1,195 @@
"""渲染引擎灰度观测指标单元测试。"""
from __future__ import annotations
import subprocess
import time
import pytest
from prometheus_client import REGISTRY, CollectorRegistry, Counter, Gauge, Histogram
@pytest.fixture(autouse=True)
def _isolate_registry(monkeypatch):
"""每个测试使用独立的 CollectorRegistry,避免互相影响。"""
from video_processing import render_metrics
test_registry = CollectorRegistry()
# 用测试 registry 重新创建所有指标
task_total = Counter(
"render_task_total",
"Total number of render tasks",
["engine", "task_type", "status"],
registry=test_registry,
)
task_duration = Histogram(
"render_task_duration_seconds",
"Render task duration in seconds",
["engine", "task_type"],
buckets=(1, 5, 10, 30, 60, 120, 300, 600, 900, 1800, 3600, 7200),
registry=test_registry,
)
error_total = Counter(
"render_error_total",
"Total number of render errors",
["engine", "task_type", "error_type"],
registry=test_registry,
)
stream_copy_total = Counter(
"render_stream_copy_total",
"Total number of stream copy attempts and results",
["result", "reason"],
registry=test_registry,
)
active_tasks = Gauge(
"render_active_tasks",
"Number of render tasks currently in progress",
["engine", "task_type"],
registry=test_registry,
)
# 替换模块中的指标对象
monkeypatch.setattr(render_metrics, "RENDER_TASK_TOTAL", task_total)
monkeypatch.setattr(render_metrics, "RENDER_TASK_DURATION_SECONDS", task_duration)
monkeypatch.setattr(render_metrics, "RENDER_ERROR_TOTAL", error_total)
monkeypatch.setattr(render_metrics, "RENDER_STREAM_COPY_TOTAL", stream_copy_total)
monkeypatch.setattr(render_metrics, "RENDER_ACTIVE_TASKS", active_tasks)
yield test_registry
def _counter_value(counter: Counter, **labels) -> float:
"""获取 Counter 指定标签的当前值。"""
return counter.labels(**labels)._value.get()
def _gauge_value(gauge: Gauge, **labels) -> float:
"""获取 Gauge 指定标签的当前值。"""
return gauge.labels(**labels)._value.get()
def _histogram_sum(histogram: Histogram, **labels) -> float:
"""获取 Histogram 指定标签的观测总和。"""
return histogram.labels(**labels)._sum.get()
def test_render_task_metrics_success():
"""成功任务应记录 success 状态 + 耗时 + 活跃数正确增减。"""
from video_processing.render_metrics import (
RENDER_ACTIVE_TASKS,
RENDER_TASK_DURATION_SECONDS,
RENDER_TASK_TOTAL,
render_task_metrics,
)
assert _counter_value(RENDER_TASK_TOTAL, engine="unified", task_type="compose_video", status="success") == 0
assert _gauge_value(RENDER_ACTIVE_TASKS, engine="unified", task_type="compose_video") == 0
with render_task_metrics(engine="unified", task_type="compose_video"):
assert _gauge_value(RENDER_ACTIVE_TASKS, engine="unified", task_type="compose_video") == 1
time.sleep(0.01)
assert _gauge_value(RENDER_ACTIVE_TASKS, engine="unified", task_type="compose_video") == 0
assert _counter_value(RENDER_TASK_TOTAL, engine="unified", task_type="compose_video", status="success") == 1
assert _counter_value(RENDER_TASK_TOTAL, engine="unified", task_type="compose_video", status="failure") == 0
assert _histogram_sum(RENDER_TASK_DURATION_SECONDS, engine="unified", task_type="compose_video") > 0
def test_render_task_metrics_failure():
"""失败任务应记录 failure 状态 + 对应错误类型。"""
from video_processing.render_metrics import (
RENDER_ERROR_TOTAL,
RENDER_TASK_TOTAL,
render_task_metrics,
)
with pytest.raises(ValueError, match="test error"):
with render_task_metrics(engine="legacy", task_type="edit_plan"):
raise ValueError("test error")
assert _counter_value(RENDER_TASK_TOTAL, engine="legacy", task_type="edit_plan", status="failure") == 1
assert _counter_value(RENDER_TASK_TOTAL, engine="legacy", task_type="edit_plan", status="success") == 0
assert _counter_value(RENDER_ERROR_TOTAL, engine="legacy", task_type="edit_plan", error_type="validation") == 1
def test_render_task_metrics_ffmpeg_error():
"""FFmpeg 执行失败应分类为 ffmpeg_error。"""
from video_processing.render_metrics import RENDER_ERROR_TOTAL, render_task_metrics
with pytest.raises(subprocess.CalledProcessError):
with render_task_metrics(engine="unified", task_type="generate_video"):
raise subprocess.CalledProcessError(returncode=1, cmd=["ffmpeg"])
assert _counter_value(RENDER_ERROR_TOTAL, engine="unified", task_type="generate_video", error_type="ffmpeg_error") == 1
def test_render_task_metrics_timeout():
"""超时异常应分类为 timeout。"""
from video_processing.render_metrics import RENDER_ERROR_TOTAL, render_task_metrics
with pytest.raises(subprocess.TimeoutExpired):
with render_task_metrics(engine="unified", task_type="compose_video"):
raise subprocess.TimeoutExpired(cmd=["ffmpeg"], timeout=10)
assert _counter_value(RENDER_ERROR_TOTAL, engine="unified", task_type="compose_video", error_type="timeout") == 1
def test_render_task_metrics_active_tasks_decrement_on_error():
"""任务异常退出时,活跃任务数也应正确递减。"""
from video_processing.render_metrics import RENDER_ACTIVE_TASKS, render_task_metrics
assert _gauge_value(RENDER_ACTIVE_TASKS, engine="legacy", task_type="compose_video") == 0
with pytest.raises(RuntimeError):
with render_task_metrics(engine="legacy", task_type="compose_video"):
raise RuntimeError("boom")
assert _gauge_value(RENDER_ACTIVE_TASKS, engine="legacy", task_type="compose_video") == 0
def test_record_stream_copy():
"""stream_copy 命中/跳过/回退应分别计数。"""
from video_processing.render_metrics import RENDER_STREAM_COPY_TOTAL, record_stream_copy
assert _counter_value(RENDER_STREAM_COPY_TOTAL, result="hit", reason="成功") == 0
assert _counter_value(RENDER_STREAM_COPY_TOTAL, result="miss", reason="编码不匹配") == 0
record_stream_copy("hit", "成功")
record_stream_copy("hit", "成功")
record_stream_copy("miss", "编码不匹配")
record_stream_copy("fallback", "ffmpeg失败")
assert _counter_value(RENDER_STREAM_COPY_TOTAL, result="hit", reason="成功") == 2
assert _counter_value(RENDER_STREAM_COPY_TOTAL, result="miss", reason="编码不匹配") == 1
assert _counter_value(RENDER_STREAM_COPY_TOTAL, result="fallback", reason="ffmpeg失败") == 1
def test_classify_error_unknown():
"""未知异常应分类为 unknown。"""
from video_processing.render_metrics import _classify_error
assert _classify_error(RuntimeError("something")) == "unknown"
def test_classify_error_io():
"""IO 异常应分类为 io_error。"""
from video_processing.render_metrics import _classify_error
assert _classify_error(OSError("disk full")) == "io_error"
def test_engine_and_task_type_labels():
"""不同 engine 和 task_type 应独立计数。"""
from video_processing.render_metrics import RENDER_TASK_TOTAL, render_task_metrics
with render_task_metrics(engine="unified", task_type="compose_video"):
pass
with render_task_metrics(engine="legacy", task_type="compose_video"):
pass
with render_task_metrics(engine="unified", task_type="edit_plan"):
pass
assert _counter_value(RENDER_TASK_TOTAL, engine="unified", task_type="compose_video", status="success") == 1
assert _counter_value(RENDER_TASK_TOTAL, engine="legacy", task_type="compose_video", status="success") == 1
assert _counter_value(RENDER_TASK_TOTAL, engine="unified", task_type="edit_plan", status="success") == 1