Skip to content

Commit 90204a2

Browse files
committed
refactor(stats): 统计一律读小时汇总,总览与图表口径统一(A2)
问题(实测):overview/by_provider 读 usage_events(明细,只留 90 天), timeline/model-timeline 读 usage_hourly(永久)。选「全部」时同屏两个数字 互相矛盾:200 天前的明细被清理后 overview=1、timeline=2。 - usage_hourly 补 reasoning_tokens / cached_tokens / cached_known 三列 (SCHEMA_VERSION 5→6,走 _MIGRATION_COLUMNS;老库历史行补 0) - overview / by_provider 改读小时汇总,与图表彻底同源 - record() 增量累加当前小时行:新请求立即可见,不必等 5 分钟一轮的 rollup (否则刚发生的请求统计页显示 0,比原来的口径矛盾更糟); rollup_hourly 仍全量重算,两者结果一致(幂等,已加测试) - 顺带修 C2:延迟均值分子分母配对——分子只累加成功请求 (此前 SUM(latency_ms) 含失败、分母只有 ok_count,实测 200ms 被算成 650ms) - 修 scripts/cleanup_invalid_stats.py 的 _DbAdapter:Database.transaction() 重构后它少了该方法,--apply 直接 AttributeError(补该脚本的测试,此前零覆盖) 测试:总览与图表在 90 天外一致、reasoning/cached 聚合、均值只算成功、 老库补列保留历史行、增量与全量重算等价。
1 parent b331d8f commit 90204a2

8 files changed

Lines changed: 286 additions & 31 deletions

File tree

‎TECHNICAL.md‎

Lines changed: 16 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -78,8 +78,8 @@ coding2api/
7878
│ │ ├── retention.py # 每 5 分钟:小时汇总重算(幂等)+ 90 天前明细清理
7979
│ │ └── runner.py # 后台任务调度,接入应用生命周期
8080
│ ├── stats/
81-
│ │ ├── collector.py # usage_events 写入(脱敏)
82-
│ │ └── query.py # overview / by-provider / events 查询(events JOIN credentials 带凭证昵称)
81+
│ │ ├── collector.py # usage_events 写入(脱敏)+ 小时汇总双写/重算
82+
│ │ └── query.py # overview / by-provider / timeline / events 查询(前者读小时汇总,events JOIN credentials 带凭证昵称)
8383
│ └── api/
8484
│ ├── deps.py # Services 容器 + require_api_key / session / csrf 依赖
8585
│ ├── chat.py # POST /v1/chat/completions
@@ -286,7 +286,12 @@ class Scheduler:
286286
| 额度探测(quota_probe.py) | 启动立即跑一轮(不节流) + 每 `QUOTA_PROBE_MINUTES`(默认 60)分钟 | 探测上游剩余额度 → `credentials.quota_*` / `quota_expiry_ladder` / `health` 写回 |
287287
| token 预刷新(refresh.py) | 每 60 分钟 | 到期前 `REFRESH_SKEW_HOURS`(默认 24h)窗口内轮换 refresh token |
288288
| 每日签到(checkin.py) | 每 10 分钟(全天) | 成功即封账该凭证当日(`日期:scope`,进程内内存态,重启重建);失败凭证持续重试,同账号多凭证共享一次 |
289-
| 明细清理(retention.py) | 每 5 分钟 | `usage_events` 全量重算小时汇总(幂等 upsert,最新小时滞后 ≤5 分钟)+ 90 天前明细清理 |
289+
| 明细清理(retention.py) | 每 5 分钟 | `usage_events` 全量重算小时汇总(幂等 upsert,与 record 的增量双写对账)+ 90 天前明细清理 |
290+
291+
**清理切点必顶对齐到小时边界**(`purge_expired`)。这不是保守取值而是正确性要求:
292+
`rollup_hourly` 对整行是 REPLACE 语义,只有保证「仍有明细的小时保有全部明细」
293+
重算才精确;若把边界小时只删一半,下一轮 rollup 会把汇总行覆盖成剩下那一半,
294+
被删部分永久丢失(明细已不在)。代价:明细最多多留 1 小时。
290295

