a2270ee1b2
新手引导视频:运营后台上传 MP4(上限 100MB,魔数校验只认 ISO BMFF), App 端在领券等候浮层前 N 次以引导视频替代广告。新增 guide_video 的 model/schema/repository/router(App 侧 + 后台侧)与播放记录表迁移。 美团券:首页「销量最高 / 智能推荐」两个 tab 改游标分页,配套两条 (city_id, dedup_key, 排序键 DESC) 复合索引,让 Postgres 顺着索引流式 去重,免掉每翻一页重排整城券的开销。美团 CPS client 在 lifespan 预热 并在关闭时释放连接池。 Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
414 lines
19 KiB
Python
414 lines
19 KiB
Python
"""美团 CPS 券列表 + 换链。
|
|
|
|
不需要登录,客户端传经纬度即可。
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
from concurrent.futures import ThreadPoolExecutor
|
|
from typing import TYPE_CHECKING, Any
|
|
|
|
from fastapi import APIRouter, Depends, HTTPException
|
|
from sqlalchemy import nullslast, select
|
|
from sqlalchemy.orm import Session
|
|
|
|
from app.core.config import settings
|
|
from app.db.session import get_db
|
|
from app.integrations.meituan import MeituanCpsError, _call, get_referral_link, query_coupon
|
|
from app.models.meituan_coupon import MeituanCoupon
|
|
from app.schemas.meituan import (
|
|
CouponCard,
|
|
CouponListRequest,
|
|
CouponListResponse,
|
|
FeedRequest,
|
|
FeedResponse,
|
|
ReferralLinkRequest,
|
|
ReferralLinkResponse,
|
|
TopSalesRequest,
|
|
)
|
|
from app.utils import mt_search_cursor
|
|
from app.utils.meituan_city import get_meituan_city
|
|
|
|
if TYPE_CHECKING: # 仅供类型标注(本模块已开 from __future__ import annotations)
|
|
from collections.abc import Callable
|
|
|
|
from sqlalchemy import ColumnElement
|
|
|
|
logger = logging.getLogger("shagua.meituan")
|
|
|
|
|
|
def _resolve_city_id(latitude: float, longitude: float) -> str:
|
|
"""经纬度 → 美团城市 ID;解析失败返 ""(调用方应降级返空)。"""
|
|
try:
|
|
return get_meituan_city(latitude, longitude).get("city_id", "")
|
|
except Exception:
|
|
logger.exception("get_meituan_city 失败")
|
|
return ""
|
|
|
|
router = APIRouter(prefix="/api/v1/meituan", tags=["meituan-cps"])
|
|
|
|
|
|
@router.post("/coupons", response_model=CouponListResponse, summary="券列表(为您推荐)")
|
|
def list_coupons(req: CouponListRequest) -> CouponListResponse:
|
|
if not settings.mt_cps_configured:
|
|
return CouponListResponse(items=[], has_next=False, search_id=None, status="degraded")
|
|
logger.info("[coupons] lon=%.6f lat=%.6f topic=%s", req.longitude, req.latitude, req.list_topic_id)
|
|
try:
|
|
raw = query_coupon(
|
|
longitude=req.longitude,
|
|
latitude=req.latitude,
|
|
platform=req.platform,
|
|
biz_line=req.biz_line,
|
|
list_topic_id=req.list_topic_id,
|
|
search_text=req.keyword,
|
|
search_id=req.search_id,
|
|
sort_field=req.sort_field,
|
|
page_no=req.page,
|
|
page_size=req.page_size,
|
|
)
|
|
except MeituanCpsError as e:
|
|
# 软降级:不再抛 502(避免前端弹报错 toast),返空 + degraded,前端显示「服务繁忙」。
|
|
logger.warning("[coupons] query_coupon 失败,降级返空: %s", e)
|
|
return CouponListResponse(items=[], has_next=False, search_id=None, status="degraded")
|
|
|
|
items: list[CouponCard] = []
|
|
for it in (raw.get("data") or []):
|
|
try:
|
|
items.append(CouponCard.from_raw(it))
|
|
except Exception: # noqa: BLE001
|
|
continue
|
|
return CouponListResponse(
|
|
items=items,
|
|
has_next=raw.get("hasNext", False),
|
|
search_id=raw.get("searchId"),
|
|
status="ok" if items else "empty",
|
|
)
|
|
|
|
|
|
_TOPIC_ROUNDS = [
|
|
(3, 3), # 爆款筛选
|
|
(2, 2), # 今日必推
|
|
(1, 5), # 精选 + 限时筛选
|
|
]
|
|
|
|
def _interleave(waimai: list[dict], daodian: list[dict]) -> list[CouponCard]:
|
|
items: list[CouponCard] = []
|
|
seen: set[str] = set()
|
|
i = j = 0
|
|
while i < len(waimai) or j < len(daodian):
|
|
for _ in range(2):
|
|
if i < len(waimai):
|
|
card = CouponCard.from_raw(waimai[i]); i += 1
|
|
if card.product_view_sign not in seen:
|
|
seen.add(card.product_view_sign); items.append(card)
|
|
if j < len(daodian):
|
|
card = CouponCard.from_raw(daodian[j]); j += 1
|
|
if card.product_view_sign not in seen:
|
|
seen.add(card.product_view_sign); items.append(card)
|
|
return items
|
|
|
|
|
|
def _commission_pct(card: CouponCard) -> float:
|
|
"""'1.4%' → 1.4;解析失败按 0(会被智能推荐过滤掉)。"""
|
|
try:
|
|
return float(card.commission_rate.rstrip("%"))
|
|
except (ValueError, AttributeError):
|
|
return 0.0
|
|
|
|
|
|
# ────────────── 离线库分页(智能推荐 / 销量最高 共用) ──────────────
|
|
# 去重+排序阶段**只投影这几列**:够 DISTINCT ON 分组、够排序、够回表定位,且全是定长小字段。
|
|
# ⚠️ 关键性能点:`raw` 是整条美团原始返回(JSONB,每行数 KB)。原实现用 select(MeituanCoupon)
|
|
# 做子查询,等于把整城几千行连 raw 一起塞进两次排序(DISTINCT ON 一次 + 分页一次),
|
|
# 体量轻松超过 work_mem → Postgres 落盘做外部归并排序,而且**每翻一页都要重来一遍**。
|
|
# 拆成「先在小列上排出本页 id,再按 id 回表取 raw」后,排序数据量降到原来的百分之几,
|
|
# JSONB 只解析当前页 ~20 行。
|
|
_DEDUP_COLS = (
|
|
MeituanCoupon.id,
|
|
MeituanCoupon.dedup_key,
|
|
MeituanCoupon.sale_volume_num,
|
|
MeituanCoupon.commission_percent,
|
|
)
|
|
|
|
|
|
def _paged_dedup_ids(
|
|
db: Session,
|
|
*,
|
|
conds: list[ColumnElement[bool]],
|
|
dedup_order: list[ColumnElement],
|
|
page_order: Callable[[Any], list[ColumnElement]],
|
|
page: int,
|
|
page_size: int,
|
|
) -> tuple[list[int], bool]:
|
|
"""DISTINCT ON(dedup_key) 跨源去重 → 整体排序 → 分页,返回 (本页 id 列表, 是否还有下一页)。
|
|
|
|
- `dedup_order`:同一个 dedup_key 的多条里留哪条(如销量最高/佣金最高)。
|
|
- `page_order`:接收去重子查询的列集合(`sub.c`),返回去重后的整体排序。
|
|
多取 1 条用于判断 has_next。
|
|
"""
|
|
deduped = (
|
|
select(*_DEDUP_COLS)
|
|
.where(*conds)
|
|
.distinct(MeituanCoupon.dedup_key)
|
|
.order_by(MeituanCoupon.dedup_key, *dedup_order)
|
|
.subquery()
|
|
)
|
|
ids = db.execute(
|
|
select(deduped.c.id)
|
|
.order_by(*page_order(deduped.c))
|
|
.offset((page - 1) * page_size)
|
|
.limit(page_size + 1)
|
|
).scalars().all()
|
|
return list(ids[:page_size]), len(ids) > page_size
|
|
|
|
|
|
def _load_raws(db: Session, ids: list[int]) -> list[dict]:
|
|
"""按给定 id 顺序取 raw(只回表本页 ~20 行)。缺行(被 ETL 清掉)静默跳过。"""
|
|
if not ids:
|
|
return []
|
|
raw_by_id = {
|
|
row_id: raw
|
|
for row_id, raw in db.execute(
|
|
select(MeituanCoupon.id, MeituanCoupon.raw).where(MeituanCoupon.id.in_(ids))
|
|
).all()
|
|
}
|
|
return [raw_by_id[i] for i in ids if i in raw_by_id]
|
|
|
|
|
|
def _cards_from_raws(raws: list[dict], *, hide_distance: bool) -> list[CouponCard]:
|
|
"""raw → CouponCard;解析失败的单条跳过,不整页失败。
|
|
|
|
hide_distance:离线库里的距离是相对「城市默认点」算的,对用户无意义且误导 —— 智能推荐 /
|
|
销量最高两个 tab 一律置空,前端「距离 店名」那行只剩店名、自动顶到最左。
|
|
"""
|
|
cards: list[CouponCard] = []
|
|
for raw in raws:
|
|
try:
|
|
card = CouponCard.from_raw(raw or {})
|
|
except Exception: # noqa: BLE001
|
|
continue
|
|
if not card.product_view_sign:
|
|
continue
|
|
if hide_distance:
|
|
card.distance_text = None
|
|
card.distance_meters = None
|
|
cards.append(card)
|
|
return cards
|
|
|
|
|
|
@router.post("/feed", response_model=FeedResponse, summary="混合feed(外卖+到店交叉);tab=rec智能推荐/distance距离最近")
|
|
def feed(req: FeedRequest, db: Session = Depends(get_db)) -> FeedResponse:
|
|
lon, lat = req.longitude, req.latitude
|
|
tab = (req.tab or "").strip()
|
|
logger.info("[feed] tab=%s page=%s lon=%.6f lat=%.6f", tab or "(default)", req.page, lon, lat)
|
|
|
|
# rec 是纯离线库查询,不依赖 MT 凭证、不实时打美团;其余 tab(distance / 默认)要实时打美团,
|
|
# 未配凭证直接降级返空(degraded),让前端区分「服务暂不可用」而非「暂无」。
|
|
if tab != "rec" and not settings.mt_cps_configured:
|
|
return FeedResponse(items=[], has_next=False, page=req.page, status="degraded")
|
|
|
|
def _fetch_topic(platform: int, biz_line: int | None, topic: int) -> tuple[list[dict], bool]:
|
|
"""返回(items, 是否调用失败)。失败标志用于把默认 feed 整页标 degraded。"""
|
|
try:
|
|
data = (query_coupon(
|
|
longitude=lon, latitude=lat,
|
|
platform=platform, biz_line=biz_line,
|
|
list_topic_id=topic, page_size=20,
|
|
).get("data") or [])
|
|
return data, False
|
|
except MeituanCpsError:
|
|
return [], True
|
|
|
|
# 距离最近:搜索召回(外卖搜"外卖" + 到店搜"美食",都 sortField=6 离我最近)一页页拉。
|
|
# 搜索翻页必须用 searchId(pageNo 翻不动),而接口是无状态的(客户端只传页码)—— 原实现因此
|
|
# 每次都从第 1 页顺序重放到第 N 页,取第 N 页要向美团发 N 次请求,越往下滑越慢。
|
|
# 现在把沿途 searchId 记进 [mt_search_cursor],稳态下每翻一页恒定 1 次请求;两路仍并行。
|
|
# 按你位置实时算距离(库里没存 POI 经纬度,只能实时)。
|
|
if tab == "distance":
|
|
lon_i, lat_i = int(lon * 1_000_000), int(lat * 1_000_000)
|
|
|
|
def _replay(
|
|
platform: int, biz_line: int | None, keyword: str,
|
|
key: mt_search_cursor.RouteKey, start: int, sid: str | None, n: int,
|
|
) -> tuple[list[dict], bool, bool]:
|
|
"""从第 start 页(用 sid 取)顺序翻到第 n 页。start==1 时 sid 应为 None(走 pageNo=1)。"""
|
|
data: list[dict] = []
|
|
has_next = False
|
|
for pg in range(start, n + 1):
|
|
body: dict = {
|
|
"platform": platform, "searchText": keyword, "sortField": 6,
|
|
"longitude": lon_i, "latitude": lat_i, "pageSize": 20,
|
|
}
|
|
if biz_line is not None:
|
|
body["bizLine"] = biz_line
|
|
if sid:
|
|
body["searchId"] = sid
|
|
else:
|
|
body["pageNo"] = 1
|
|
try:
|
|
r = _call("/cps_open/common/api/v1/query_coupon", body)
|
|
except MeituanCpsError:
|
|
return [], False, True # 调用失败(用于标 degraded)
|
|
data = r.get("data") or []
|
|
sid = r.get("searchId")
|
|
has_next = bool(r.get("hasNext")) and bool(data)
|
|
# 记下「下一页要用哪个 searchId」;没有下一页就别记,免得存进死游标。
|
|
if sid and has_next:
|
|
mt_search_cursor.remember(key, pg + 1, sid)
|
|
if not data or (not has_next and pg < n):
|
|
return [], False, False # 没那么多页了(非错误)
|
|
return data, has_next, False
|
|
|
|
def _search_page_n(platform: int, biz_line: int | None, keyword: str, n: int) -> tuple[list[dict], bool, bool]:
|
|
"""取第 n 页,返回(第 n 页 items, 是否还有下一页, 是否调用失败)。
|
|
|
|
优先用缓存游标一发直达;缓存未命中/过期才从最近的已知页往后重放,并把沿途游标补进缓存。
|
|
"""
|
|
key = mt_search_cursor.route_key(lat, lon, platform, keyword)
|
|
start, sid = mt_search_cursor.lookup(key, n)
|
|
data, has_next, failed = _replay(platform, biz_line, keyword, key, start, sid, n)
|
|
# 用缓存游标却打不通,多半是上游 searchId 过期:作废整条路线,回到第 1 页重放一次。
|
|
if failed and start > 1:
|
|
mt_search_cursor.drop(key)
|
|
data, has_next, failed = _replay(platform, biz_line, keyword, key, 1, None, n)
|
|
return data, has_next, failed
|
|
|
|
with ThreadPoolExecutor(max_workers=2) as pool:
|
|
f_wm = pool.submit(_search_page_n, 1, None, "外卖", req.page)
|
|
f_dd = pool.submit(_search_page_n, 2, 1, "美食", req.page)
|
|
(wm_data, wm_hn, wm_fail), (dd_data, dd_hn, dd_fail) = f_wm.result(), f_dd.result()
|
|
|
|
seen: set[str] = set()
|
|
cards: list[CouponCard] = []
|
|
for it in wm_data + dd_data:
|
|
try:
|
|
card = CouponCard.from_raw(it)
|
|
except Exception: # noqa: BLE001
|
|
continue
|
|
if card.product_view_sign and card.product_view_sign not in seen:
|
|
seen.add(card.product_view_sign)
|
|
cards.append(card)
|
|
cards.sort(key=lambda c: c.distance_meters if c.distance_meters is not None else float("inf"))
|
|
# 两路都失败且无数据 → degraded;否则有数据 ok / 无数据 empty
|
|
status = "degraded" if (not cards and wm_fail and dd_fail) else ("ok" if cards else "empty")
|
|
return FeedResponse(items=cards, has_next=wm_hn or dd_hn, page=req.page, status=status)
|
|
|
|
# 智能推荐(rec):走【离线库】筛佣金率 ≥ 3%,按城市过滤,分页返回(SQL 侧去重+排序+分页,秒级、不打美团)。
|
|
# 实测库里佣金≥3% 去重后仅 ~578 条(几乎全是外卖;到店团购佣金普遍 <3%):实时按"同城热销榜单"
|
|
# 拉既撞限流、又填不满(该榜单中位佣金 ~0.8%,筛完每页剩 0-1 条),故从库出。佣金阈值逻辑不变。
|
|
if tab == "rec":
|
|
city_id = _resolve_city_id(lat, lon)
|
|
if not city_id:
|
|
return FeedResponse(items=[], has_next=False, page=req.page, status="degraded")
|
|
PAGE = 20
|
|
try:
|
|
ids, has_next = _paged_dedup_ids(
|
|
db,
|
|
conds=[
|
|
MeituanCoupon.commission_percent >= 3.0,
|
|
MeituanCoupon.city_id == city_id,
|
|
],
|
|
# 同一去重键留佣金最高那条
|
|
dedup_order=[MeituanCoupon.commission_percent.desc()],
|
|
# 销量高的优先(无销量档排后),同档佣金高优先,id 兜底稳定分页
|
|
page_order=lambda c: [
|
|
nullslast(c.sale_volume_num.desc()), c.commission_percent.desc(), c.id,
|
|
],
|
|
page=req.page, page_size=PAGE,
|
|
)
|
|
raws = _load_raws(db, ids)
|
|
except Exception: # noqa: BLE001
|
|
logger.exception("[feed] rec 库查询失败,降级返空")
|
|
return FeedResponse(items=[], has_next=False, page=req.page, status="degraded")
|
|
cards = _cards_from_raws(raws, hide_distance=True)
|
|
if not cards and req.page == 1:
|
|
# 命中城市却 0 券:该城确无 ≥3% 券,或 ETL 灌的 city_id 与 city_dict 口径不一致。
|
|
logger.info("[feed] rec city_id=%s 命中 0 券(该城确无券?或 ETL/city_dict 的 city_id 口径不一致)", city_id)
|
|
return FeedResponse(items=cards, has_next=has_next, page=req.page,
|
|
status="ok" if cards else "empty")
|
|
|
|
# 默认(老客户端不传 tab):沿用逐轮分页的混合 feed,不筛。
|
|
page_idx = req.page - 1
|
|
if page_idx >= len(_TOPIC_ROUNDS):
|
|
return FeedResponse(items=[], has_next=False, page=req.page, status="empty")
|
|
|
|
wm_topic, dd_topic = _TOPIC_ROUNDS[page_idx]
|
|
with ThreadPoolExecutor(max_workers=2) as pool:
|
|
f_wm = pool.submit(_fetch_topic, 1, None, wm_topic)
|
|
f_dd = pool.submit(_fetch_topic, 2, 1, dd_topic)
|
|
(waimai, wm_fail), (daodian, dd_fail) = f_wm.result(), f_dd.result()
|
|
|
|
items = _interleave(waimai, daodian)
|
|
has_next = page_idx + 1 < len(_TOPIC_ROUNDS)
|
|
# 两路都失败且无数据 → degraded;否则有数据 ok / 无数据 empty
|
|
status = "degraded" if (not items and wm_fail and dd_fail) else ("ok" if items else "empty")
|
|
return FeedResponse(items=items, has_next=has_next, page=req.page, status=status)
|
|
|
|
|
|
@router.post("/referral-link", response_model=ReferralLinkResponse, summary="换取推广链接(点抢时调)")
|
|
def referral_link(req: ReferralLinkRequest) -> ReferralLinkResponse:
|
|
if not settings.mt_cps_configured:
|
|
return ReferralLinkResponse(link="", link_map={})
|
|
try:
|
|
raw = get_referral_link(
|
|
product_view_sign=req.product_view_sign,
|
|
platform=req.platform,
|
|
biz_line=req.biz_line,
|
|
sid=req.sid,
|
|
link_type_list=req.link_type_list,
|
|
)
|
|
except MeituanCpsError as e:
|
|
logger.error("get_referral_link failed: %s", e)
|
|
raise HTTPException(status_code=502, detail=f"meituan: {e}") from e
|
|
|
|
link_map = raw.get("referralLinkMap") or {}
|
|
link = raw.get("data") or link_map.get("1") or link_map.get("3") or next(iter(link_map.values()), "")
|
|
return ReferralLinkResponse(link=link, link_map=link_map)
|
|
|
|
|
|
@router.post("/top-sales", response_model=CouponListResponse,
|
|
summary="销量最高(从离线库 meituan_coupon 按销量降序 + 跨源去重,按城市过滤,不实时打美团)")
|
|
def top_sales(req: TopSalesRequest, db: Session = Depends(get_db)) -> CouponListResponse:
|
|
# 按设备经纬度定位城市,只查同城券;老客户端不带坐标 → 降级返空(不 422、不误返全城)。
|
|
if req.latitude is None or req.longitude is None:
|
|
return CouponListResponse(items=[], has_next=False, search_id=None, status="degraded")
|
|
city_id = _resolve_city_id(req.latitude, req.longitude)
|
|
if not city_id:
|
|
return CouponListResponse(items=[], has_next=False, search_id=None, status="degraded")
|
|
|
|
# 去重 + 排序 + 分页全在 SQL 做,每页只回表并解析当前页 ~20 条(见 _paged_dedup_ids 的性能说明)。
|
|
# 库为空(prod 刚部署 / ETL 未跑完)时返空 + status=empty,不崩;库查询异常降级 degraded。
|
|
conds = [
|
|
MeituanCoupon.sale_volume_num.isnot(None),
|
|
MeituanCoupon.city_id == city_id,
|
|
]
|
|
if req.platform is not None:
|
|
conds.append(MeituanCoupon.platform == req.platform)
|
|
try:
|
|
ids, has_next = _paged_dedup_ids(
|
|
db,
|
|
conds=conds,
|
|
# 每个去重键(品牌|名|价)只留销量最高那条(同销量再按佣金)
|
|
dedup_order=[
|
|
MeituanCoupon.sale_volume_num.desc(),
|
|
MeituanCoupon.commission_percent.desc(),
|
|
],
|
|
# 加 id 作稳定 tiebreaker:同销量同佣金的并列项排序确定,避免跨页重复/漏项
|
|
page_order=lambda c: [
|
|
c.sale_volume_num.desc(), c.commission_percent.desc(), c.id,
|
|
],
|
|
page=req.page, page_size=req.page_size,
|
|
)
|
|
raws = _load_raws(db, ids)
|
|
except Exception: # noqa: BLE001
|
|
logger.exception("[top-sales] 库查询失败,降级返空")
|
|
return CouponListResponse(items=[], has_next=False, search_id=None, status="degraded")
|
|
|
|
# 不显示距离:库里的距离是相对城市默认点的(对用户无意义、且误导),与推荐流口径一致。
|
|
cards = _cards_from_raws(raws, hide_distance=True)
|
|
if not cards and req.page == 1:
|
|
# 命中城市却 0 券:可能该城确无券,也可能 ETL 灌的 city_id 与 city_dict 口径不一致(静默降级的隐患)。
|
|
logger.info("[top-sales] city_id=%s 命中 0 券(该城确无券?或 ETL/city_dict 的 city_id 口径不一致)", city_id)
|
|
return CouponListResponse(items=cards, has_next=has_next, search_id=None,
|
|
status="ok" if cards else "empty")
|