Compare commits

..

1 Commits

Author SHA1 Message Date
unknown f767530741 优化:修正广告收益看板统计与明细口径 2026-07-26 19:40:19 +08:00
20 changed files with 730 additions and 147 deletions
+105 -3
View File
@@ -12,10 +12,11 @@
"""
from __future__ import annotations
from sqlalchemy import func, select
from sqlalchemy import func, or_, select
from sqlalchemy.orm import Session
from app.core import rewards
from app.models.ad_ecpm import AdEcpmRecord
from app.models.ad_feed_reward import AdFeedRewardRecord
from app.models.ad_reward import AdRewardRecord
from app.repositories.ad_feed_reward import FEED_REWARD_UNIT_SECONDS
@@ -55,10 +56,22 @@ def _reward_video_rows(
if user_id is not None:
stmt = stmt.where(AdRewardRecord.user_id == user_id)
records = list(db.execute(stmt).scalars())
# S2S 发奖回调本身不携带实际填充的 ADN/底层 rit;按客户端在展示时上报的
# ad_session_id 回填。这样“纯发奖”行也能在运营后台追溯到真实广告网络。
session_ids = {rec.ad_session_id for rec in records if rec.ad_session_id}
impression_by_session = {
(rec.user_id, rec.ad_session_id): rec
for rec in db.execute(
select(AdEcpmRecord).where(AdEcpmRecord.ad_session_id.in_(session_ids))
).scalars()
} if session_ids else {}
# 用本日之前的累计份数做起点,当日 granted 在其上继续递增 → 与 _granted_cumulative+1 对齐
granted_n: dict[int, int] = _prior_granted_counts(db, date=date, user_id=user_id)
rows: list[dict] = []
for rec in db.execute(stmt).scalars():
for rec in records:
impression = impression_by_session.get((rec.user_id, rec.ad_session_id))
if rec.status == "granted":
nth = granted_n.get(rec.user_id, 0) + 1
granted_n[rec.user_id] = nth
@@ -68,6 +81,8 @@ def _reward_video_rows(
"record_id": rec.id,
"user_id": rec.user_id,
"ad_session_id": rec.ad_session_id,
"adn": impression.adn if impression is not None else None,
"slot_id": impression.slot_id if impression is not None else None,
"app_env": rec.app_env,
"our_code_id": rec.our_code_id,
"created_at": rec.created_at,
@@ -90,6 +105,8 @@ def _reward_video_rows(
"record_id": rec.id,
"user_id": rec.user_id,
"ad_session_id": rec.ad_session_id,
"adn": impression.adn if impression is not None else None,
"slot_id": impression.slot_id if impression is not None else None,
"app_env": rec.app_env,
"our_code_id": rec.our_code_id,
"created_at": rec.created_at,
@@ -149,6 +166,81 @@ def _feed_scene_matches(rec: AdFeedRewardRecord, scene: str | None) -> bool:
return True
def _nonblank(value: str | None) -> str | None:
value = value.strip() if value else None
return value or None
def _unique_ad_source(records: list[AdEcpmRecord]) -> tuple[str | None, str | None]:
"""仅在候选展示记录指向唯一 ADN 时回填来源,绝不把一次多广告流程猜成某一个网络。"""
adns = {_nonblank(record.adn) for record in records}
adns.discard(None)
if len(adns) != 1:
return None, None
slots = {_nonblank(record.slot_id) for record in records}
slots.discard(None)
return next(iter(adns)), next(iter(slots)) if len(slots) == 1 else None
def _feed_source_fallbacks(
db: Session, records: list[AdFeedRewardRecord]
) -> tuple[dict[tuple[int, str], tuple[str | None, str | None]], dict[tuple[int, str, str], tuple[str | None, str | None]]]:
"""构建信息流来源回填索引。
新客户端会把 ADN 直接随 feed-reward 上报;旧记录可能缺失。展示收益记录的
``ad_session_id`` 是每条 impressionId,而发奖记录保留的是整场会话 ID,因此先按
会话精确匹配;匹配不到时仅允许按 ``user + trace_id + 原始 eCPM`` 回填,且候选 ADN
必须唯一。trace 内存在多个网络时保持空值,避免错误归因。
"""
session_ids = {record.ad_session_id for record in records if record.ad_session_id}
trace_ids = {record.trace_id for record in records if record.trace_id}
if not session_ids and not trace_ids:
return {}, {}
filters = []
if session_ids:
filters.append(AdEcpmRecord.ad_session_id.in_(session_ids))
if trace_ids:
filters.append(AdEcpmRecord.trace_id.in_(trace_ids))
impressions = list(db.execute(select(AdEcpmRecord).where(or_(*filters))).scalars())
by_session: dict[tuple[int, str], list[AdEcpmRecord]] = {}
by_trace_ecpm: dict[tuple[int, str, str], list[AdEcpmRecord]] = {}
for impression in impressions:
if impression.ad_session_id:
by_session.setdefault((impression.user_id, impression.ad_session_id), []).append(impression)
if impression.trace_id:
by_trace_ecpm.setdefault(
(impression.user_id, impression.trace_id, impression.ecpm_raw), []
).append(impression)
return (
{key: _unique_ad_source(value) for key, value in by_session.items()},
{key: _unique_ad_source(value) for key, value in by_trace_ecpm.items()},
)
def _feed_source(
record: AdFeedRewardRecord,
*,
by_session: dict[tuple[int, str], tuple[str | None, str | None]],
by_trace_ecpm: dict[tuple[int, str, str], tuple[str | None, str | None]],
) -> tuple[str | None, str | None]:
"""取得本条发奖广告的真实来源;无唯一证据时返回原始空值。"""
adn, slot_id = _nonblank(record.adn), _nonblank(record.slot_id)
if adn and slot_id:
return adn, slot_id
candidate = by_session.get((record.user_id, record.ad_session_id or ""))
if candidate is None and record.trace_id:
candidate = by_trace_ecpm.get((record.user_id, record.trace_id, record.ecpm_raw))
if candidate is None:
return adn, slot_id
candidate_adn, candidate_slot_id = candidate
return adn or candidate_adn, slot_id or candidate_slot_id
def _feed_rows(
db: Session, *, date: str, user_id: int | None, scene: str | None = None
) -> list[dict]:
@@ -167,11 +259,17 @@ def _feed_rows(
if user_id is not None:
stmt = stmt.where(AdFeedRewardRecord.user_id == user_id)
records = list(db.execute(stmt).scalars())
by_session, by_trace_ecpm = _feed_source_fallbacks(db, records)
# 本日之前的累计**条数**做起点,与发奖侧 granted_unit_total(COUNT granted)对齐
granted_count: dict[int, int] = _feed_prior_granted_count(db, date=date, user_id=user_id)
rows: list[dict] = []
for rec in db.execute(stmt).scalars():
for rec in records:
keep = _feed_scene_matches(rec, scene) # 累计照常推进,这里只决定是否展示本行
adn, slot_id = _feed_source(
rec, by_session=by_session, by_trace_ecpm=by_trace_ecpm
)
if rec.status == "granted":
# 一条广告 = 1 份(与 grant_feed_reward 同口径:看满一份即发该条满额,不按 unit_count 累加)。
# nth = 账号累计第几**条**(含本日之前),与发奖侧 granted_unit_total+1 对齐;累计照常推进
@@ -188,6 +286,8 @@ def _feed_rows(
"record_id": rec.id,
"user_id": rec.user_id,
"ad_session_id": rec.ad_session_id,
"adn": adn,
"slot_id": slot_id,
"trace_id": rec.trace_id,
"app_env": rec.app_env,
"our_code_id": rec.our_code_id,
@@ -214,6 +314,8 @@ def _feed_rows(
"record_id": rec.id,
"user_id": rec.user_id,
"ad_session_id": rec.ad_session_id,
"adn": adn,
"slot_id": slot_id,
"trace_id": rec.trace_id,
"app_env": rec.app_env,
"our_code_id": rec.our_code_id,
+94 -56
View File
@@ -22,7 +22,7 @@ report_date / reward_date 归日。
"""
from __future__ import annotations
from datetime import UTC, datetime, timedelta
from datetime import UTC, datetime, time, timedelta
from datetime import date as _date
from sqlalchemy import select
@@ -81,9 +81,9 @@ def _date_range(date_from: str, date_to: str) -> list[str]:
# ad_feed_reward_record,由 audit 内部按 ad_type 区分(feed 含历史 NULL,draw 仅 ad_type=="draw")。
_AUDIT_SCENES = {"reward_video", "feed", "draw"}
# 激励视频未满足有效播放条件时不计客户端预估收益。客户端仍会在 onAdShow
# 上报 eCPM,随后才在关闭时补报以下终态,因此必须在展示/发奖合并后修正收益
_ZERO_REVENUE_REWARD_VIDEO_STATUSES = frozenset({"closed_early", "too_short"})
# GroMore 官方说明第三方 ADN 的 Reporting API 最晚约 13:50 更新。只有 D+1 14:00
# 之后完成的同步才标记为「API 同步窗口完成」;这不代表覆盖全部 ADN 或最终结算
_PANGLE_API_FINAL_SYNC_TIME = time(hour=14)
# 发奖复算明细字段(展开下钻看「金币怎么算出来的」)——从 audit 行原样取这些 key。
@@ -96,7 +96,32 @@ _REWARD_DETAIL_KEYS = (
def _reward_detail(row: dict) -> dict:
"""从 audit 行抽出发奖复算明细(给前端展开行渲染因子1/因子2/份数/LT/应发实发)。"""
return {k: row[k] for k in _REWARD_DETAIL_KEYS}
detail = {k: row[k] for k in _REWARD_DETAIL_KEYS}
# 发奖明细必须保留自己的广告网络,不能复用整场聚合父行的来源:
# 同一次比价/领券可能先后由不同 ADN 填充。
detail["adn"] = row.get("adn")
detail["slot_id"] = row.get("slot_id")
return detail
def _as_cn(dt: datetime) -> datetime:
"""数据库 synced_at → 北京时间;SQLite naive 值按 UTC 处理。"""
if dt.tzinfo is None:
dt = dt.replace(tzinfo=UTC)
return dt.astimezone(rewards.CN_TZ)
def _pangle_api_day_complete(day: str, aggregate: dict) -> bool:
"""某天 API 收益是否已在 D+1 14:00 后同步(仅表示同步窗口完成)。"""
synced_at = aggregate.get("synced_at")
if aggregate.get("api_revenue_yuan") is None or synced_at is None:
return False
cutoff = datetime.combine(
_date.fromisoformat(day) + timedelta(days=1),
_PANGLE_API_FINAL_SYNC_TIME,
tzinfo=rewards.CN_TZ,
)
return _as_cn(synced_at) >= cutoff
def ad_revenue_report(
@@ -186,12 +211,10 @@ def ad_revenue_report(
"has_impression": True,
"impressions": 1,
"ecpm": rec.ecpm_raw,
# 单次展示收益(元)= eCPM元 ÷ 1000(每千次→单次)。eCPM 先钳到 AD_ECPM_MAX_FEN(¥500 CPM)
# 再折收益,与发奖口径 [rewards.calculate_ad_reward_coin] 一致(2026-06-29 修:原裸 parse_ecpm_yuan
# 不钳,伪造/异常天价 eCPM 会把报表预估收益冲到任意大;金币侧已钳、收益侧漏钳)
"revenue_yuan": round(
min(rewards.parse_ecpm_yuan(rec.ecpm_raw), rewards.AD_ECPM_MAX_FEN / 100.0) / 1000.0, 6,
),
# 客户端 SDK 展示预估收益(元)= 后端留存 getEcpm 元/千次 ÷ 1000。
# 这里不能复用发奖防作弊的 ¥500 CPM 钳顶:钳顶只限制金币成本,不改变广告已产生的
# 收入估值。onAdShow 已发生即计展示收入,是否看满只影响发奖,不影响广告收入
"revenue_yuan": round(rewards.parse_ecpm_yuan(rec.ecpm_raw) / 1000.0, 6),
"adn": rec.adn,
"slot_id": rec.slot_id,
"sub_rewards": [],
@@ -206,11 +229,6 @@ def ad_revenue_report(
"matched": bool(rwd["matched"]),
"reward_detail": _reward_detail(rwd),
})
if (
rec.ad_type == "reward_video"
and rwd["status"] in _ZERO_REVENUE_REWARD_VIDEO_STATUSES
):
ev["revenue_yuan"] = 0.0
else:
# 纯展示(信息流逐条展示、激励视频缺发奖记录):不计对账,matched=True。
ev.update({
@@ -244,8 +262,8 @@ def ad_revenue_report(
"impressions": 0,
"ecpm": row["ecpm"],
"revenue_yuan": 0.0,
"adn": None,
"slot_id": None,
"adn": row.get("adn"),
"slot_id": row.get("slot_id"),
"has_reward": True,
"status": row["status"],
"expected_coin": int(row["expected_coin"]),
@@ -271,10 +289,10 @@ def ad_revenue_report(
# 父行 eCPM:组内各条 eCPM(分)均值(展示用,各条不同);无有效值则取代表条
ecpm_fens = [rewards.parse_ecpm_fen(g["ecpm"]) for g in group if g.get("ecpm")]
avg_ecpm = str(round(sum(ecpm_fens) / len(ecpm_fens))) if ecpm_fens else rep.get("ecpm")
# 主表逐行显示用:这次发奖广告的预估收益之和(发奖侧 eCPM 折算,钳顶同展示侧)。只放进
# 主表逐行显示用:这次发奖广告的预估收益之和(发奖侧 eCPM 折算)。只放进
# row_revenue_yuan 给主表逐行展示,不进 revenue_yuan/合计/趋势——避免与展示侧 total 重复计。
row_revenue = round(sum(
min(rewards.parse_ecpm_yuan(g["ecpm"]), rewards.AD_ECPM_MAX_FEN / 100.0) / 1000.0
rewards.parse_ecpm_yuan(g["ecpm"]) / 1000.0
for g in group if g.get("ecpm")
), 6)
events.append({
@@ -364,13 +382,15 @@ def ad_revenue_report(
for d in sorted(daily_map.values(), key=lambda x: x["date"])
]
# 穿山甲后台收益(GroMore 数据 API,T+1 入库 ad_pangle_daily_revenue):汇总 + 按天趋势级展示,
# GroMore 排序价预估 / ADN Reporting API 收益(T+1 入库):汇总 + 按天趋势级展示,
# 与上面客户端自报 eCPM 折算的预估并列对照(看 gap)。穿山甲数据**无用户/场景/类型维度**,故仅在
# 「全量视图」(未按 user_id / ad_type / feed_scene 过滤)给值;一旦带这些过滤,穿山甲数无法对应口径
# → 置 None,前端显示「-」并提示。逐条事件行不动(仍是客户端预估)。
pangle_filterable = user_id is None and ad_type is None and feed_scene is None
total_pangle_revenue_yuan: float | None = None
total_pangle_api_revenue_yuan: float | None = None
pangle_api_revenue_complete = False
pangle_latest_synced_at: datetime | None = None
if pangle_filterable:
pangle_aggs = ad_pangle_revenue.aggregate_by_date(
db,
@@ -388,6 +408,12 @@ def ad_revenue_report(
total_pangle_revenue_yuan = round(sum(a["revenue_yuan"] for a in pangle_aggs), 6)
api_vals = [a["api_revenue_yuan"] for a in pangle_aggs if a["api_revenue_yuan"] is not None]
total_pangle_api_revenue_yuan = round(sum(api_vals), 6) if api_vals else None
sync_times = [a["synced_at"] for a in pangle_aggs if a["synced_at"] is not None]
pangle_latest_synced_at = max(sync_times) if sync_times else None
pangle_api_revenue_complete = all(
day in by_date and _pangle_api_day_complete(day, by_date[day])
for day in _date_range(date_from, date_to)
)
# 按小时汇总(全量,不受分页 limit/offset 影响):供前端按小时趋势图(单日 granularity=hour 时用)。
# 只在 by_hour 下聚合(此时每个 event 带 hour);否则空。前端按天趋势仍用 daily。
@@ -412,41 +438,50 @@ def ad_revenue_report(
for hd in sorted(hour_map.values(), key=lambda x: x["hour"])
]
# 分广告类型小计(按 ad_type:展示条数 + 预估收益;eCPM 由前端用 收益÷展示×1000 算)。
# 基于全量(已按 feed_scene 过滤)events;前端只取 draw / reward_video 两类展示。
type_map: dict[str, dict] = {}
for e in events:
t = type_map.get(e["ad_type"])
if t is None:
t = {"impressions": 0, "revenue_yuan": 0.0}
type_map[e["ad_type"]] = t
t["impressions"] += e["impressions"]
t["revenue_yuan"] += e["revenue_yuan"]
type_stats = {
k: {"impressions": v["impressions"], "revenue_yuan": round(v["revenue_yuan"], 6)}
for k, v in type_map.items()
}
def _aggregate_stats(bucket_of) -> dict[str, dict]:
"""按展示事件聚合收益 / 加权 SDK eCPM,避免前端漏合并历史类型。"""
stat_map: dict[str, dict] = {}
for e in events:
bucket = bucket_of(e)
if bucket is None:
continue
stat = stat_map.setdefault(bucket, {
"impressions": 0,
"revenue_yuan": 0.0,
"ecpm_fen_sum": 0.0,
})
impressions = int(e["impressions"])
stat["impressions"] += impressions
stat["revenue_yuan"] += e["revenue_yuan"]
# eCPM 必须以每次真实展示为权重;纯发奖父行 impressions=0,不能参与分母或均值。
stat["ecpm_fen_sum"] += rewards.parse_ecpm_fen(e["ecpm"]) * impressions
return {
key: {
"impressions": value["impressions"],
"revenue_yuan": round(value["revenue_yuan"], 6),
"ecpm_yuan": round(
value["ecpm_fen_sum"] / value["impressions"] / 100.0,
6,
) if value["impressions"] else 0.0,
}
for key, value in stat_map.items()
}
# 分场景小计(按 feed_scene:展示条数 + 预估收益),同 type_stats 基于全量 events——
# 供数据大盘「领券广告 / 比价广告」卡用。此前大盘是在分页 items 里按 feed_scene 现算,
# 2026-07-02 起信息流逐条展示行(唯一带收益 + 场景的行)不再进主表 items,现算恒为 0;
# 改为服务端在全量上聚合下发(也顺带不受 limit 分页截断影响)。feed_scene 为空(激励视频 /
# 旧数据)不计入任何场景桶。
scene_map: dict[str, dict] = {}
for e in events:
sc = e.get("feed_scene")
if not sc:
continue
s = scene_map.get(sc)
if s is None:
s = {"impressions": 0, "revenue_yuan": 0.0}
scene_map[sc] = s
s["impressions"] += e["impressions"]
s["revenue_yuan"] += e["revenue_yuan"]
scene_stats = {
k: {"impressions": v["impressions"], "revenue_yuan": round(v["revenue_yuan"], 6)}
for k, v in scene_map.items()
}
# 原始 ad_type 小计,供明细筛选和排查使用。
type_stats = _aggregate_stats(lambda e: e["ad_type"])
# 经营看板使用的规范分类:Draw 包含历史 feed;看视频包含福利与提现视频。
# 这两个集合与筛选逻辑保持一致,避免只取 draw / reward_video 而漏算历史或提现数据。
category_stats = _aggregate_stats(
lambda e: (
"draw" if e["ad_type"] in {"draw", "feed"}
else "video" if e["ad_type"] in {"reward_video", "withdrawal_video"}
else None
)
)
# 分场景小计,同 type_stats 基于全量 events,供数据大盘「领券广告 / 比价广告」卡使用。
# feed_scene 为空的激励视频 / 历史数据不计入任何场景桶。
scene_stats = _aggregate_stats(lambda e: e.get("feed_scene"))
# DAU:复用数据大盘活跃用户口径(登录 + 开始比价 + 开始领券,按用户去重),按所选日期区间
# 统计(含今日),历史 / 多天区间同样有值。ARPU = 区间预估收益 ÷ 区间活跃用户。全局口径,
@@ -471,9 +506,11 @@ def ad_revenue_report(
"truncated": len(main_rows) > offset + limit,
"total_impressions": total_impressions,
"total_revenue_yuan": total_revenue_yuan,
# 穿山甲后台收益合计(元):预估 revenue + 收益Api;非全量视图(带 user/类型/场景过滤)或无数据为 None。
# GroMore 排序价预估 + ADN Reporting API 收益;非全量视图或无数据为 None。
"total_pangle_revenue_yuan": total_pangle_revenue_yuan,
"total_pangle_api_revenue_yuan": total_pangle_api_revenue_yuan,
"pangle_api_revenue_complete": pangle_api_revenue_complete,
"pangle_latest_synced_at": pangle_latest_synced_at,
"pangle_revenue_available": total_pangle_revenue_yuan is not None,
"total_expected_coin": total_expected_coin,
"total_actual_coin": total_actual_coin,
@@ -481,6 +518,7 @@ def ad_revenue_report(
"daily": daily,
"hourly": hourly,
"type_stats": type_stats,
"category_stats": category_stats,
"scene_stats": scene_stats,
"dau": dau,
"items": main_rows[offset:offset + limit],
+1 -9
View File
@@ -1228,11 +1228,6 @@ def _cn_wall_to_utc(dt: datetime) -> datetime:
return dt.replace(tzinfo=rewards.CN_TZ).astimezone(timezone.utc).replace(tzinfo=None)
def _coin_record_sort_key(row: dict) -> datetime:
"""金币明细跨数据源排序键:兼容 SQLite naive 与 PostgreSQL aware 时间。"""
return _as_utc(row["created_at"])
def user_coin_records(
db: Session,
user_id: int,
@@ -1324,10 +1319,7 @@ def user_coin_records(
"coin": rec.amount,
})
# SQLite 常返回 naive datetimePostgreSQL timestamptz 返回 aware datetime
# 统一成 aware UTC 排序,避免线上合并广告记录与签到记录时抛
# “can't compare offset-naive and offset-aware datetimes”。
rows.sort(key=_coin_record_sort_key, reverse=True)
rows.sort(key=lambda r: r["created_at"], reverse=True)
has_more = len(rows) > offset + limit
# 总数 = 三源在窗口内 granted 计数之和(供前端页码分页渲染页码/共 N 条)
+3
View File
@@ -99,6 +99,7 @@ def get_ad_revenue_report(
daily=[AdRevenueDaily(**d) for d in result["daily"]],
hourly=[AdRevenueHourly(**h) for h in result["hourly"]],
type_stats={k: AdRevenueTypeStat(**v) for k, v in result["type_stats"].items()},
category_stats={k: AdRevenueTypeStat(**v) for k, v in result["category_stats"].items()},
scene_stats={k: AdRevenueTypeStat(**v) for k, v in result["scene_stats"].items()},
dau=result["dau"],
total=result["total"],
@@ -107,6 +108,8 @@ def get_ad_revenue_report(
total_revenue_yuan=result["total_revenue_yuan"],
total_pangle_revenue_yuan=result["total_pangle_revenue_yuan"],
total_pangle_api_revenue_yuan=result["total_pangle_api_revenue_yuan"],
pangle_api_revenue_complete=result["pangle_api_revenue_complete"],
pangle_latest_synced_at=result["pangle_latest_synced_at"],
pangle_revenue_available=result["pangle_revenue_available"],
total_expected_coin=result["total_expected_coin"],
total_actual_coin=result["total_actual_coin"],
+1 -1
View File
@@ -86,7 +86,7 @@ def get_user_reward_stats(
date_to: Annotated[datetime | None, Query()] = None,
) -> UserRewardStats:
"""提现详情抽屉「用户统计区」。date_from/date_to 都不传 = 注册至今(全量)。"""
if not user_repo.user_exists(db, user_id):
if user_repo.get_user_by_id(db, user_id) is None:
raise HTTPException(status_code=404, detail="用户不存在")
return UserRewardStats(
**queries.user_reward_stats(db, user_id, date_from=date_from, date_to=date_to)
+27 -12
View File
@@ -40,6 +40,8 @@ class AdRevenueRecord(BaseModel):
expected_coin: int = Field(..., description="按公式复算应发金币")
actual_coin: int = Field(..., description="实际入账金币")
matched: bool = Field(..., description="复算与实发是否一致")
adn: str | None = Field(None, description="本条发奖对应的实际填充 ADN 子渠道")
slot_id: str | None = Field(None, description="本条发奖对应的底层 mediation rit")
class AdRevenueDaily(BaseModel):
@@ -47,12 +49,12 @@ class AdRevenueDaily(BaseModel):
date: str = Field(..., description="北京时间 YYYY-MM-DD")
impressions: int = Field(..., description="当天展示条数合计")
revenue_yuan: float = Field(..., description="当天客户端有效预估收益合计(元;eCPM 折算)")
revenue_yuan: float = Field(..., description="当天客户端 SDK 展示预估合计(元;后端留存 eCPM 折算)")
pangle_revenue_yuan: float | None = Field(
None, description="当天穿山甲后台预估收益(元;GroMore revenue);非全量视图/无数据为空"
None, description="当天 GroMore 排序价预估(元;revenue,非结算收入);非全量视图/无数据为空"
)
pangle_api_revenue_yuan: float | None = Field(
None, description="当天穿山甲收益Api(元;GroMore api_revenue,更接近结算);未配/当天/无数据为空"
None, description="当天 ADN Reporting API 收益(元;GroMore api_revenue);未配/当天/无数据为空"
)
expected_coin: int = Field(..., description="当天应发金币合计")
actual_coin: int = Field(..., description="当天实发金币合计")
@@ -69,10 +71,11 @@ class AdRevenueHourly(BaseModel):
class AdRevenueTypeStat(BaseModel):
"""按广告类型(ad_type)的小计:展示条数 + 预估收益(eCPM 由前端用 收益÷展示×1000 算)"""
"""展示条数、SDK 展示预估收益与按展示次数加权的 SDK eCPM"""
impressions: int = Field(..., description="该类型展示条数合计")
revenue_yuan: float = Field(..., description="该类型预估收益合计(元)")
ecpm_yuan: float = Field(..., description="按展示次数加权的 SDK eCPM(元/千次)")
class AdRevenueRow(BaseModel):
@@ -98,15 +101,15 @@ class AdRevenueRow(BaseModel):
ecpm: str | None = Field(None, description="eCPM 原始值(分/千次);展示行取展示值,纯发奖行取发奖采用值")
revenue_yuan: float = Field(
...,
description="本次有效展示预估收益(元)= eCPM元 ÷ 1000;纯发奖、激励视频提前关闭/时长不足=0",
description="本次 SDK 展示预估收益(元)=后端留存 eCPM 元 ÷ 1000;是否满足发奖条件不改变展示收入预估",
)
row_revenue_yuan: float | None = Field(
None,
description="主表逐行展示用的预估收益(元):一次比价/领券聚合行=该次发奖广告 eCPM 折算之和;"
"其它行为空(前端回退取 revenue_yuan)。不进合计/趋势,避免与展示侧重复计",
)
adn: str | None = Field(None, description="实际填充 ADN 子渠道(pangle/gdt…);纯发奖行为空")
slot_id: str | None = Field(None, description="底层 mediation rit(非我们配置的广告位 ID);纯发奖行为空")
adn: str | None = Field(None, description="实际填充 ADN 子渠道(pangle/gdt…);历史或未上报展示来源为空")
slot_id: str | None = Field(None, description="底层 mediation rit(非我们配置的广告位 ID);历史或未上报展示来源为空")
# ── 发奖侧 ──
has_reward: bool = Field(..., description="是否有发奖记录(激励视频合并行 / 信息流整场发奖行=True;纯展示=False)")
status: str | None = Field(None, description="发奖状态 granted/closed_early/too_short/…;纯展示为空")
@@ -140,7 +143,11 @@ class AdRevenueReportOut(BaseModel):
)
type_stats: dict[str, AdRevenueTypeStat] = Field(
default_factory=dict,
description="按广告类型(ad_type)小计 {ad_type: {impressions, revenue_yuan}};前端取 draw / reward_video 做分类大盘",
description="原始广告类型(ad_type)小计,供筛选与排查使用",
)
category_stats: dict[str, AdRevenueTypeStat] = Field(
default_factory=dict,
description="按经营分类小计:draw=draw+历史 feedvideo=reward_video+withdrawal_video",
)
scene_stats: dict[str, AdRevenueTypeStat] = Field(
default_factory=dict,
@@ -156,20 +163,28 @@ class AdRevenueReportOut(BaseModel):
total: int = Field(..., description="广告事件总数(全量,不受分页影响;= 当前筛选下的分页总条数)")
truncated: bool = Field(..., description="当前页之后是否还有更多事件(len(events) > offset + limit)")
total_impressions: int = Field(..., description="全量展示条数合计")
total_revenue_yuan: float = Field(..., description="全量客户端有效预估收益合计(元;eCPM 折算)")
total_revenue_yuan: float = Field(..., description="全量客户端 SDK 展示预估合计(元;后端留存 eCPM 折算)")
total_pangle_revenue_yuan: float | None = Field(
None,
description="全量穿山甲后台预估收益合计(元;GroMore revenue)。穿山甲无用户/类型/场景维度,"
description="全量 GroMore 排序价预估合计(元;revenue,非结算收入)。GroMore 无用户/类型/场景维度,"
"仅「全量视图」(未按 user_id/ad_type/feed_scene 过滤)时有值,否则为 null",
)
total_pangle_api_revenue_yuan: float | None = Field(
None,
description="全量穿山甲收益Api合计(元;GroMore api_revenue,各 ADN 回传、更接近结算);"
description="全量 ADN Reporting API 收益合计(元;GroMore api_revenue,仅已配置回传的 ADN);"
"未配 Reporting / 查当天 / 非全量视图 时为 null",
)
pangle_api_revenue_complete: bool = Field(
False,
description="所选每一天是否都已在 D+1 14:00 后完成 API 同步窗口;不代表覆盖全部 ADN 或最终结算",
)
pangle_latest_synced_at: datetime | None = Field(
None,
description="所选范围穿山甲/GroMore 日报最近同步时间",
)
pangle_revenue_available: bool = Field(
False,
description="本次结果是否带穿山甲后台收益(=全量视图且已同步到数据)。false 时前端「穿山甲收益」显示「-」",
description="本次结果是否带 GroMore/ADN 收益(=全量视图且已同步到数据)。false 时前端显示「-」",
)
total_expected_coin: int = Field(..., description="全量应发金币合计")
total_actual_coin: int = Field(..., description="全量实发金币合计")
+2 -1
View File
@@ -289,7 +289,8 @@ def ecpm_report(payload: EcpmReportIn, user: CurrentUser, db: DbSession) -> Ecpm
"""客户端在广告展示后(onAdShow 读 getShowEcpm)上报 eCPM,落库做内部收益统计/对账。
Bearer 鉴权,user_id 取自 JWT(不信 body)。best-effort:落库即 ok,客户端 fire-and-forget,
丢一两条不影响业务(穿山甲后台报表是结算权威)。eCPM 与发奖(S2S)是两条独立流,不逐条关联。
丢一两条不影响发奖业务(收入另由 ADN Reporting API 对账)。eCPM 与发奖(S2S)是两条独立流,
不逐条关联。
"""
attributed_trace_id = crud_ecpm.attributable_trace_id(
db,
+1 -1
View File
@@ -318,7 +318,7 @@ class Settings(BaseSettings):
# ===== 穿山甲 GroMore 数据 API(报表收益拉取,T+1)=====
# ⚠️ 与上面发奖回调的 m-key 是【两套完全不同的凭证】:这三样在穿山甲后台
# 「接入中心 → GroMore-API → 聚合数据报告 API」文档页领取(user_id / role_id / Security Key),
# 仅用于按天拉 GroMore 收益报表(revenue 预估收益 + api_revenue 收益Api),不参与发奖。
# 仅用于按天拉 GroMore 报表(revenue 排序价预估 + api_revenue ADN Reporting 收益),不参与发奖。
# 该 API 只能查【GroMore 聚合代码位】的数据(=我们 useMediation 的口径),非穿山甲 SDK 数据;
# 且不提供用户/设备维度(官方明确),故收益只能落到 日期×代码位 汇总,不能挂到逐条事件。
# 子账号(role_id≠user_id)需主账号在「角色管理」授予「查看全部数据」权限,否则查不到
+2 -2
View File
@@ -14,8 +14,8 @@
- 只返回【GroMore 聚合代码位】在 GroMore 内的数据(=我们 useMediation 的口径),
查不到穿山甲 SDK 自身的数据;
- **不提供分用户/设备维度**(官方 FAQ 明确拒绝),最细到 日期×应用×代码位×广告源;
- `revenue` = 预估收益(元,所有 ADN 都有);`api_revenue` = 收益Api(各 ADN 经 Reporting
回传、按实时汇率折算账号币种,更接近结算),需后台为该 ADN 配置 Reporting 才有、且不支持当天;
- `revenue` = 排序价/竞价实时价预估(元,非结算收入);`api_revenue` = 各 ADN 经 Reporting
回传、按实时汇率折算账号币种的收益,需后台为该 ADN 配置 Reporting 才有、且不支持当天;
- 「今天」与「今天以前」必须分开查;天级跨度 ≤ 1 个月、不早于 12 个月。
"""
from __future__ import annotations
+8 -8
View File
@@ -1,15 +1,15 @@
"""穿山甲 GroMore 天级收益报表(后台结算口径,定时拉取入库)。
"""GroMore 天级排序价预估与 ADN Reporting 收益(定时拉取入库)。
每行 = GroMore 数据 API 返回的一条「日期 × 应用 × 代码位」聚合收益(`integrations/pangle_report`
+ `scripts/sync_pangle_revenue` 落库)。**权威/预估收益的来源**,与 `ad_ecpm_record`(客户端自报
eCPM 折算的预估)互为对照:
+ `scripts/sync_pangle_revenue` 落库)与 `ad_ecpm_record`(客户端 SDK eCPM 折算的预估)
互为对照:
- `revenue_yuan` ← 接口 `revenue`(预估收益,元;排序价×展示/1000,所有 ADN 都有);
- `api_revenue_yuan` ← 接口 `api_revenue`(收益Api,元;各 ADN Reporting 回传更接近结算;
- `revenue_yuan` ← 接口 `revenue`(排序价/竞价实时价预估,元,不是结算收入);
- `api_revenue_yuan` ← 接口 `api_revenue`(各 ADN Reporting 回传收益,元,更接近结算;
未配置该 ADN 的 Reporting 或查当天时为空)。
⚠️ 穿山甲不提供分用户/设备维度,故本表最细只到 日期×应用×代码位,**无法挂到逐条广告事件**;
广告收益报表里只用于汇总/趋势级的「穿山甲后台收益」,不改逐条行的客户端预估。
广告收益报表里只用于汇总/趋势级的 GroMore/ADN 对账,不改逐条行的客户端预估。
"""
from __future__ import annotations
@@ -51,9 +51,9 @@ class AdPangleDailyRevenue(Base):
our_code_id: Mapped[str] = mapped_column(String(64), index=True, nullable=False)
# 广告源(接口 network 数字→名,如 pangle/gdt);"" = 未分广告源的代码位汇总行(当前默认口径)。
adn: Mapped[str] = mapped_column(String(16), nullable=False, default="")
# 预估收益(元)← 接口 revenue。
# 排序价/竞价实时价预估(元)← 接口 revenue,非结算收入
revenue_yuan: Mapped[float] = mapped_column(Float, nullable=False, default=0.0)
# 收益Api(元)← 接口 api_revenue;未配 Reporting / 当天 等情况接口不返回 → NULL。
# ADN Reporting API 收益(元)← api_revenue;未配 Reporting / 当天等情况不返回 → NULL。
api_revenue_yuan: Mapped[float | None] = mapped_column(Float, nullable=True)
# 预估 eCPM 原值(接口 ecpm,单位元/千次,**与客户端 getEcpm 的「分」不同**),参考用原样存。
ecpm: Mapped[str | None] = mapped_column(String(32), nullable=True)
+3 -3
View File
@@ -1,8 +1,8 @@
"""广告 eCPM 上报 CRUD(内部收益统计/对账)。
客户端在广告展示后(onAdShow)读到 eCPM,经鉴权接口上报,这里落库。鉴权接口已确保
user 存在(JWT),故不做 UnknownUser 校验。best-effort 上报:丢一两条不影响业务,
穿山甲后台报表是结算权威兜底
user 存在(JWT),故不做 UnknownUser 校验。best-effort 上报:丢一两条不影响发奖业务;
汇总收入以 ADN Reporting API 和最终结算单为准
"""
from __future__ import annotations
@@ -96,7 +96,7 @@ def create_ecpm_record(
db.rollback()
# 撞唯一约束 uq_ad_ecpm_record_session(全局按 ad_session_id、不含 user_id):并发同会话重复上报,
# 或同一 ad_session_id 已被先到的上报占用。本接口 fire-and-forget、best-effort —— 丢一条不影响业务
# (穿山甲后台才是结算权威),绝不向客户端抛 500。兜底查找须与唯一约束**同口径**(只按 ad_session_id、
# (收入另由 ADN Reporting API 对账),绝不向客户端抛 500。兜底查找须与唯一约束**同口径**(只按 ad_session_id、
# 不带 user_id):否则不同 user 上报了同一 ad_session_id 时,带 user_id 的查找会漏掉那条别人的记录 →
# 旧逻辑在此 raise 成 500(本应静默吞掉)。
existing = _find_by_session_global(db, ad_session_id)
+6 -2
View File
@@ -1,12 +1,13 @@
"""穿山甲 GroMore 天级收益 读写(`ad_pangle_daily_revenue` 表)。
`scripts/sync_pangle_revenue` 拉数后调 `upsert_daily_rows` 落库(同一(日期×应用×代码位×广告源)
幂等覆盖,T+1 订正可重跑);admin 广告收益报表调 `aggregate_by_date` 取「穿山甲后台收益
幂等覆盖,T+1 订正可重跑);admin 广告收益报表调 `aggregate_by_date` 取 GroMore/ADN 收益做
汇总/趋势级展示。穿山甲无用户维度,故这里不涉及 user_id。
"""
from __future__ import annotations
from collections.abc import Collection
from datetime import datetime
from typing import Any, TypedDict
from sqlalchemy import func, select
@@ -23,6 +24,7 @@ class PangleDateAgg(TypedDict):
revenue_yuan: float
api_revenue_yuan: float | None
impressions: int
synced_at: datetime | None
def upsert_daily_rows(db: Session, rows: list[dict[str, Any]]) -> dict[str, int]:
@@ -87,6 +89,7 @@ def aggregate_by_date(
func.sum(AdPangleDailyRevenue.revenue_yuan),
func.sum(AdPangleDailyRevenue.api_revenue_yuan),
func.sum(AdPangleDailyRevenue.impressions),
func.max(AdPangleDailyRevenue.synced_at),
)
.where(
AdPangleDailyRevenue.report_date >= date_from,
@@ -103,11 +106,12 @@ def aggregate_by_date(
stmt = stmt.where(AdPangleDailyRevenue.our_code_id.in_(our_code_ids))
out: list[PangleDateAgg] = []
for report_date, rev, api_rev, imp in db.execute(stmt).all():
for report_date, rev, api_rev, imp, synced_at in db.execute(stmt).all():
out.append(PangleDateAgg(
date=report_date,
revenue_yuan=round(float(rev or 0.0), 6),
api_revenue_yuan=(round(float(api_rev), 6) if api_rev is not None else None),
impressions=int(imp or 0),
synced_at=synced_at,
))
return out
-5
View File
@@ -93,11 +93,6 @@ def get_user_by_id(db: Session, user_id: int) -> User | None:
return db.get(User, user_id)
def user_exists(db: Session, user_id: int) -> bool:
"""只查主键判断用户是否存在,避免只读统计接口依赖完整用户表结构。"""
return db.scalar(select(User.id).where(User.id == user_id)) is not None
def get_user_by_phone(db: Session, phone: str) -> User | None:
stmt = select(User).where(User.phone == phone)
return db.execute(stmt).scalar_one_or_none()
+10 -10
View File
@@ -1,13 +1,13 @@
# 穿山甲 GroMore 收益拉取 定时任务 — 运维手册
> 对象:维护「每天拉穿山甲后台收益入库」这套定时任务的同事。
> 对象:维护「每天拉 GroMore / ADN 收益入库」这套定时任务的同事。
> 🔒 服务器登录信息见**私密交接清单**,不入库。
## 它是什么
admin「广告收益报表」里的「穿山甲后台收益(T+1)」读的是**本地表 `ad_pangle_daily_revenue` 的快照,不是实时查穿山甲**。穿山甲只通过 GroMore 数据 API 给数、且 **T+1**(次日约 10:00 出昨天的数),所以每天得拉一次入库,报表才会往前走
admin「广告收益报表」里的 GroMore / ADN 收益读的是**本地表 `ad_pangle_daily_revenue` 的快照,不是实时查询**。GroMore 的 T+1 初值约 10:00 可用,但第三方 ADN Reporting 数据可能到 13:50 才更新,所以需要早晚各拉一次
- 每天 10:30 跑一轮 `scripts/sync_pangle_revenue.py`,默认 `--days 3` 回补近 3 天。
- 维度 = 日期 × 应用(site_id)× 广告位(ad_unit_id);指标 = `revenue`(预估)+ `api_revenue`(结算口径)。
- 每天 10:30 拉初值、14:30 拉日终值,均由 `scripts/sync_pangle_revenue.py` `--days 3` 回补近 3 天。
- 维度 = 日期 × 应用(site_id)× 广告位(ad_unit_id);指标 = `revenue`(排序价预估)+ `api_revenue`(ADN Reporting 回传,更接近结算)。
- **幂等 upsert**:同一(日期×应用×代码位)重跑只覆盖、不重复,故回补 / 重跑 / catch-up 都安全。
- 穿山甲无用户/设备维度 → 只能落「汇总/趋势级」,报表带 user_id 过滤时这块收益置空(显示「-」)。
@@ -30,15 +30,15 @@ admin「广告收益报表」里的「穿山甲后台收益(T+1)」读的是**
```bash
sudo cp deploy/pangle-revenue.{service,timer} /etc/systemd/system/
sudo systemctl daemon-reload && sudo systemctl enable --now pangle-revenue.timer
systemctl list-timers pangle-revenue.timer # 确认下次触发时间(应是次日 10:30)
systemctl list-timers pangle-revenue.timer # 确认下次触发时间(10:30 或 14:30)
```
## 怎么看健康 / 手动跑一次
```bash
sudo systemctl start pangle-revenue.service # 立即手动跑一轮(不等 10:30)
journalctl -u pangle-revenue -n 30 --no-pager # 看日志:拉取区间 / 入库行数 / 新增更新 / 预估收益合计
journalctl -u pangle-revenue -n 30 --no-pager # 看日志:拉取区间 / 入库行数 / 新增更新 / 收益合计
sudo systemctl start pangle-revenue.service # 立即手动跑一轮
```
成功日志形如:`✅ 完成:接口 N 行 → 入库 M 行(跳过 x),新增 a / 更新 b;预估收益合计 ¥19.42`
成功日志形如:`✅ 完成:接口 N 行 → 入库 M 行(跳过 x),新增 a / 更新 b;排序价预估合计 ¥19.42`
> 看不到收益、提示 `PANGLE_REPORT_* 未配置`→ 回「上线前置」补 `.env`;报 118 → 子账号没授「查看全部数据」。
## 本机 Windows 开发(无 systemd)
@@ -57,11 +57,11 @@ journalctl -u pangle-revenue -n 30 --no-pager # 看日志:拉取区间 / 入
- `--start / --end`:指定闭区间(跨度 ≤ 31 天,接口上限 1 个月,超了报 114)。
## 注意事项
- **触发时间**:`OnCalendar=*-*-* 10:30:00`。穿山甲 ~10:00 出数,故别早于 10:00 跑(会拉到空/不全)
- **触发时间**:10:30 提供初值,14:30 覆盖为日终值;报表只把 D+1 14:00 后同步的数据标记为日终
- **catch-up**:`Persistent=true` 补跑错过的那一轮;叠加 `--days 3`,漏一两天重新触发即自愈。
- **今天 / 今天以前要分开查**:脚本默认只拉昨天及更早,不混查今天(接口约束),无需关心。
- **join key 是 `ad_unit_id`(我们配的 104xxx)不是 `code_id`**:`code_id` 是底层各 ADN 代码位,对不上口径;`ad_unit_id='-1'` 是未归因桶。改维度时务必注意(详见脚本头注释)。
- **`api_revenue` 很稀疏**:测试应用 ADN 没配 Reporting → 全 0,仅 prod 个别位有;`revenue`(预估)才是稳的主力
- **`api_revenue` 依赖 ADN Reporting 配置**:未配置的测试应用可能为空或 0;`revenue` 只是排序价估算,不能当结算收入
- **DB 无关**:sqlite / postgres 均可(upsert 逐行 select-then-write,不像美团 ETL 需要 PG)。
- **别和别的触发方式双跑**:本 systemd timer 与「手动 cron / 进程内任务」二选一,虽幂等不会重复入库,纯属多余。
- **改脚本 / 改部署**:走 git + PR,由有 root 的人部署。
+2 -2
View File
@@ -1,5 +1,5 @@
# 每天拉穿山甲 GroMore T+1 天级收益入库 —— 单轮跑,由 pangle-revenue.timer 每天 10:30 触发。
# 落 ad_pangle_daily_revenue 表,供 admin 广告收益报表的「穿山甲后台收益(T+1)」区块。
# GroMore T+1 天级收益入库 —— 单轮跑,由 timer 每天 10:30、14:30 触发。
# 落 ad_pangle_daily_revenue 表,供 admin 广告收益报表的 GroMore/ADN 对账区块。
#
# 仅用于 Linux 服务器;本机 Windows 开发无 systemd,直接手动跑脚本即可:
# .venv\Scripts\python -m scripts.sync_pangle_revenue # 拉昨天(北京时间)
+4 -3
View File
@@ -1,11 +1,12 @@
# 每天 10:30 触发一次穿山甲 GroMore T+1 收益拉取入库(Linux 服务器用)。
# 每天 10:30 首次拉取、14:30 终值复拉 GroMore T+1 收益(Linux 服务器用)。
# 见 pangle-revenue.service 顶部注释的部署步骤。
[Unit]
Description=Run Pangle GroMore daily revenue sync at 10:30
Description=Run Pangle GroMore daily revenue sync at 10:30 and 14:30
[Timer]
# 穿山甲 T+1、次日约 10:00 出数;10:30 触发留 30min 余量。要错开整点扎堆可微调到 10:35
# 10:30 尽早展示初值;第三方 ADN Reporting 最晚约 13:50 更新,14:30 再拉一次作为日终值
OnCalendar=*-*-* 10:30:00
OnCalendar=*-*-* 14:30:00
# 服务器宕机/重启后,补跑错过的那一轮(而不是干等次日);叠加 --days 3 回补,漏一两天能自愈。
Persistent=true
AccuracySec=1min
+211
View File
@@ -0,0 +1,211 @@
"""生成本地 admin「广告收益」与数据大盘用的可重复 mock 数据。
只处理 ``local-admin-revenue-mock-`` 前缀的数据,重跑会替换自身数据,不会触碰真实本地记录。
会覆盖 Draw(含一条历史 feed)、福利激励视频、提现视频,以及对应的金币流水。
用法:
python -m scripts.seed_admin_revenue_mock
"""
from __future__ import annotations
from datetime import datetime, time, timedelta
from sqlalchemy import delete, select
from app.core.config import settings
from app.core.rewards import CN_TZ, cn_today
from app.db.session import SessionLocal
from app.models.ad_ecpm import AdEcpmRecord
from app.models.ad_feed_reward import AdFeedRewardRecord
from app.models.ad_reward import AdRewardRecord
from app.models.user import User
from app.models.wallet import CoinAccount, CoinTransaction
PREFIX = "local-admin-revenue-mock-"
PHONE = "19900009002"
USERNAME = "80000009002"
def _at(day_offset: int, hour: int, minute: int) -> datetime:
day = cn_today() - timedelta(days=day_offset)
return datetime.combine(day, time(hour, minute), tzinfo=CN_TZ)
def _add_coin(
db,
*,
user_id: int,
amount: int,
biz_type: str,
ref_id: str,
created_at: datetime,
balance_after: int,
) -> None:
db.add(CoinTransaction(
user_id=user_id,
amount=amount,
balance_after=balance_after,
biz_type=biz_type,
ref_id=ref_id,
remark="本地运营后台广告收益 Mock",
created_at=created_at,
))
def seed() -> dict[str, int]:
if settings.APP_ENV == "prod":
raise RuntimeError("Refusing to seed admin revenue mock data in production")
with SessionLocal() as db:
# 清理顺序按外键依赖从流水/奖励到展示;只碰本脚本自己的稳定前缀。
db.execute(delete(CoinTransaction).where(CoinTransaction.ref_id.like(f"{PREFIX}%")))
db.execute(delete(AdFeedRewardRecord).where(
AdFeedRewardRecord.client_event_id.like(f"{PREFIX}%")
))
db.execute(delete(AdRewardRecord).where(AdRewardRecord.trans_id.like(f"{PREFIX}%")))
db.execute(delete(AdEcpmRecord).where(AdEcpmRecord.ad_session_id.like(f"{PREFIX}%")))
user = db.execute(select(User).where(User.phone == PHONE)).scalar_one_or_none()
if user is None:
user = User(
phone=PHONE,
username=USERNAME,
nickname="运营收益 Mock 用户",
register_channel="sms",
status="active",
)
db.add(user)
db.flush()
else:
user.nickname = "运营收益 Mock 用户"
user.status = "active"
balance = 0
event_count = 0
reward_count = 0
# 近四天的数据既能覆盖单日,也能覆盖近 7 天趋势与分类合计。
for day_offset in range(4):
suffix = f"d{day_offset}"
compare_trace = "local-invite-mock-compare-success"
coupon_trace = "mock-coupon-repeat-prod-second"
draw_events = [
("draw", "comparison", compare_trace, 2800 + day_offset * 100, 10, 10),
("draw", "coupon", coupon_trace, 1750 + day_offset * 100, 10, 28),
]
# 历史 feed 必须被 Draw 分类一起计算,用于走查兼容逻辑。
if day_offset == 1:
draw_events.append(("feed", "coupon", coupon_trace, 1250, 11, 12))
for index, (ad_type, scene, trace_id, ecpm, hour, minute) in enumerate(draw_events, start=1):
session = f"{PREFIX}{suffix}-draw-{index}"
created_at = _at(day_offset, hour, minute)
db.add(AdEcpmRecord(
user_id=user.id,
ad_type=ad_type,
feed_scene=scene,
trace_id=trace_id,
ad_session_id=session,
adn="pangle" if index == 1 else "gdt",
slot_id="mock-draw-rit",
app_env="prod",
our_code_id="104098712",
ecpm_raw=str(ecpm),
report_date=created_at.date().isoformat(),
created_at=created_at,
))
coin = 18 + day_offset * 2
db.add(AdFeedRewardRecord(
client_event_id=f"{PREFIX}{suffix}-feed-reward-{index}",
ad_session_id=session,
user_id=user.id,
reward_date=created_at.date().isoformat(),
duration_seconds=20,
unit_count=2,
ecpm_raw=str(ecpm),
adn="pangle" if index == 1 else "gdt",
slot_id="mock-draw-rit",
ad_type=ad_type,
feed_scene=scene,
trace_id=trace_id,
app_env="prod",
our_code_id="104098712",
coin=coin,
status="granted",
created_at=created_at + timedelta(seconds=20),
))
balance += coin
_add_coin(
db,
user_id=user.id,
amount=coin,
biz_type="feed_ad_reward",
ref_id=f"{PREFIX}{suffix}-feed-coin-{index}",
created_at=created_at + timedelta(seconds=20),
balance_after=balance,
)
event_count += 1
reward_count += 1
for ad_type, ecpm, hour, coin in (
("reward_video", 13200 + day_offset * 500, 14, 66),
("withdrawal_video", 32000 + day_offset * 800, 18, 0),
):
session = f"{PREFIX}{suffix}-{ad_type}"
created_at = _at(day_offset, hour, 6)
db.add(AdEcpmRecord(
user_id=user.id,
ad_type=ad_type,
ad_session_id=session,
adn="ks" if ad_type == "reward_video" else "baidu",
slot_id="mock-video-rit",
app_env="prod",
our_code_id="104099389",
ecpm_raw=str(ecpm),
report_date=created_at.date().isoformat(),
created_at=created_at,
))
event_count += 1
if ad_type == "reward_video":
db.add(AdRewardRecord(
trans_id=f"{PREFIX}{suffix}-reward-video",
user_id=user.id,
coin=coin,
status="granted",
reward_scene="reward_video",
ad_session_id=session,
ecpm_raw=str(ecpm),
app_env="prod",
our_code_id="104099389",
reward_date=created_at.date().isoformat(),
reward_name="Mock 福利视频",
created_at=created_at + timedelta(seconds=35),
))
balance += coin
_add_coin(
db,
user_id=user.id,
amount=coin,
biz_type="reward_video",
ref_id=f"{PREFIX}{suffix}-reward-video-coin",
created_at=created_at + timedelta(seconds=35),
balance_after=balance,
)
reward_count += 1
account = db.get(CoinAccount, user.id)
if account is None:
account = CoinAccount(user_id=user.id)
db.add(account)
account.coin_balance = balance
account.total_coin_earned = balance
db.commit()
return {"events": event_count, "rewards": reward_count, "coin": balance}
if __name__ == "__main__":
result = seed()
print(
"Seeded local admin revenue mock: "
f"{result['events']} impressions, {result['rewards']} rewards, {result['coin']} coins"
)
+4 -4
View File
@@ -1,7 +1,7 @@
"""每日拉取穿山甲 GroMore 天级收益报表入库(供 admin 广告收益报表的「穿山甲后台收益」)。
"""每日拉取 GroMore 排序价预估与 ADN Reporting 收益入库(供 admin 广告收益对账)。
GroMore 数据 API 为 T+1:次日穿山甲约 10:00 出数。建议线上每天 ~10:30 由 systemd timer 跑一次
(默认拉【昨天】);穿山甲对历史数据可能订正,故支持回补近 N 天(幂等 upsert,重跑无害)
GroMore 数据 API 为 T+1:次日约 10:00 出初值,第三方 ADN Reporting 最晚约 13:50 更新。
线上由 systemd timer 在 10:30、14:30 各跑一次;历史数据可能订正,故支持回补近 N 天。
用法:
python -m scripts.sync_pangle_revenue # 拉昨天(北京时间)
@@ -117,7 +117,7 @@ def sync_range(start_date: str, end_date: str) -> None:
stats = repo.upsert_daily_rows(db, rows)
total_rev = round(sum(r["revenue_yuan"] for r in rows), 4)
print(f"✅ 完成:接口 {len(raw)} 行 → 入库 {len(rows)} 行(跳过 {skipped}),"
f"新增 {stats['inserted']} / 更新 {stats['updated']};预估收益合计 ¥{total_rev}")
f"新增 {stats['inserted']} / 更新 {stats['updated']};排序价预估合计 ¥{total_rev}")
def main() -> None:
+245 -10
View File
@@ -8,12 +8,15 @@ from sqlalchemy import delete
from app.admin.repositories import ad_revenue
from app.db.session import SessionLocal
from app.models.ad_ecpm import AdEcpmRecord
from app.models.ad_feed_reward import AdFeedRewardRecord
from app.models.ad_pangle_revenue import AdPangleDailyRevenue
from app.models.ad_reward import AdRewardRecord
from app.models.user import User
REPORT_DATE = "2040-02-03"
PLAYBACK_DATE = "2040-02-04"
DETAIL_DATE = "2040-02-06"
SOURCE_FALLBACK_DATE = "2040-02-07"
def test_business_scope_filters_client_and_pangle_by_env_and_code(monkeypatch) -> None:
@@ -50,18 +53,22 @@ def test_business_scope_filters_client_and_pangle_by_env_and_code(monkeypatch) -
AdPangleDailyRevenue(
report_date=REPORT_DATE, app_env="prod", our_code_id="prod-reward",
adn="", revenue_yuan=1.5, api_revenue_yuan=1.2, impressions=10,
synced_at=datetime(2040, 2, 4, 6, 30, tzinfo=UTC),
),
AdPangleDailyRevenue(
report_date=REPORT_DATE, app_env="prod", our_code_id="prod-demo",
adn="", revenue_yuan=8.0, api_revenue_yuan=7.0, impressions=40,
synced_at=datetime(2040, 2, 4, 6, 30, tzinfo=UTC),
),
AdPangleDailyRevenue(
report_date=REPORT_DATE, app_env="prod", our_code_id="104098712",
adn="", revenue_yuan=2.5, api_revenue_yuan=2.0, impressions=20,
synced_at=datetime(2040, 2, 4, 6, 30, tzinfo=UTC),
),
AdPangleDailyRevenue(
report_date=REPORT_DATE, app_env="test", our_code_id="104127529",
adn="", revenue_yuan=9.0, api_revenue_yuan=8.0, impressions=50,
synced_at=datetime(2040, 2, 4, 6, 30, tzinfo=UTC),
),
])
db.commit()
@@ -86,6 +93,8 @@ def test_business_scope_filters_client_and_pangle_by_env_and_code(monkeypatch) -
assert business["total_revenue_yuan"] == 0.5
assert business["total_pangle_revenue_yuan"] == 4.0
assert business["total_pangle_api_revenue_yuan"] == 3.2
assert business["pangle_api_revenue_complete"] is True
assert business["pangle_latest_synced_at"] is not None
all_codes = ad_revenue.ad_revenue_report(
db,
@@ -131,7 +140,7 @@ def test_business_scope_filters_client_and_pangle_by_env_and_code(monkeypatch) -
db.close()
def test_reward_video_incomplete_playback_has_zero_revenue() -> None:
def test_reward_video_impression_revenue_is_independent_of_reward_status_and_cap() -> None:
db = SessionLocal()
phone = "18800009992"
sessions = {
@@ -147,13 +156,17 @@ def test_reward_video_incomplete_playback_has_zero_revenue() -> None:
for index, (status, session_id) in enumerate(sessions.items(), start=1):
created_at = datetime(2040, 2, 4, index, tzinfo=UTC)
# 发奖状态 capped 的广告故意使用 ¥1000 CPM,验证收入不套用金币侧 ¥500 CPM 封顶。
ecpm_raw = "100000" if status == "capped" else "10000"
db.add(AdEcpmRecord(
user_id=user.id,
ad_type="reward_video",
ad_session_id=session_id,
adn=f"adn-{status}",
slot_id=f"rit-{status}",
app_env="prod",
our_code_id="prod-reward",
ecpm_raw="10000",
ecpm_raw=ecpm_raw,
report_date=PLAYBACK_DATE,
created_at=created_at,
))
@@ -166,7 +179,7 @@ def test_reward_video_incomplete_playback_has_zero_revenue() -> None:
ad_session_id=session_id,
app_env="prod",
our_code_id="prod-reward",
ecpm_raw="10000",
ecpm_raw=ecpm_raw,
reward_date=PLAYBACK_DATE,
created_at=created_at,
))
@@ -185,22 +198,33 @@ def test_reward_video_incomplete_playback_has_zero_revenue() -> None:
revenue_by_status = {row["status"]: row["revenue_yuan"] for row in result["items"]}
assert revenue_by_status == {
"closed_early": 0.0,
"too_short": 0.0,
"capped": 0.1,
"closed_early": 0.1,
"too_short": 0.1,
"capped": 1.0,
"granted": 0.1,
}
assert result["total_impressions"] == 4
assert result["total_revenue_yuan"] == 0.2
assert result["total_revenue_yuan"] == 1.3
assert len(result["daily"]) == 1
assert result["daily"][0]["date"] == PLAYBACK_DATE
assert result["daily"][0]["impressions"] == 4
assert result["daily"][0]["revenue_yuan"] == 0.2
assert sum(row["revenue_yuan"] for row in result["hourly"]) == 0.2
assert result["daily"][0]["revenue_yuan"] == 1.3
assert sum(row["revenue_yuan"] for row in result["hourly"]) == 1.3
assert result["type_stats"]["reward_video"] == {
"impressions": 4,
"revenue_yuan": 0.2,
"revenue_yuan": 1.3,
"ecpm_yuan": 325.0,
}
detail_by_status = {
row["status"]: row["reward_detail"] for row in result["items"]
}
assert detail_by_status["granted"]["adn"] == "adn-granted"
assert detail_by_status["granted"]["slot_id"] == "rit-granted"
# 未进入发奖的记录可保留展示收入,但不能凭空生成奖励因子或占用 LT 累计。
for status in ("closed_early", "too_short", "capped"):
assert detail_by_status[status]["ecpm_factor"] is None
assert detail_by_status[status]["lt_factor_start"] is None
assert detail_by_status[status]["lt_index_start"] is None
finally:
db.rollback()
db.execute(delete(AdRewardRecord).where(AdRewardRecord.reward_date == PLAYBACK_DATE))
@@ -208,3 +232,214 @@ def test_reward_video_incomplete_playback_has_zero_revenue() -> None:
db.execute(delete(User).where(User.phone == phone))
db.commit()
db.close()
def test_category_stats_merge_legacy_feed_and_withdrawal_video() -> None:
db = SessionLocal()
phone = "18800009993"
category_date = "2040-02-05"
try:
user = User(phone=phone, username="29999999993", register_channel="sms")
db.add(user)
db.flush()
db.add_all([
# Draw 经营分类必须包含新 draw 与历史 feed。
AdEcpmRecord(
user_id=user.id, ad_type="draw", ad_session_id="category-draw",
app_env="prod", our_code_id="prod-draw", ecpm_raw="10000",
report_date=category_date, created_at=datetime(2040, 2, 5, 1, tzinfo=UTC),
),
AdEcpmRecord(
user_id=user.id, ad_type="feed", ad_session_id="category-feed",
app_env="prod", our_code_id="prod-draw", ecpm_raw="20000",
report_date=category_date, created_at=datetime(2040, 2, 5, 2, tzinfo=UTC),
),
# 看视频经营分类必须包含福利与提现两个视频入口。
AdEcpmRecord(
user_id=user.id, ad_type="reward_video", ad_session_id="category-reward",
app_env="prod", our_code_id="prod-reward", ecpm_raw="30000",
report_date=category_date, created_at=datetime(2040, 2, 5, 3, tzinfo=UTC),
),
AdEcpmRecord(
user_id=user.id, ad_type="withdrawal_video", ad_session_id="category-withdraw",
app_env="prod", our_code_id="prod-reward", ecpm_raw="50000",
report_date=category_date, created_at=datetime(2040, 2, 5, 4, tzinfo=UTC),
),
])
db.commit()
result = ad_revenue.ad_revenue_report(
db,
date_from=category_date,
date_to=category_date,
user_id=user.id,
app_env="prod",
revenue_scope="all",
)
assert result["category_stats"] == {
"draw": {"impressions": 2, "revenue_yuan": 0.3, "ecpm_yuan": 150.0},
"video": {"impressions": 2, "revenue_yuan": 0.8, "ecpm_yuan": 400.0},
}
finally:
db.rollback()
db.execute(delete(AdEcpmRecord).where(AdEcpmRecord.report_date == category_date))
db.execute(delete(User).where(User.phone == phone))
db.commit()
db.close()
def test_feed_reward_detail_keeps_each_record_adn_instead_of_parent_adn() -> None:
db = SessionLocal()
phone = "18800009994"
try:
user = User(phone=phone, username="29999999994", register_channel="sms")
db.add(user)
db.flush()
db.add_all([
AdFeedRewardRecord(
client_event_id="detail-adn-pangle",
user_id=user.id,
reward_date=DETAIL_DATE,
duration_seconds=20,
unit_count=1,
ecpm_raw="12000",
adn="pangle",
slot_id="rit-pangle",
ad_type="draw",
feed_scene="coupon",
trace_id="detail-adn-trace",
app_env="prod",
our_code_id="104098712",
coin=12,
status="granted",
created_at=datetime(2040, 2, 6, 1, tzinfo=UTC),
),
AdFeedRewardRecord(
client_event_id="detail-adn-gdt",
user_id=user.id,
reward_date=DETAIL_DATE,
duration_seconds=20,
unit_count=1,
ecpm_raw="25000",
adn="gdt",
slot_id="rit-gdt",
ad_type="draw",
feed_scene="coupon",
trace_id="detail-adn-trace",
app_env="prod",
our_code_id="104098712",
coin=25,
status="granted",
created_at=datetime(2040, 2, 6, 2, tzinfo=UTC),
),
])
db.commit()
result = ad_revenue.ad_revenue_report(
db,
date_from=DETAIL_DATE,
date_to=DETAIL_DATE,
user_id=user.id,
app_env="prod",
revenue_scope="all",
)
item = next(row for row in result["items"] if row["event_key"].startswith("feedgrp-"))
assert item["adn"] is None
assert [detail["adn"] for detail in item["sub_rewards"]] == ["pangle", "gdt"]
assert [detail["slot_id"] for detail in item["sub_rewards"]] == ["rit-pangle", "rit-gdt"]
finally:
db.rollback()
db.execute(delete(AdFeedRewardRecord).where(AdFeedRewardRecord.reward_date == DETAIL_DATE))
db.execute(delete(User).where(User.phone == phone))
db.commit()
db.close()
def test_feed_reward_source_can_fallback_to_unique_trace_impression_only() -> None:
"""发奖会话是整场 ID、展示会话是 impressionId 时,只在 trace+eCPM 唯一时回填 ADN。"""
db = SessionLocal()
phone = "18800009995"
try:
user = User(phone=phone, username="29999999995", register_channel="sms")
db.add(user)
db.flush()
db.add_all([
AdFeedRewardRecord(
client_event_id="source-fallback-unique", user_id=user.id,
reward_date=SOURCE_FALLBACK_DATE, duration_seconds=3, unit_count=0,
ecpm_raw="4700", ad_session_id="flow-session", trace_id="source-trace",
app_env="prod", our_code_id="104098712", coin=0, status="too_short",
ad_type="draw", feed_scene="comparison",
created_at=datetime(2040, 2, 7, 1, tzinfo=UTC),
),
AdFeedRewardRecord(
client_event_id="source-fallback-ambiguous", user_id=user.id,
reward_date=SOURCE_FALLBACK_DATE, duration_seconds=3, unit_count=0,
ecpm_raw="4800", ad_session_id="flow-session", trace_id="source-trace",
app_env="prod", our_code_id="104098712", coin=0, status="too_short",
ad_type="draw", feed_scene="comparison",
created_at=datetime(2040, 2, 7, 2, tzinfo=UTC),
),
AdEcpmRecord(
user_id=user.id, ad_type="draw", ad_session_id="impression-unique",
trace_id="source-trace", ecpm_raw="4700", adn="baidu", slot_id="rit-baidu",
app_env="prod", our_code_id="104098712", report_date=SOURCE_FALLBACK_DATE,
created_at=datetime(2040, 2, 7, 1, tzinfo=UTC),
),
AdEcpmRecord(
user_id=user.id, ad_type="draw", ad_session_id="impression-ambiguous-a",
trace_id="source-trace", ecpm_raw="4800", adn="baidu", slot_id="rit-baidu",
app_env="prod", our_code_id="104098712", report_date=SOURCE_FALLBACK_DATE,
created_at=datetime(2040, 2, 7, 2, tzinfo=UTC),
),
AdEcpmRecord(
user_id=user.id, ad_type="draw", ad_session_id="impression-ambiguous-b",
trace_id="source-trace", ecpm_raw="4800", adn="ks", slot_id="rit-ks",
app_env="prod", our_code_id="104098712", report_date=SOURCE_FALLBACK_DATE,
created_at=datetime(2040, 2, 7, 2, 1, tzinfo=UTC),
),
])
db.commit()
result = ad_revenue.ad_revenue_report(
db,
date_from=SOURCE_FALLBACK_DATE,
date_to=SOURCE_FALLBACK_DATE,
user_id=user.id,
app_env="prod",
revenue_scope="all",
)
item = next(row for row in result["items"] if row["event_key"].startswith("feedgrp-"))
details = {detail["ecpm"]: detail for detail in item["sub_rewards"]}
assert details["4700"]["adn"] == "baidu"
assert details["4700"]["slot_id"] == "rit-baidu"
assert details["4800"]["adn"] is None
assert details["4800"]["slot_id"] is None
finally:
db.rollback()
db.execute(delete(AdFeedRewardRecord).where(AdFeedRewardRecord.reward_date == SOURCE_FALLBACK_DATE))
db.execute(delete(AdEcpmRecord).where(AdEcpmRecord.report_date == SOURCE_FALLBACK_DATE))
db.execute(delete(User).where(User.phone == phone))
db.commit()
db.close()
def test_pangle_api_day_is_provisional_before_14_beijing_time() -> None:
assert ad_revenue._pangle_api_day_complete(
REPORT_DATE,
{
"api_revenue_yuan": 3.2,
# D+1 10:30 北京时间。
"synced_at": datetime(2040, 2, 4, 2, 30, tzinfo=UTC),
},
) is False
assert ad_revenue._pangle_api_day_complete(
REPORT_DATE,
{
"api_revenue_yuan": 3.2,
# D+1 14:00 北京时间,达到日终判定线。
"synced_at": datetime(2040, 2, 4, 6, 0, tzinfo=UTC),
},
) is True
+1 -15
View File
@@ -1,7 +1,7 @@
"""Admin M2 读接口测试:大盘聚合 + 用户/流水/提现/反馈列表 + 鉴权拦截。"""
from __future__ import annotations
from datetime import UTC, datetime
from datetime import datetime
import pytest
from fastapi.testclient import TestClient
@@ -9,7 +9,6 @@ from sqlalchemy import event
from app.admin.main import admin_app
from app.admin.repositories import admin_user as admin_repo
from app.admin.repositories import queries
from app.db.session import SessionLocal, engine
from app.models.comparison import ComparisonRecord
from app.models.feedback import Feedback
@@ -145,8 +144,6 @@ def test_user_reward_detail_does_not_select_unrelated_new_ad_columns(
) -> None:
if "ad_reward_record.boost_round_id" in statement:
raise AssertionError("提现详情不应查询未使用的 boost_round_id")
if "FROM user" in statement and "user.phone" in statement:
raise AssertionError("奖励统计的用户存在性检查不应展开完整 user 表")
event.listen(engine, "before_cursor_execute", reject_full_ad_reward_projection)
try:
@@ -165,17 +162,6 @@ def test_user_reward_detail_does_not_select_unrelated_new_ad_columns(
assert records.status_code == 200, records.text
def test_user_coin_record_sort_accepts_mixed_timezone_datetimes() -> None:
"""线上 PostgreSQL 返回 awareSQLite/历史转换可能返回 naive,二者必须可混排。"""
naive = datetime(2038, 1, 1, 8, 0)
aware = datetime(2038, 1, 1, 7, 0, tzinfo=UTC)
rows = [{"created_at": aware}, {"created_at": naive}]
rows.sort(key=queries._coin_record_sort_key, reverse=True)
assert rows == [{"created_at": naive}, {"created_at": aware}]
def test_user_filter_by_status(admin_client: TestClient, admin_token: str) -> None:
_seed_user_with_data("13800000003")
r = admin_client.get("/admin/api/users", params={"status": "active"}, headers=_auth(admin_token))