291296
---
292297

@@ -360,6 +365,14 @@ fixture 断言两个方向:**解析正确**(样本 → 期望 Event)与**
360365
- **手写 SQL 而非 ORM**:8 张表规模下 ORM 收益为负
361366
- **polling OAuth 不转回调**(Q17=C):上游协议决定;TRAE 回调走主端口 + PUBLIC_BASE_URL
362367
- **v1 无 Anthropic**(Q8=A):Event 层已预留,v1.1 只加 `compat/anthropic/` 适配器
368+
- **统计一律以 `usage_hourly` 为准**:`overview` / `by_provider` / `timeline` /
369+
`model-timeline` 均读小时汇总,只有 `events`(逐请求明细)读 `usage_events`。
370+
统一口径是为了让选「全部」时总览与图表同值(明细只留 90 天,汇总永久)
371+
- **小时汇总双写**:`record()` 写入明细的同时增量累加当前小时行,所以新请求
372+
立即可见于统计页(不依赖 5 分钟一轮的 rollup);`rollup_hourly` 仍每 5 分钟
373+
全量重算作对账,两者结果一致(幂等)
374+
- **延迟均值只算成功请求**:分子 `SUM(latency_ms WHERE ok=1)` 与分母 `ok_count`
375+
配对;失败请求的耗时不能拉偏「典型耗时」(与图表口径一致)
363376
- **两套数据源共存(已知不一致)**:`overview` / `by_provider` 读 `usage_events`(即时,
364377
仅覆盖 90 天明细),`timeline` / `model-timeline` 读 `usage_hourly`(≤5 分钟滞后,永久)。
365378
时间范围 ≤90 天时两者一致(汇总由同一批明细算出);选「全部」时总览会小于图表,

‎scripts/cleanup_invalid_stats.py‎

Lines changed: 12 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@
2626
import sqlite3
2727
import sys
2828
import time
29+
from contextlib import contextmanager
2930
from pathlib import Path
3031

3132
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
@@ -77,14 +78,24 @@ def sync_hourly(conn) -> None:
7778

7879

7980
class _DbAdapter:
80-
"""把裸连接包装成 Database.connect() 的形态(rollup 只用这一个方法)。"""
81+
"""把裸连接包装成 Database 的形态(rollup 只用 connect/transaction)。"""
8182

8283
def __init__(self, conn) -> None:
8384
self._conn = conn
8485

8586
def connect(self):
8687
return self._conn
8788

89+
@contextmanager
90+
def transaction(self):
91+
"""与 Database.transaction 同语义:正常提交、异常回滚。"""
92+
try:
93+
yield self._conn
94+
except BaseException:
95+
self._conn.rollback()
96+
raise
97+
self._conn.commit()
98+
8899

89100
def backup(db_path: Path) -> Path:
90101
"""删除前备份主库文件到同目录(.bak-<时间戳>),失败即中止。"""

‎src/db/migrate.py‎

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -16,7 +16,7 @@
1616
SCHEMA_NAME = "schema.sql"
1717

1818
# 当前 schema 版本。新增列/表、删表时 +1,并在下方对应元组里补增量。
19-
SCHEMA_VERSION = 5
19+
SCHEMA_VERSION = 6
2020

2121
# (表, 列定义):历史库升级时逐条补列
2222
_MIGRATION_COLUMNS: tuple[tuple[str, str], ...] = (
@@ -26,6 +26,11 @@
2626
"ttfb_sum INTEGER NOT NULL DEFAULT 0"), # 首字延迟聚合(图表维度,老库补 0)
2727
("credentials",
2828
"quota_expiry_ladder TEXT"), # 到期阶梯 JSON:选号按窗口内到期积分排序
29+
# 总览改为读小时汇总后,这两项也必须能在小时表里聚合(老库历史行补 0,
30+
# 历史小时的数值无法回填:明细已不在,只能接受旧时段显示 0)
31+
("usage_hourly", "reasoning_tokens INTEGER NOT NULL DEFAULT 0"),
32+
("usage_hourly", "cached_tokens INTEGER NOT NULL DEFAULT 0"),
33+
("usage_hourly", "cached_known INTEGER NOT NULL DEFAULT 0"),
2934
)
3035

