"""美团 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")