fix(analytics): id-tiebreak snapshot ordering + rates/cn_day tests + breakdown ValueError
This commit is contained in:
@@ -9,7 +9,7 @@ from __future__ import annotations
|
||||
from collections import defaultdict
|
||||
from datetime import UTC, datetime
|
||||
|
||||
from sqlalchemy import and_, func, select
|
||||
from sqlalchemy import func, select
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from app.core import rewards
|
||||
@@ -31,7 +31,7 @@ def diff_snapshots(rows: list[dict]) -> list[dict]:
|
||||
|
||||
out: list[dict] = []
|
||||
for group in parts.values():
|
||||
group.sort(key=lambda r: r["created_at"])
|
||||
group.sort(key=lambda r: (r["created_at"], r.get("id", 0)))
|
||||
prev = {k: 0 for k in _COUNTS}
|
||||
for r in group:
|
||||
deltas = {f"d_{k}": max(0, int(r[k]) - prev[k]) for k in _COUNTS}
|
||||
@@ -74,7 +74,7 @@ def _sum_deltas(deltas: list[dict]) -> dict:
|
||||
def _fetch_rows(db: Session, date_from: datetime, date_to: datetime) -> list[dict]:
|
||||
"""取 [from, to) 区间行 + 每分区在 from 左侧的最后一条基线行(供第一条区间增量做差)。"""
|
||||
cols = (
|
||||
H.device_id, H.epoch_id, E.event, H.created_at,
|
||||
H.id, H.device_id, H.epoch_id, E.event, H.created_at,
|
||||
H.app_ver, H.oem, H.os,
|
||||
E.attempted, E.drop_capture, E.delivered, E.drop_undelivered,
|
||||
)
|
||||
@@ -83,8 +83,10 @@ def _fetch_rows(db: Session, date_from: datetime, date_to: datetime) -> list[dic
|
||||
.where(H.created_at >= date_from, H.created_at < date_to)
|
||||
).mappings().all()
|
||||
|
||||
# 注:基线子查询无下界扫 from 左侧全量(spec §7 已接受的取舍;量级变大再上物化 rollup)。
|
||||
# 用 max(id) 而非 max(created_at) 选"最新一条":id 严格单调,避免 SQLite 秒级时间戳撞车时选歧义。
|
||||
sub = (
|
||||
select(H.device_id, H.epoch_id, E.event, func.max(H.created_at).label("mx"))
|
||||
select(H.device_id, H.epoch_id, E.event, func.max(H.id).label("max_id"))
|
||||
.join(E, E.snapshot_id == H.id)
|
||||
.where(H.created_at < date_from)
|
||||
.group_by(H.device_id, H.epoch_id, E.event)
|
||||
@@ -92,13 +94,7 @@ def _fetch_rows(db: Session, date_from: datetime, date_to: datetime) -> list[dic
|
||||
)
|
||||
baseline = db.execute(
|
||||
select(*cols).join(E, E.snapshot_id == H.id).join(
|
||||
sub,
|
||||
and_(
|
||||
sub.c.device_id == H.device_id,
|
||||
sub.c.epoch_id == H.epoch_id,
|
||||
sub.c.event == E.event,
|
||||
sub.c.mx == H.created_at,
|
||||
),
|
||||
sub, sub.c.max_id == H.id
|
||||
)
|
||||
).mappings().all()
|
||||
|
||||
@@ -128,7 +124,8 @@ def trend(db: Session, date_from: datetime, date_to: datetime) -> list[dict]:
|
||||
|
||||
|
||||
def breakdown(db: Session, date_from: datetime, date_to: datetime, dim: str) -> list[dict]:
|
||||
assert dim in ("event", "app_ver", "oem")
|
||||
if dim not in ("event", "app_ver", "oem"):
|
||||
raise ValueError(f"invalid dim: {dim!r}")
|
||||
deltas = _in_range_deltas(db, date_from, date_to)
|
||||
by_key: dict[str, list[dict]] = defaultdict(list)
|
||||
for d in deltas:
|
||||
|
||||
Reference in New Issue
Block a user