3136
# 已废弃的表:schema.sql 里已删定义,但老库里可能还留着,必须显式清理。

‎src/db/schema.sql‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -75,6 +75,9 @@ CREATE TABLE IF NOT EXISTS usage_hourly (
7575
ok_count INTEGER NOT NULL DEFAULT 0,
7676
input_tokens INTEGER NOT NULL DEFAULT 0,
7777
output_tokens INTEGER NOT NULL DEFAULT 0,
78+
reasoning_tokens INTEGER NOT NULL DEFAULT 0,
79+
cached_tokens INTEGER NOT NULL DEFAULT 0, -- 命中缓存的输入 token 之和
80+
cached_known INTEGER NOT NULL DEFAULT 0, -- 上报过 cached_tokens 的条数(=0 时该值不可信)
7881
credit_sum REAL,
7982
credit_known INTEGER NOT NULL DEFAULT 0,
8083
latency_sum INTEGER NOT NULL DEFAULT 0,

‎src/stats/collector.py‎

Lines changed: 53 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -109,6 +109,47 @@ def record(
109109
event.output_tokens, event.reasoning_tokens, event.cached_tokens, event.credit,
110110
event.latency_ms, event.ttfb_ms),
111111
)
112+
# 当前小时增量累加:总览/图表都读小时表,不能等 5 分钟一轮的
113+
# retention rollup 才可见(否则刚发生的请求统计页面显示 0)。
114+
# 增量累加与全量重算等价:rollup 后续会把这一小时算成同样的值。
115+
self._bump_hourly(conn, event)
116+
117+
@staticmethod
118+
def _bump_hourly(conn, event: UsageEvent) -> None:
119+
"""把一条明细增量累加进所属小时行(同 key 则累加,不存在则插入)。"""
120+
conn.execute(
121+
"""
122+
INSERT INTO usage_hourly (hour_utc, username, provider, model, requests, ok_count,
123+
input_tokens, output_tokens, reasoning_tokens,
124+
cached_tokens, cached_known,
125+
credit_sum, credit_known, latency_sum, ttfb_sum)
126+
VALUES (?,?,?,?,1,?,?,?,?,?,?,?,?,?,?)
127+
ON CONFLICT(hour_utc, username, provider, model) DO UPDATE SET
128+
requests = requests + 1,
129+
ok_count = ok_count + excluded.ok_count,
130+
input_tokens = input_tokens + excluded.input_tokens,
131+
output_tokens = output_tokens + excluded.output_tokens,
132+
reasoning_tokens = reasoning_tokens + excluded.reasoning_tokens,
133+
cached_tokens = cached_tokens + excluded.cached_tokens,
134+
cached_known = cached_known + excluded.cached_known,
135+
credit_sum = COALESCE(credit_sum, 0) + excluded.credit_sum,
136+
credit_known = credit_known + excluded.credit_known,
137+
latency_sum = latency_sum + excluded.latency_sum,
138+
ttfb_sum = ttfb_sum + excluded.ttfb_sum
139+
""",
140+
# NULL 一律折成 0 参与累加:列定义是 NOT NULL DEFAULT 0,
141+
# `x + NULL` 会交出 NULL,把整行污染掉
142+
((event.ts // 3600) * 3600, event.username, event.provider, event.model,
143+
int(event.ok),
144+
event.input_tokens or 0, event.output_tokens or 0,
145+
event.reasoning_tokens or 0,
146+
event.cached_tokens or 0,
147+
0 if event.cached_tokens is None else 1,
148+
event.credit or 0.0,
149+
0 if event.credit is None else 1,
150+
(event.latency_ms or 0) if event.ok else 0,
151+
(event.ttfb_ms or 0) if event.ok else 0),
152+
)
112153

113154
def purge_expired(self, retention_days: int = 90, now: int | None = None) -> int:
114155
"""明细保留 90 天;小时汇总永久(PROPOSAL §8)。
@@ -139,20 +180,28 @@ def rollup_hourly(self, since: int | None = None) -> int:
139180
cursor = conn.execute(
140181
"""
141182
INSERT INTO usage_hourly (hour_utc, username, provider, model, requests, ok_count,
142-
input_tokens, output_tokens, credit_sum, credit_known,
143-
latency_sum, ttfb_sum)
183+
input_tokens, output_tokens, reasoning_tokens,
184+
cached_tokens, cached_known,
185+
credit_sum, credit_known, latency_sum, ttfb_sum)
144186
SELECT (ts / 3600) * 3600 AS hour_utc, username, provider, model,
145187
COUNT(*), SUM(ok),
146188
COALESCE(SUM(input_tokens), 0), COALESCE(SUM(output_tokens), 0),
147-
SUM(credit), SUM(CASE WHEN credit IS NULL THEN 0 ELSE 1 END),
148-
COALESCE(SUM(latency_ms), 0), COALESCE(SUM(ttfb_ms), 0)
189+
COALESCE(SUM(reasoning_tokens), 0),
190+
COALESCE(SUM(cached_tokens), 0),
191+
SUM(CASE WHEN cached_tokens IS NULL THEN 0 ELSE 1 END),
192+
COALESCE(SUM(credit), 0), SUM(CASE WHEN credit IS NULL THEN 0 ELSE 1 END),
193+
COALESCE(SUM(CASE WHEN ok = 1 THEN latency_ms END), 0),
194+
COALESCE(SUM(CASE WHEN ok = 1 THEN ttfb_ms END), 0)
149195
FROM usage_events WHERE (? IS NULL OR ts >= ?)
150196
GROUP BY hour_utc, username, provider, model
151197
ON CONFLICT(hour_utc, username, provider, model) DO UPDATE SET
152198
requests = excluded.requests,
153199
ok_count = excluded.ok_count,
154200
input_tokens = excluded.input_tokens,
155201
output_tokens = excluded.output_tokens,
202+
reasoning_tokens = excluded.reasoning_tokens,
203+
cached_tokens = excluded.cached_tokens,
204+
cached_known = excluded.cached_known,
156205
credit_sum = excluded.credit_sum,
157206
credit_known = excluded.credit_known,
158207
latency_sum = excluded.latency_sum,

