diff --git a/app/admin/repositories/coupon_data.py b/app/admin/repositories/coupon_data.py index 422136e..5ae22d4 100644 --- a/app/admin/repositories/coupon_data.py +++ b/app/admin/repositories/coupon_data.py @@ -1,7 +1,7 @@ """admin「领券数据」看板聚合:发起/完成数、耗时均值与分位、按天/小时趋势、逐条明细。 -数据源 coupon_session(一次领券一行,客户端 /api/v1/coupon/session 两段上报)。量级不大,全量拉 -区间数据后 Python 聚合(分位 SQLite 无 percentile,统一 Python 算,PG 上也一致)。 +数据源 coupon_session(一次领券一行,客户端 /api/v1/coupon/session 两段上报)。生产 PostgreSQL +使用 percentile_cont 聚合耗时分位;SQLite 本地/测试环境回退读取耗时单列计算。 - 发起数 = 区间内全部 session(含 started/completed/failed/abandoned),= 流失统计的基数。 - 完成数 / 耗时均值 / 分位 = 仅 status==completed 子集(成功跑完才有可比的"领券耗时")。 - 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) +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: return round(sum(vals) / len(vals)) if vals else None @@ -215,30 +284,22 @@ def coupon_data_report( if not user_ids: return _empty_result() - stmt = select(CouponSession).where( + conditions = [ CouponSession.started_date >= d_from, CouponSession.started_date <= d_to, - ) + ] if app_env is not None: - stmt = stmt.where(CouponSession.app_env == app_env) + conditions.append(CouponSession.app_env == app_env) if statuses: - stmt = stmt.where(CouponSession.status.in_(statuses)) + conditions.append(CouponSession.status.in_(statuses)) 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()) # ── 汇总卡 ── - completed_elapsed = sorted( - r.elapsed_ms for r in rows if r.status == "completed" and r.elapsed_ms is not None - ) summary = { - "started_count": len(rows), - "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), + **_coupon_summary_aggregates(db, conditions), **_success_rates(rows), } @@ -323,7 +384,7 @@ def coupon_data_report( "summary": summary, "daily": daily, "hourly": hourly, - "total": len(rows), + "total": summary["started_count"], "items": items, } diff --git a/app/admin/repositories/queries.py b/app/admin/repositories/queries.py index 26a36db..33b510d 100644 --- a/app/admin/repositories/queries.py +++ b/app/admin/repositories/queries.py @@ -6,6 +6,7 @@ from __future__ import annotations from datetime import date, datetime, time, timedelta, timezone +from decimal import ROUND_HALF_UP, Decimal from zoneinfo import ZoneInfo 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) +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( db: Session, *, @@ -332,20 +385,18 @@ def comparison_records_summary( success = int(row[2] or 0) lower_price = int(row[4] or 0) cancelled = int(row[5] or 0) - success_durations = sorted(db.execute( - select(ComparisonRecord.total_ms).where( - *conditions, - ComparisonRecord.status == "success", - ComparisonRecord.total_ms.is_not(None), - ) - ).scalars().all()) - cancelled_durations = sorted(db.execute( - select(ComparisonRecord.total_ms).where( - *conditions, - ComparisonRecord.status == "cancelled", - ComparisonRecord.total_ms.is_not(None), - ) - ).scalars().all()) + success_duration_stats = _comparison_duration_aggregates( + db, + conditions=conditions, + status="success", + quantiles=(0.05, 0.5, 0.95, 0.99), + ) + cancelled_duration_stats = _comparison_duration_aggregates( + db, + conditions=conditions, + status="cancelled", + quantiles=(0.05, 0.5, 0.95), + ) success_rate_denominator = started - cancelled return { "started": started, @@ -354,19 +405,16 @@ def comparison_records_summary( "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, "lower_price_rate": lower_price / success if success else None, - "avg_duration_ms": ( - int(sum(success_durations) / len(success_durations) + 0.5) - if success_durations else None - ), - "p5_duration_ms": _comparison_percentile(success_durations, 0.05), - "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), + "avg_duration_ms": success_duration_stats[0], + "p5_duration_ms": success_duration_stats[1], + "p50_duration_ms": success_duration_stats[2], + "p95_duration_ms": success_duration_stats[3], + "p99_duration_ms": success_duration_stats[4], "cancelled": cancelled, "cancelled_rate": cancelled / started if started else None, - "cancelled_p5_ms": _comparison_percentile(cancelled_durations, 0.05), - "cancelled_p50_ms": _comparison_percentile(cancelled_durations, 0.5), - "cancelled_p95_ms": _comparison_percentile(cancelled_durations, 0.95), + "cancelled_p5_ms": cancelled_duration_stats[1], + "cancelled_p50_ms": cancelled_duration_stats[2], + "cancelled_p95_ms": cancelled_duration_stats[3], } diff --git a/tests/test_comparison_admin_summary.py b/tests/test_comparison_admin_summary.py index 5006fd6..2f86e24 100644 --- a/tests/test_comparison_admin_summary.py +++ b/tests/test_comparison_admin_summary.py @@ -4,12 +4,28 @@ from __future__ import annotations from datetime import UTC, date, datetime import pytest +from sqlalchemy.dialects import postgresql from app.admin.repositories import queries from app.db.session import SessionLocal 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: db = SessionLocal() try: diff --git a/tests/test_coupon_platform_success.py b/tests/test_coupon_platform_success.py index adc04a9..89ecf9b 100644 --- a/tests/test_coupon_platform_success.py +++ b/tests/test_coupon_platform_success.py @@ -7,8 +7,12 @@ from __future__ import annotations from datetime import UTC, date, datetime 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.models.coupon_state import CouponSession from app.repositories.coupon_state import ( @@ -19,7 +23,12 @@ from app.repositories.coupon_state import ( 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: """构造一条聚合测试用 session(started_date 固定 2020-01-02、app_env=prod,不 commit)。""" return CouponSession( @@ -31,9 +40,53 @@ def _agg_session( platform_success=platform_success, started_at=datetime(2020, 1, 2, tzinfo=UTC), 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: row = CouponSession( trace_id=trace_id,