Compare commits

..

1 Commits

Author SHA1 Message Date
unknown 25c65280e5 feat(cps): schedule daily order reconciliation 2026-07-22 19:55:22 +08:00
4 changed files with 45 additions and 223 deletions
+18 -79
View File
@@ -1,7 +1,7 @@
"""admin「领券数据」看板聚合:发起/完成数、耗时均值与分位、按天/小时趋势、逐条明细。
数据源 coupon_session(一次领券一行,客户端 /api/v1/coupon/session 两段上报)。生产 PostgreSQL
使用 percentile_cont 聚合耗时分位SQLite 本地/测试环境回退读取耗时单列计算
数据源 coupon_session(一次领券一行,客户端 /api/v1/coupon/session 两段上报)。量级不大,全量拉
区间数据后 Python 聚合(分位 SQLite 无 percentile,统一 Python 算,PG 上也一致)
- 发起数 = 区间内全部 session(含 started/completed/failed/abandoned),= 流失统计的基数。
- 完成数 / 耗时均值 / 分位 = 仅 status==completed 子集(成功跑完才有可比的"领券耗时")。
- summary/daily/hourly/total 在全量上算,不受分页;items 为排序后当前页。
@@ -45,75 +45,6 @@ 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
@@ -284,22 +215,30 @@ def coupon_data_report(
if not user_ids:
return _empty_result()
conditions = [
stmt = select(CouponSession).where(
CouponSession.started_date >= d_from,
CouponSession.started_date <= d_to,
]
)
if app_env is not None:
conditions.append(CouponSession.app_env == app_env)
stmt = stmt.where(CouponSession.app_env == app_env)
if statuses:
conditions.append(CouponSession.status.in_(statuses))
stmt = stmt.where(CouponSession.status.in_(statuses))
if user_ids is not None:
conditions.append(CouponSession.user_id.in_(user_ids))
stmt = select(CouponSession).where(*conditions)
stmt = stmt.where(CouponSession.user_id.in_(user_ids))
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 = {
**_coupon_summary_aggregates(db, conditions),
"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),
**_success_rates(rows),
}
@@ -384,7 +323,7 @@ def coupon_data_report(
"summary": summary,
"daily": daily,
"hourly": hourly,
"total": summary["started_count"],
"total": len(rows),
"items": items,
}
+25 -73
View File
@@ -6,7 +6,6 @@
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
@@ -298,58 +297,6 @@ 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,
*,
@@ -385,18 +332,20 @@ def comparison_records_summary(
success = int(row[2] or 0)
lower_price = int(row[4] or 0)
cancelled = int(row[5] or 0)
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_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_rate_denominator = started - cancelled
return {
"started": started,
@@ -405,16 +354,19 @@ 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": 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],
"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),
"cancelled": cancelled,
"cancelled_rate": cancelled / started if started else None,
"cancelled_p5_ms": cancelled_duration_stats[1],
"cancelled_p50_ms": cancelled_duration_stats[2],
"cancelled_p95_ms": cancelled_duration_stats[3],
"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),
}
-16
View File
@@ -4,28 +4,12 @@ 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:
+2 -55
View File
@@ -7,12 +7,8 @@ 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_summary_aggregate_stmt,
coupon_data_report,
)
from app.admin.repositories.coupon_data import coupon_data_report
from app.db.session import SessionLocal
from app.models.coupon_state import CouponSession
from app.repositories.coupon_state import (
@@ -23,12 +19,7 @@ from app.repositories.coupon_state import (
def _agg_session(
trace: str,
platforms,
platform_success,
*,
status: str = "completed",
elapsed_ms: int | None = None,
trace: str, platforms, platform_success, *, status: str = "completed"
) -> CouponSession:
"""构造一条聚合测试用 session(started_date 固定 2020-01-02、app_env=prod,不 commit)。"""
return CouponSession(
@@ -40,53 +31,9 @@ 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,