Compare commits

...

1 Commits

Author SHA1 Message Date
unknown 7f28355378 perf(admin): aggregate record duration percentiles in PostgreSQL 2026-07-22 19:37:30 +08:00
4 changed files with 223 additions and 45 deletions
+79 -18
View File
@@ -1,7 +1,7 @@
"""admin「领券数据」看板聚合:发起/完成数、耗时均值与分位、按天/小时趋势、逐条明细。 """admin「领券数据」看板聚合:发起/完成数、耗时均值与分位、按天/小时趋势、逐条明细。
数据源 coupon_session(一次领券一行,客户端 /api/v1/coupon/session 两段上报)。量级不大,全量拉 数据源 coupon_session(一次领券一行,客户端 /api/v1/coupon/session 两段上报)。生产 PostgreSQL
区间数据后 Python 聚合(分位 SQLite 无 percentile,统一 Python 算,PG 上也一致) 使用 percentile_cont 聚合耗时分位SQLite 本地/测试环境回退读取耗时单列计算
- 发起数 = 区间内全部 session(含 started/completed/failed/abandoned),= 流失统计的基数。 - 发起数 = 区间内全部 session(含 started/completed/failed/abandoned),= 流失统计的基数。
- 完成数 / 耗时均值 / 分位 = 仅 status==completed 子集(成功跑完才有可比的"领券耗时")。 - 完成数 / 耗时均值 / 分位 = 仅 status==completed 子集(成功跑完才有可比的"领券耗时")。
- summary/daily/hourly/total 在全量上算,不受分页;items 为排序后当前页。 - summary/daily/hourly/total 在全量上算,不受分页;items 为排序后当前页。
@@ -45,6 +45,75 @@ def _percentile(sorted_vals: list[int], q: float) -> int | None:
return round(sorted_vals[lo] * (1 - frac) + sorted_vals[hi] * frac) return round(sorted_vals[lo] * (1 - frac) + sorted_vals[hi] * frac)
def _round_duration_ms(value) -> int | None:
"""将数据库聚合结果按既有 Python round 口径转为整数毫秒。"""
if value is None:
return None
return int(round(value))
def _coupon_summary_aggregate_stmt(conditions: list):
"""PostgreSQL 汇总卡聚合语句;计数、均值与四个分位一次返回。"""
completed = CouponSession.status == "completed"
completed_elapsed = completed & CouponSession.elapsed_ms.is_not(None)
return select(
func.count(CouponSession.id),
func.sum(case((completed, 1), else_=0)),
func.avg(CouponSession.elapsed_ms).filter(completed_elapsed),
*(
func.percentile_cont(q)
.within_group(CouponSession.elapsed_ms)
.filter(completed_elapsed)
for q in (0.05, 0.5, 0.95, 0.99)
),
).where(*conditions)
def _coupon_summary_aggregates(db: Session, conditions: list) -> dict:
"""汇总卡基础指标;生产 PG 全部在数据库内完成,SQLite 仅作测试回退。"""
if db.bind is not None and db.bind.dialect.name == "postgresql":
row = db.execute(_coupon_summary_aggregate_stmt(conditions)).one()
return {
"started_count": int(row[0] or 0),
"completed_count": int(row[1] or 0),
"avg_elapsed_ms": _round_duration_ms(row[2]),
"p5_ms": _round_duration_ms(row[3]),
"p50_ms": _round_duration_ms(row[4]),
"p95_ms": _round_duration_ms(row[5]),
"p99_ms": _round_duration_ms(row[6]),
}
counts = db.execute(
select(
func.count(CouponSession.id),
func.sum(case((CouponSession.status == "completed", 1), else_=0)),
).where(*conditions)
).one()
# SQLite 没有 percentile_cont;本地/测试只回退读取耗时单列,不加载完整记录。
completed_elapsed = list(
db.execute(
select(CouponSession.elapsed_ms)
.where(
*conditions,
CouponSession.status == "completed",
CouponSession.elapsed_ms.is_not(None),
)
.order_by(CouponSession.elapsed_ms)
).scalars()
)
return {
"started_count": int(counts[0] or 0),
"completed_count": int(counts[1] or 0),
"avg_elapsed_ms": _round_duration_ms(
sum(completed_elapsed) / len(completed_elapsed)
) if completed_elapsed else None,
"p5_ms": _percentile(completed_elapsed, 5),
"p50_ms": _percentile(completed_elapsed, 50),
"p95_ms": _percentile(completed_elapsed, 95),
"p99_ms": _percentile(completed_elapsed, 99),
}
def _avg(vals: list[int]) -> int | None: def _avg(vals: list[int]) -> int | None:
return round(sum(vals) / len(vals)) if vals else None return round(sum(vals) / len(vals)) if vals else None
@@ -215,30 +284,22 @@ def coupon_data_report(
if not user_ids: if not user_ids:
return _empty_result() return _empty_result()
stmt = select(CouponSession).where( conditions = [
CouponSession.started_date >= d_from, CouponSession.started_date >= d_from,
CouponSession.started_date <= d_to, CouponSession.started_date <= d_to,
) ]
if app_env is not None: if app_env is not None:
stmt = stmt.where(CouponSession.app_env == app_env) conditions.append(CouponSession.app_env == app_env)
if statuses: if statuses:
stmt = stmt.where(CouponSession.status.in_(statuses)) conditions.append(CouponSession.status.in_(statuses))
if user_ids is not None: if user_ids is not None:
stmt = stmt.where(CouponSession.user_id.in_(user_ids)) conditions.append(CouponSession.user_id.in_(user_ids))
stmt = select(CouponSession).where(*conditions)
rows = list(db.execute(stmt).scalars()) rows = list(db.execute(stmt).scalars())
# ── 汇总卡 ── # ── 汇总卡 ──
completed_elapsed = sorted(
r.elapsed_ms for r in rows if r.status == "completed" and r.elapsed_ms is not None
)
summary = { summary = {
"started_count": len(rows), **_coupon_summary_aggregates(db, conditions),
"completed_count": sum(1 for r in rows if r.status == "completed"),
"avg_elapsed_ms": _avg(completed_elapsed),
"p5_ms": _percentile(completed_elapsed, 5),
"p50_ms": _percentile(completed_elapsed, 50),
"p95_ms": _percentile(completed_elapsed, 95),
"p99_ms": _percentile(completed_elapsed, 99),
**_success_rates(rows), **_success_rates(rows),
} }
@@ -323,7 +384,7 @@ def coupon_data_report(
"summary": summary, "summary": summary,
"daily": daily, "daily": daily,
"hourly": hourly, "hourly": hourly,
"total": len(rows), "total": summary["started_count"],
"items": items, "items": items,
} }
+73 -25
View File
@@ -6,6 +6,7 @@
from __future__ import annotations from __future__ import annotations
from datetime import date, datetime, time, timedelta, timezone from datetime import date, datetime, time, timedelta, timezone
from decimal import ROUND_HALF_UP, Decimal
from zoneinfo import ZoneInfo from zoneinfo import ZoneInfo
from sqlalchemy import Select, asc, case, desc, func, or_, select from sqlalchemy import Select, asc, case, desc, func, or_, select
@@ -297,6 +298,58 @@ def _comparison_percentile(sorted_values: list[int], q: float) -> int | None:
return int(value + 0.5) return int(value + 0.5)
def _round_duration_ms(value) -> int | None:
"""将数据库聚合结果按既有口径四舍五入为整数毫秒。"""
if value is None:
return None
return int(Decimal(str(value)).quantize(Decimal("1"), rounding=ROUND_HALF_UP))
def _comparison_duration_aggregate_stmt(conditions: list, status: str, quantiles: tuple[float, ...]):
"""PostgreSQL 耗时聚合语句;每种状态只返回一行。"""
return select(
func.avg(ComparisonRecord.total_ms),
*(
func.percentile_cont(q).within_group(ComparisonRecord.total_ms)
for q in quantiles
),
).where(
*conditions,
ComparisonRecord.status == status,
ComparisonRecord.total_ms.is_not(None),
)
def _comparison_duration_aggregates(
db: Session,
*,
conditions: list,
status: str,
quantiles: tuple[float, ...],
) -> list[int | None]:
"""返回平均值和各分位数;生产 PG 在数据库内聚合,SQLite 仅作测试回退。"""
if db.bind is not None and db.bind.dialect.name == "postgresql":
row = db.execute(
_comparison_duration_aggregate_stmt(conditions, status, quantiles)
).one()
return [_round_duration_ms(value) for value in row]
# SQLite 没有 percentile_cont;本地/测试只回退读取耗时单列,不加载完整记录。
values = list(
db.execute(
select(ComparisonRecord.total_ms)
.where(
*conditions,
ComparisonRecord.status == status,
ComparisonRecord.total_ms.is_not(None),
)
.order_by(ComparisonRecord.total_ms)
).scalars()
)
average = _round_duration_ms(sum(values) / len(values)) if values else None
return [average, *(_comparison_percentile(values, q) for q in quantiles)]
def comparison_records_summary( def comparison_records_summary(
db: Session, db: Session,
*, *,
@@ -332,20 +385,18 @@ def comparison_records_summary(
success = int(row[2] or 0) success = int(row[2] or 0)
lower_price = int(row[4] or 0) lower_price = int(row[4] or 0)
cancelled = int(row[5] or 0) cancelled = int(row[5] or 0)
success_durations = sorted(db.execute( success_duration_stats = _comparison_duration_aggregates(
select(ComparisonRecord.total_ms).where( db,
*conditions, conditions=conditions,
ComparisonRecord.status == "success", status="success",
ComparisonRecord.total_ms.is_not(None), quantiles=(0.05, 0.5, 0.95, 0.99),
) )
).scalars().all()) cancelled_duration_stats = _comparison_duration_aggregates(
cancelled_durations = sorted(db.execute( db,
select(ComparisonRecord.total_ms).where( conditions=conditions,
*conditions, status="cancelled",
ComparisonRecord.status == "cancelled", quantiles=(0.05, 0.5, 0.95),
ComparisonRecord.total_ms.is_not(None), )
)
).scalars().all())
success_rate_denominator = started - cancelled success_rate_denominator = started - cancelled
return { return {
"started": started, "started": started,
@@ -354,19 +405,16 @@ def comparison_records_summary(
"success_rate": success / success_rate_denominator if success_rate_denominator else None, "success_rate": success / success_rate_denominator if success_rate_denominator else None,
"avg_token_cost": float(row[3]) if row[3] is not None else None, "avg_token_cost": float(row[3]) if row[3] is not None else None,
"lower_price_rate": lower_price / success if success else None, "lower_price_rate": lower_price / success if success else None,
"avg_duration_ms": ( "avg_duration_ms": success_duration_stats[0],
int(sum(success_durations) / len(success_durations) + 0.5) "p5_duration_ms": success_duration_stats[1],
if success_durations else None "p50_duration_ms": success_duration_stats[2],
), "p95_duration_ms": success_duration_stats[3],
"p5_duration_ms": _comparison_percentile(success_durations, 0.05), "p99_duration_ms": success_duration_stats[4],
"p50_duration_ms": _comparison_percentile(success_durations, 0.5),
"p95_duration_ms": _comparison_percentile(success_durations, 0.95),
"p99_duration_ms": _comparison_percentile(success_durations, 0.99),
"cancelled": cancelled, "cancelled": cancelled,
"cancelled_rate": cancelled / started if started else None, "cancelled_rate": cancelled / started if started else None,
"cancelled_p5_ms": _comparison_percentile(cancelled_durations, 0.05), "cancelled_p5_ms": cancelled_duration_stats[1],
"cancelled_p50_ms": _comparison_percentile(cancelled_durations, 0.5), "cancelled_p50_ms": cancelled_duration_stats[2],
"cancelled_p95_ms": _comparison_percentile(cancelled_durations, 0.95), "cancelled_p95_ms": cancelled_duration_stats[3],
} }
+16
View File
@@ -4,12 +4,28 @@ from __future__ import annotations
from datetime import UTC, date, datetime from datetime import UTC, date, datetime
import pytest import pytest
from sqlalchemy.dialects import postgresql
from app.admin.repositories import queries from app.admin.repositories import queries
from app.db.session import SessionLocal from app.db.session import SessionLocal
from app.models.comparison import ComparisonRecord from app.models.comparison import ComparisonRecord
def test_postgresql_duration_summary_uses_ordered_set_aggregates() -> None:
stmt = queries._comparison_duration_aggregate_stmt(
[], "success", (0.05, 0.5, 0.95, 0.99)
)
sql = str(
stmt.compile(
dialect=postgresql.dialect(),
compile_kwargs={"literal_binds": True},
)
)
assert sql.count("percentile_cont") == 4
assert "comparison_record.status = 'success'" in sql
def test_summary_uses_only_success_durations_and_filters_beijing_date() -> None: def test_summary_uses_only_success_durations_and_filters_beijing_date() -> None:
db = SessionLocal() db = SessionLocal()
try: try:
+55 -2
View File
@@ -7,8 +7,12 @@ from __future__ import annotations
from datetime import UTC, date, datetime from datetime import UTC, date, datetime
from sqlalchemy import delete, func, select from sqlalchemy import delete, func, select
from sqlalchemy.dialects import postgresql
from app.admin.repositories.coupon_data import coupon_data_report from app.admin.repositories.coupon_data import (
_coupon_summary_aggregate_stmt,
coupon_data_report,
)
from app.db.session import SessionLocal from app.db.session import SessionLocal
from app.models.coupon_state import CouponSession from app.models.coupon_state import CouponSession
from app.repositories.coupon_state import ( from app.repositories.coupon_state import (
@@ -19,7 +23,12 @@ from app.repositories.coupon_state import (
def _agg_session( def _agg_session(
trace: str, platforms, platform_success, *, status: str = "completed" trace: str,
platforms,
platform_success,
*,
status: str = "completed",
elapsed_ms: int | None = None,
) -> CouponSession: ) -> CouponSession:
"""构造一条聚合测试用 session(started_date 固定 2020-01-02、app_env=prod,不 commit)。""" """构造一条聚合测试用 session(started_date 固定 2020-01-02、app_env=prod,不 commit)。"""
return CouponSession( return CouponSession(
@@ -31,9 +40,53 @@ def _agg_session(
platform_success=platform_success, platform_success=platform_success,
started_at=datetime(2020, 1, 2, tzinfo=UTC), started_at=datetime(2020, 1, 2, tzinfo=UTC),
started_date=date(2020, 1, 2), started_date=date(2020, 1, 2),
elapsed_ms=elapsed_ms,
) )
def test_postgresql_coupon_summary_uses_ordered_set_aggregates() -> None:
sql = str(
_coupon_summary_aggregate_stmt([]).compile(
dialect=postgresql.dialect(),
compile_kwargs={"literal_binds": True},
)
)
assert sql.count("percentile_cont") == 4
assert "FILTER (WHERE coupon_session.status = 'completed'" in sql
def test_coupon_duration_summary_uses_only_completed_rows() -> None:
db = SessionLocal()
try:
db.add_all([
_agg_session("duration-a", [], [], elapsed_ms=1000),
_agg_session("duration-b", [], [], elapsed_ms=3000),
_agg_session(
"duration-failed", [], [], status="failed", elapsed_ms=100_000
),
])
db.flush()
summary = coupon_data_report(
db,
date_from="2020-01-02",
date_to="2020-01-02",
app_env="prod",
)["summary"]
assert summary["started_count"] == 3
assert summary["completed_count"] == 2
assert summary["avg_elapsed_ms"] == 2000
assert summary["p5_ms"] == 1100
assert summary["p50_ms"] == 2000
assert summary["p95_ms"] == 2900
assert summary["p99_ms"] == 2980
finally:
db.rollback()
db.close()
def _make_session(db, trace_id: str, **kw) -> CouponSession: def _make_session(db, trace_id: str, **kw) -> CouponSession:
row = CouponSession( row = CouponSession(
trace_id=trace_id, trace_id=trace_id,