Compare commits
2 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 1ee120653e | |||
| f481bcd396 |
+214
@@ -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:
|
||||
|
||||
@@ -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)
|
||||
|
||||
Executable
+195
@@ -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
|
||||
Reference in New Issue
Block a user