‎src/stats/query.py‎

Lines changed: 31 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -52,48 +52,61 @@ def _where(*, username: str | None = None, provider: str | None = None,
5252

5353
def overview(self, *, username: str | None = None, provider: str | None = None,
5454
since: int | None = None) -> dict[str, Any]:
55-
"""总览。admin 传 username=None 看全局;普通用户只看自己。"""
56-
where, params = self._where(username=username, provider=provider, since=since)
55+
"""总览。admin 传 username=None 看全局;普通用户只看自己。
56+
57+
数据源是 `usage_hourly`(与 timeline/model-timeline 同源):
58+
明细只留 90 天,读明细会让选「全部」时总览小于图表;小时汇总永久保留。
59+
代价:最近 ≤5 分钟未进汇总的请求不计入(retention 每 5 分钟 rollup),
60+
刷新一次即可。延迟类均值统一按 `ok_count` 归一(与图表同口径)。
61+
"""
62+
where, params = self._where(username=username, provider=provider,
63+
since=since, since_col="hour_utc")
5764

5865
row = self._db.connect().execute(
5966
f"""
60-
SELECT COUNT(*) AS requests,
61-
COALESCE(SUM(ok), 0) AS ok_count,
67+
SELECT COALESCE(SUM(requests), 0) AS requests,
68+
COALESCE(SUM(ok_count), 0) AS ok_count,
6269
COALESCE(SUM(input_tokens), 0) AS input_tokens,
6370
COALESCE(SUM(output_tokens), 0) AS output_tokens,
6471
COALESCE(SUM(reasoning_tokens), 0) AS reasoning_tokens,
6572
SUM(cached_tokens) AS cached_tokens,
66-
SUM(CASE WHEN cached_tokens IS NULL THEN 0 ELSE 1 END) AS cached_known,
67-
SUM(credit) AS credit_sum,
68-
SUM(CASE WHEN credit IS NULL THEN 0 ELSE 1 END) AS credit_known,
69-
COALESCE(AVG(latency_ms), 0) AS avg_latency,
70-
COALESCE(AVG(ttfb_ms), 0) AS avg_ttfb
71-
FROM usage_events {where}
73+
SUM(cached_known) AS cached_known,
74+
SUM(credit_sum) AS credit_sum,
75+
SUM(credit_known) AS credit_known,
76+
SUM(latency_sum) AS latency_sum,
77+
SUM(ttfb_sum) AS ttfb_sum
78+
FROM usage_hourly {where}
7279
""", params).fetchone()
7380
requests = row["requests"] or 0
81+
ok_count = row["ok_count"] or 0
82+
# 均值口径:只除成功请求(失败请求的 latency 会拉偏"典型耗时",
83+
# 且与图表 latency_sum/ok_count 同一致),无成功则 None
7484
return {
7585
"requests": requests,
76-
"ok_count": row["ok_count"],
77-
"success_rate": (row["ok_count"] / requests) if requests else None,
86+
"ok_count": ok_count,
87+
"success_rate": (ok_count / requests) if requests else None,
7888
"input_tokens": row["input_tokens"],
7989
"output_tokens": row["output_tokens"],
8090
"reasoning_tokens": row["reasoning_tokens"],
8191
# 缓存命中:任一明细上报过才可信,否则 None(前端显示 —)
8292
"cached_tokens": row["cached_tokens"] if row["cached_known"] else None,
8393
"credit": row["credit_sum"] if row["credit_known"] else None,
84-
"avg_latency_ms": round(row["avg_latency"]) if requests else None,
85-
"avg_ttfb_ms": round(row["avg_ttfb"]) if requests else None,
94+
"avg_latency_ms": round(row["latency_sum"] / ok_count) if ok_count else None,
95+
"avg_ttfb_ms": round(row["ttfb_sum"] / ok_count) if ok_count else None,
8696
}
8797

8898
def by_provider(self, *, username: str | None = None,
8999
since: int | None = None) -> list[dict[str, Any]]:
90-
where, params = self._where(username=username, since=since)
100+
"""按渠道聚合(同样读小时汇总,与总览/图表同源)。"""
101+
where, params = self._where(username=username, since=since, since_col="hour_utc")
91102
rows = self._db.connect().execute(
92103
f"""
93-
SELECT provider, COUNT(*) AS requests, COALESCE(SUM(ok), 0) AS ok_count,
104+
SELECT provider, COALESCE(SUM(requests), 0) AS requests,
105+
COALESCE(SUM(ok_count), 0) AS ok_count,
94106
COALESCE(SUM(input_tokens), 0) AS input_tokens,
95-
COALESCE(SUM(output_tokens), 0) AS output_tokens, SUM(credit) AS credit_sum
96-
FROM usage_events {where} GROUP BY provider ORDER BY provider
107+
COALESCE(SUM(output_tokens), 0) AS output_tokens,
108+
SUM(credit_sum) AS credit_sum
109+
FROM usage_hourly {where} GROUP BY provider ORDER BY provider
97110
""", params).fetchall()
98111
return [
99112
{"provider": row["provider"], "requests": row["requests"],

0 commit comments

Comments
 (0)