Compare commits
5 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 179303d875 | |||
| 356d6e52b9 | |||
| 4faf0b8e16 | |||
| c9578f5ba2 | |||
| cfeacb4bab |
@@ -40,9 +40,6 @@ MT_CPS_APP_KEY=
|
||||
MT_CPS_APP_SECRET=
|
||||
# 默认渠道追踪标识(sid),用于区分不同 app 的 CPS 数据
|
||||
MT_CPS_DEFAULT_SID=sgbjia
|
||||
# 美团调用走的代理。⚠️ 本机/内网开发直连美团会 SSL EOF,必须填本地代理(如 http://127.0.0.1:7897);
|
||||
# 线上国内服务器留空(=直连)。留空且本机直连失败时 /feed、/coupons、/top-sales 会返回空。
|
||||
MT_CPS_PROXY=
|
||||
|
||||
# ===== Pricebot 上游 (领券/比价业务透传目标) =====
|
||||
# 客户端调本服务的 /api/v1/coupon/step 等,我们透传到 pricebot-backend。
|
||||
|
||||
@@ -1,26 +0,0 @@
|
||||
"""merge meituan_coupon 与 ops/price_observation 迁移头
|
||||
|
||||
Revision ID: 11a1d08c6f55
|
||||
Revises: meituan_coupon_table, opsrename01
|
||||
Create Date: 2026-06-08 00:06:01.633796
|
||||
|
||||
"""
|
||||
from typing import Sequence, Union
|
||||
|
||||
from alembic import op
|
||||
import sqlalchemy as sa
|
||||
|
||||
|
||||
# revision identifiers, used by Alembic.
|
||||
revision: str = '11a1d08c6f55'
|
||||
down_revision: Union[str, Sequence[str], None] = ('meituan_coupon_table', 'opsrename01')
|
||||
branch_labels: Union[str, Sequence[str], None] = None
|
||||
depends_on: Union[str, Sequence[str], None] = None
|
||||
|
||||
|
||||
def upgrade() -> None:
|
||||
pass
|
||||
|
||||
|
||||
def downgrade() -> None:
|
||||
pass
|
||||
@@ -1,77 +0,0 @@
|
||||
"""coupon state tables (领券今日状态:弹窗频控 engagement + 领券记录 claim)
|
||||
|
||||
Revision ID: coupon_state_tables
|
||||
Revises: 11a1d08c6f55
|
||||
Create Date: 2026-06-08 00:00:00.000000
|
||||
|
||||
"""
|
||||
from typing import Sequence, Union
|
||||
|
||||
from alembic import op
|
||||
import sqlalchemy as sa
|
||||
from sqlalchemy.dialects import postgresql
|
||||
|
||||
|
||||
# revision identifiers, used by Alembic.
|
||||
revision: str = 'coupon_state_tables'
|
||||
down_revision: Union[str, Sequence[str], None] = '11a1d08c6f55'
|
||||
branch_labels: Union[str, Sequence[str], None] = None
|
||||
depends_on: Union[str, Sequence[str], None] = None
|
||||
|
||||
# PG 上为 JSONB,其它(SQLite dev)为 JSON——与模型层 with_variant 对齐
|
||||
_JSON = sa.JSON().with_variant(postgresql.JSONB(), 'postgresql')
|
||||
|
||||
|
||||
def upgrade() -> None:
|
||||
# 领券记录(资产层,当前不参与去重判断)
|
||||
op.create_table(
|
||||
'coupon_claim_record',
|
||||
sa.Column('id', sa.Integer(), autoincrement=True, nullable=False),
|
||||
sa.Column('device_id', sa.String(length=64), nullable=False),
|
||||
sa.Column('user_id', sa.Integer(), nullable=True),
|
||||
sa.Column('coupon_id', sa.String(length=64), nullable=False),
|
||||
sa.Column('claim_date', sa.Date(), nullable=False),
|
||||
sa.Column('status', sa.String(length=24), nullable=False),
|
||||
sa.Column('vendor', sa.String(length=48), nullable=True),
|
||||
sa.Column('coupon_name', sa.String(length=128), nullable=True),
|
||||
sa.Column('claimed_count', sa.Integer(), nullable=True),
|
||||
sa.Column('trace_id', sa.String(length=64), nullable=True),
|
||||
sa.Column('reason', sa.String(length=255), nullable=True),
|
||||
sa.Column('extra', _JSON, nullable=True),
|
||||
sa.Column('created_at', sa.DateTime(timezone=True), server_default=sa.text('(CURRENT_TIMESTAMP)'), nullable=False),
|
||||
sa.Column('updated_at', sa.DateTime(timezone=True), server_default=sa.text('(CURRENT_TIMESTAMP)'), nullable=False),
|
||||
sa.PrimaryKeyConstraint('id'),
|
||||
sa.UniqueConstraint('device_id', 'coupon_id', 'claim_date', name='uq_coupon_claim_device_coupon_date'),
|
||||
)
|
||||
with op.batch_alter_table('coupon_claim_record', schema=None) as batch_op:
|
||||
batch_op.create_index('ix_coupon_claim_device_date', ['device_id', 'claim_date'], unique=False)
|
||||
batch_op.create_index(batch_op.f('ix_coupon_claim_record_user_id'), ['user_id'], unique=False)
|
||||
batch_op.create_index(batch_op.f('ix_coupon_claim_record_trace_id'), ['trace_id'], unique=False)
|
||||
|
||||
# 弹窗频控(今天 engage 过就不再弹)
|
||||
op.create_table(
|
||||
'coupon_prompt_engagement',
|
||||
sa.Column('id', sa.Integer(), autoincrement=True, nullable=False),
|
||||
sa.Column('device_id', sa.String(length=64), nullable=False),
|
||||
sa.Column('user_id', sa.Integer(), nullable=True),
|
||||
sa.Column('engage_date', sa.Date(), nullable=False),
|
||||
sa.Column('engage_type', sa.String(length=16), nullable=False),
|
||||
sa.Column('created_at', sa.DateTime(timezone=True), server_default=sa.text('(CURRENT_TIMESTAMP)'), nullable=False),
|
||||
sa.Column('updated_at', sa.DateTime(timezone=True), server_default=sa.text('(CURRENT_TIMESTAMP)'), nullable=False),
|
||||
sa.PrimaryKeyConstraint('id'),
|
||||
sa.UniqueConstraint('device_id', 'engage_date', name='uq_coupon_engage_device_date'),
|
||||
)
|
||||
with op.batch_alter_table('coupon_prompt_engagement', schema=None) as batch_op:
|
||||
batch_op.create_index(batch_op.f('ix_coupon_prompt_engagement_user_id'), ['user_id'], unique=False)
|
||||
|
||||
|
||||
def downgrade() -> None:
|
||||
with op.batch_alter_table('coupon_prompt_engagement', schema=None) as batch_op:
|
||||
batch_op.drop_index(batch_op.f('ix_coupon_prompt_engagement_user_id'))
|
||||
op.drop_table('coupon_prompt_engagement')
|
||||
|
||||
with op.batch_alter_table('coupon_claim_record', schema=None) as batch_op:
|
||||
batch_op.drop_index(batch_op.f('ix_coupon_claim_record_trace_id'))
|
||||
batch_op.drop_index(batch_op.f('ix_coupon_claim_record_user_id'))
|
||||
batch_op.drop_index('ix_coupon_claim_device_date')
|
||||
op.drop_table('coupon_claim_record')
|
||||
@@ -1,62 +0,0 @@
|
||||
"""meituan_coupon table(美团 CPS 券本地缓存,供销量/佣金排序从库里捞、本地排序)
|
||||
|
||||
Revision ID: meituan_coupon_table
|
||||
Revises: withdraw_review_ad_watch
|
||||
Create Date: 2026-06-06 12:00:00.000000
|
||||
|
||||
"""
|
||||
from typing import Sequence, Union
|
||||
|
||||
from alembic import op
|
||||
import sqlalchemy as sa
|
||||
from sqlalchemy.dialects import postgresql
|
||||
|
||||
|
||||
# revision identifiers, used by Alembic.
|
||||
revision: str = 'meituan_coupon_table'
|
||||
down_revision: Union[str, Sequence[str], None] = 'withdraw_review_ad_watch'
|
||||
branch_labels: Union[str, Sequence[str], None] = None
|
||||
depends_on: Union[str, Sequence[str], None] = None
|
||||
|
||||
|
||||
def upgrade() -> None:
|
||||
op.create_table(
|
||||
'meituan_coupon',
|
||||
sa.Column('id', sa.Integer(), autoincrement=True, nullable=False),
|
||||
sa.Column('source', sa.String(length=16), nullable=False),
|
||||
sa.Column('platform', sa.Integer(), nullable=False),
|
||||
sa.Column('biz_line', sa.Integer(), nullable=True),
|
||||
sa.Column('city_id', sa.String(length=32), nullable=False),
|
||||
sa.Column('product_view_sign', sa.String(length=128), nullable=False),
|
||||
sa.Column('sku_view_id', sa.String(length=128), nullable=True),
|
||||
sa.Column('name', sa.String(length=256), nullable=True),
|
||||
sa.Column('brand_name', sa.String(length=128), nullable=True),
|
||||
sa.Column('sell_price_cents', sa.Integer(), nullable=True),
|
||||
sa.Column('original_price_cents', sa.Integer(), nullable=True),
|
||||
sa.Column('head_url', sa.String(length=512), nullable=True),
|
||||
sa.Column('sale_volume', sa.String(length=32), nullable=True),
|
||||
sa.Column('sale_volume_num', sa.Integer(), nullable=True),
|
||||
sa.Column('commission_percent', sa.Float(), nullable=True),
|
||||
sa.Column('commission_amount_cents', sa.Integer(), nullable=True),
|
||||
sa.Column('poi_name', sa.String(length=128), nullable=True),
|
||||
sa.Column('available_poi_num', sa.Integer(), nullable=True),
|
||||
sa.Column('delivery_distance_m', sa.Float(), nullable=True),
|
||||
sa.Column('dedup_key', sa.String(length=64), nullable=False),
|
||||
sa.Column('raw', sa.JSON().with_variant(postgresql.JSONB(), 'postgresql'), nullable=False),
|
||||
sa.Column('first_seen', sa.DateTime(timezone=True), server_default=sa.text('(CURRENT_TIMESTAMP)'), nullable=False),
|
||||
sa.Column('last_seen', sa.DateTime(timezone=True), server_default=sa.text('(CURRENT_TIMESTAMP)'), nullable=False),
|
||||
sa.Column('updated_at', sa.DateTime(timezone=True), server_default=sa.text('(CURRENT_TIMESTAMP)'), nullable=False),
|
||||
sa.PrimaryKeyConstraint('id'),
|
||||
sa.UniqueConstraint('source', 'product_view_sign', name='uq_meituan_coupon_source_sign'),
|
||||
)
|
||||
op.create_index('ix_meituan_coupon_source', 'meituan_coupon', ['source'])
|
||||
op.create_index('ix_meituan_coupon_city_id', 'meituan_coupon', ['city_id'])
|
||||
op.create_index('ix_meituan_coupon_brand_name', 'meituan_coupon', ['brand_name'])
|
||||
op.create_index('ix_meituan_coupon_sale_volume_num', 'meituan_coupon', ['sale_volume_num'])
|
||||
op.create_index('ix_meituan_coupon_commission_percent', 'meituan_coupon', ['commission_percent'])
|
||||
op.create_index('ix_meituan_coupon_dedup_key', 'meituan_coupon', ['dedup_key'])
|
||||
op.create_index('ix_meituan_coupon_last_seen', 'meituan_coupon', ['last_seen'])
|
||||
|
||||
|
||||
def downgrade() -> None:
|
||||
op.drop_table('meituan_coupon')
|
||||
+1
-107
@@ -16,66 +16,15 @@ from typing import Any
|
||||
|
||||
import httpx
|
||||
from fastapi import APIRouter, HTTPException, Request, status
|
||||
from fastapi.concurrency import run_in_threadpool
|
||||
|
||||
from app.api.deps import DbSession
|
||||
from app.core.config import settings
|
||||
from app.core.pricebot_router import pick_pricebot
|
||||
from app.db.session import SessionLocal
|
||||
from app.repositories import coupon_state as coupon_repo
|
||||
from app.schemas.coupon_state import CouponPromptDismissIn, CouponPromptShouldShowOut
|
||||
|
||||
logger = logging.getLogger("shagua.coupon")
|
||||
|
||||
router = APIRouter(prefix="/api/v1/coupon", tags=["coupon"])
|
||||
|
||||
|
||||
def _to_int(v: object) -> int | None:
|
||||
"""user_id 协议是字符串(登录态才带),转 int 存库;缺失/非法 → None。"""
|
||||
if v is None:
|
||||
return None
|
||||
try:
|
||||
return int(v) # type: ignore[arg-type]
|
||||
except (TypeError, ValueError):
|
||||
return None
|
||||
|
||||
|
||||
def _extract_coupon_results(resp_json: dict) -> list[dict]:
|
||||
"""从 pricebot 响应抽本帧券结果,并**按 coupon_id 去重**。
|
||||
|
||||
⚠️ done/单券帧同时带 last_coupon_result(最后一张)+ action.params.coupon_results
|
||||
(全量含最后一张)→ 最后一张出现两次。必须去重:同批里同 coupon_id 两次,
|
||||
record_claims 在 autoflush=False 下两次 select 都查不到刚 add 的行 → 重复 add →
|
||||
commit 撞唯一约束 IntegrityError 回滚整批(done/单券记录全丢)。
|
||||
coupon_results(全量权威)覆盖 last_coupon_result。"""
|
||||
by_id: dict[str, dict] = {}
|
||||
last = resp_json.get("last_coupon_result")
|
||||
if isinstance(last, dict) and last.get("coupon_id"):
|
||||
by_id[last["coupon_id"]] = last
|
||||
params = (resp_json.get("action") or {}).get("params") or {}
|
||||
cr = params.get("coupon_results")
|
||||
if isinstance(cr, list):
|
||||
for x in cr:
|
||||
if isinstance(x, dict) and x.get("coupon_id"):
|
||||
by_id[x["coupon_id"]] = x
|
||||
return list(by_id.values())
|
||||
|
||||
|
||||
def _mark_engagement_blocking(
|
||||
device_id: str, user_id: int | None, engage_type: str
|
||||
) -> None:
|
||||
"""独立 session 写 engagement(async 端点经 run_in_threadpool 调,不阻塞事件循环)。"""
|
||||
with SessionLocal() as db:
|
||||
coupon_repo.mark_engagement(db, device_id, user_id, engage_type)
|
||||
|
||||
|
||||
def _record_claims_blocking(
|
||||
device_id: str, user_id: int | None, trace_id: str | None, results: list[dict]
|
||||
) -> None:
|
||||
with SessionLocal() as db:
|
||||
coupon_repo.record_claims(db, device_id, user_id, trace_id, results)
|
||||
|
||||
|
||||
@router.post("/step", summary="领券任务步进 (透传到 pricebot)")
|
||||
async def coupon_step(
|
||||
request: Request,
|
||||
@@ -96,21 +45,6 @@ async def coupon_step(
|
||||
if not isinstance(meta, dict):
|
||||
meta = {}
|
||||
|
||||
device_id = meta.get("device_id")
|
||||
user_id = _to_int(meta.get("user_id")) # 登录态才带;判断不靠它,资产留痕用
|
||||
trace_id = meta.get("trace_id")
|
||||
|
||||
# 领券任务首帧(step=0)= 用户已发起领券 → 记一条今日 engagement(claim_started),
|
||||
# 今天这台设备不再弹引导窗(对齐前台"点一键领取即 markEngaged")。写库失败绝不能
|
||||
# 连累领券主流程,整段吞掉。
|
||||
if device_id and meta.get("step") == 0:
|
||||
try:
|
||||
await run_in_threadpool(
|
||||
_mark_engagement_blocking, device_id, user_id, "claim_started"
|
||||
)
|
||||
except Exception as e: # noqa: BLE001
|
||||
logger.warning("coupon engagement write failed: %s", e)
|
||||
|
||||
# 按 trace_id 一致性 hash 选 pricebot 实例(同一领券任务的所有帧落同一进程)
|
||||
base = pick_pricebot(meta.get("trace_id"))
|
||||
url = f"{base.rstrip('/')}/api/coupon/step"
|
||||
@@ -146,44 +80,4 @@ async def coupon_step(
|
||||
detail=f"pricebot upstream returned {resp.status_code}",
|
||||
)
|
||||
|
||||
resp_json = resp.json()
|
||||
|
||||
# 领券结果沉淀(资产):中间帧 last_coupon_result + done 帧 coupon_results 幂等写库。
|
||||
# 当前只记录、不参与"要不要领"判断(MVP 先不去重)。写库失败不影响返回。
|
||||
if device_id:
|
||||
results = _extract_coupon_results(resp_json)
|
||||
if results:
|
||||
try:
|
||||
await run_in_threadpool(
|
||||
_record_claims_blocking, device_id, user_id, trace_id, results
|
||||
)
|
||||
except Exception as e: # noqa: BLE001
|
||||
logger.warning("coupon claim write failed: %s", e)
|
||||
|
||||
return resp_json
|
||||
|
||||
|
||||
@router.post("/prompt/dismiss", summary="用户拒绝/关闭领券引导窗(记今日已 engage)")
|
||||
def coupon_prompt_dismiss(payload: CouponPromptDismissIn, db: DbSession) -> dict[str, bool]:
|
||||
"""客户端点关闭引导窗时调用 → 记一条今日 engagement(dismissed),今天不再弹。
|
||||
|
||||
server 在透传链路里看不到"用户拒绝"(拒绝不发起领券),故必须客户端通知。
|
||||
MVP 不鉴权,按 device_id 记。
|
||||
"""
|
||||
coupon_repo.mark_engagement(db, payload.device_id, payload.user_id, "dismissed")
|
||||
return {"ok": True}
|
||||
|
||||
|
||||
@router.get(
|
||||
"/prompt/should-show",
|
||||
response_model=CouponPromptShouldShowOut,
|
||||
summary="切到外卖 App 时是否还应弹领券引导窗",
|
||||
)
|
||||
def coupon_prompt_should_show(
|
||||
device_id: str, db: DbSession
|
||||
) -> CouponPromptShouldShowOut:
|
||||
"""今天这台设备已 engage(领或拒)过 → should_show=false。客户端据此决定弹不弹
|
||||
(前台 SP 缓存做快速路径,这里是权威)。"""
|
||||
return CouponPromptShouldShowOut(
|
||||
should_show=not coupon_repo.has_engaged_today(db, device_id)
|
||||
)
|
||||
return resp.json()
|
||||
|
||||
+4
-75
@@ -7,14 +7,10 @@ from __future__ import annotations
|
||||
import logging
|
||||
from concurrent.futures import ThreadPoolExecutor, as_completed
|
||||
|
||||
from fastapi import APIRouter, Depends, HTTPException
|
||||
from sqlalchemy import select
|
||||
from sqlalchemy.orm import Session
|
||||
from fastapi import APIRouter, HTTPException
|
||||
|
||||
from app.core.config import settings
|
||||
from app.db.session import get_db
|
||||
from app.integrations.meituan import MeituanCpsError, get_referral_link, query_coupon
|
||||
from app.models.meituan_coupon import MeituanCoupon
|
||||
from app.schemas.meituan import (
|
||||
CouponCard,
|
||||
CouponListRequest,
|
||||
@@ -23,7 +19,6 @@ from app.schemas.meituan import (
|
||||
FeedResponse,
|
||||
ReferralLinkRequest,
|
||||
ReferralLinkResponse,
|
||||
TopSalesRequest,
|
||||
)
|
||||
|
||||
logger = logging.getLogger("shagua.meituan")
|
||||
@@ -84,21 +79,13 @@ def _interleave(waimai: list[dict], daodian: list[dict]) -> list[CouponCard]:
|
||||
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
|
||||
|
||||
|
||||
@router.post("/feed", response_model=FeedResponse, summary="混合feed(外卖+到店交叉);tab=rec智能推荐/distance距离最近")
|
||||
@router.post("/feed", response_model=FeedResponse, summary="混合feed(外卖+到店交叉, 无限流)")
|
||||
def feed(req: FeedRequest) -> FeedResponse:
|
||||
if not settings.mt_cps_configured:
|
||||
return FeedResponse(items=[], has_next=False, page=req.page)
|
||||
page_idx = req.page - 1
|
||||
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)
|
||||
logger.info("[feed] page=%s lon=%.6f lat=%.6f", req.page, lon, lat)
|
||||
|
||||
def _fetch_topic(platform: int, biz_line: int | None, topic: int) -> list[dict]:
|
||||
try:
|
||||
@@ -110,28 +97,6 @@ def feed(req: FeedRequest) -> FeedResponse:
|
||||
except MeituanCpsError:
|
||||
return []
|
||||
|
||||
# 距离最近:拉齐全部榜单轮次,合并去重,后端在【完整池子】上全局按距离由近及远排,一次性返回。
|
||||
# (距离排序必须在完整池上做、不能逐页排——这正是之前放前端不合理的根因。)
|
||||
if tab == "distance":
|
||||
with ThreadPoolExecutor(max_workers=len(_TOPIC_ROUNDS) * 2) as pool:
|
||||
futs = []
|
||||
for wm_topic, dd_topic in _TOPIC_ROUNDS:
|
||||
futs.append(pool.submit(_fetch_topic, 1, None, wm_topic))
|
||||
futs.append(pool.submit(_fetch_topic, 2, 1, dd_topic))
|
||||
raws = [f.result() for f in futs]
|
||||
seen: set[str] = set()
|
||||
cards: list[CouponCard] = []
|
||||
for raw_list in raws:
|
||||
for it in raw_list:
|
||||
card = CouponCard.from_raw(it)
|
||||
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"))
|
||||
return FeedResponse(items=cards, has_next=False, page=1)
|
||||
|
||||
# 智能推荐(rec,默认):沿用逐轮分页的混合 feed,后端过滤掉佣金率 < 3%。
|
||||
page_idx = req.page - 1
|
||||
if page_idx >= len(_TOPIC_ROUNDS):
|
||||
return FeedResponse(items=[], has_next=False, page=req.page)
|
||||
|
||||
@@ -142,8 +107,6 @@ def feed(req: FeedRequest) -> FeedResponse:
|
||||
waimai, daodian = f_wm.result(), f_dd.result()
|
||||
|
||||
items = _interleave(waimai, daodian)
|
||||
if tab == "rec":
|
||||
items = [c for c in items if _commission_pct(c) >= 3.0]
|
||||
has_next = page_idx + 1 < len(_TOPIC_ROUNDS)
|
||||
return FeedResponse(items=items, has_next=has_next, page=req.page)
|
||||
|
||||
@@ -167,37 +130,3 @@ def referral_link(req: ReferralLinkRequest) -> ReferralLinkResponse:
|
||||
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:
|
||||
# 只取有销量档的(美团很多券没销量,排不了);按销量降序、同销量再按佣金降序
|
||||
stmt = select(MeituanCoupon).where(MeituanCoupon.sale_volume_num.isnot(None))
|
||||
if req.platform is not None:
|
||||
stmt = stmt.where(MeituanCoupon.platform == req.platform)
|
||||
stmt = stmt.order_by(
|
||||
MeituanCoupon.sale_volume_num.desc(),
|
||||
MeituanCoupon.commission_percent.desc(),
|
||||
)
|
||||
rows = db.execute(stmt).scalars().all()
|
||||
|
||||
# 跨源去重(dedup_key = 品牌|名|价):按销量降序遍历,每个 dedup_key 只留第一条(=销量最高那条)。
|
||||
# 用 raw(整条原始返回)重建 CouponCard,字段与实时接口完全一致,前端无需改渲染。
|
||||
seen: set[str] = set()
|
||||
cards: list[CouponCard] = []
|
||||
for row in rows:
|
||||
if row.dedup_key in seen:
|
||||
continue
|
||||
seen.add(row.dedup_key)
|
||||
try:
|
||||
card = CouponCard.from_raw(row.raw or {})
|
||||
except Exception: # noqa: BLE001
|
||||
continue
|
||||
if card.product_view_sign:
|
||||
cards.append(card)
|
||||
|
||||
start = (req.page - 1) * req.page_size
|
||||
page_items = cards[start:start + req.page_size]
|
||||
has_next = start + req.page_size < len(cards)
|
||||
return CouponListResponse(items=page_items, has_next=has_next, search_id=None)
|
||||
|
||||
@@ -81,9 +81,6 @@ class Settings(BaseSettings):
|
||||
MT_CPS_HOST: str = "https://media.meituan.com"
|
||||
MT_CPS_TIMEOUT_SEC: int = 15
|
||||
MT_CPS_DEFAULT_SID: str = "sgbjia"
|
||||
# 美团调用走的代理。本机开发直连美团会 SSL EOF,需填 http://127.0.0.1:7897;
|
||||
# 线上国内服务器留空(=直连)。见 .env.example 与 integrations/meituan.py。
|
||||
MT_CPS_PROXY: str = ""
|
||||
|
||||
@property
|
||||
def mt_cps_configured(self) -> bool:
|
||||
|
||||
@@ -62,12 +62,8 @@ def _call(path: str, body_obj: dict[str, Any]) -> dict[str, Any]:
|
||||
}
|
||||
|
||||
url = f"{settings.MT_CPS_HOST}{path}"
|
||||
# 美团调用走 MT_CPS_PROXY(本机开发直连会 SSL EOF,必须走代理;线上留空=直连)。
|
||||
# trust_env=False:不读进程环境的 HTTP_PROXY,只认配置,避免被错误/失效代理误导。
|
||||
proxy = settings.MT_CPS_PROXY or None
|
||||
try:
|
||||
with httpx.Client(proxy=proxy, trust_env=False, timeout=settings.MT_CPS_TIMEOUT_SEC) as client:
|
||||
resp = client.post(url, content=body, headers=headers)
|
||||
resp = httpx.post(url, content=body, headers=headers, timeout=settings.MT_CPS_TIMEOUT_SEC)
|
||||
except httpx.HTTPError as e:
|
||||
logger.exception("[MT] http error calling %s", url)
|
||||
raise MeituanCpsError(f"meituan http error: {e}") from e
|
||||
|
||||
@@ -7,12 +7,7 @@ from app.models.admin import AdminAuditLog, AdminUser # noqa: F401
|
||||
from app.models.app_config import AppConfig # noqa: F401
|
||||
from app.models.comparison import ComparisonRecord # noqa: F401
|
||||
from app.models.comparison_milestone import ComparisonMilestoneClaim # noqa: F401
|
||||
from app.models.coupon_state import ( # noqa: F401
|
||||
CouponClaimRecord,
|
||||
CouponPromptEngagement,
|
||||
)
|
||||
from app.models.feedback import Feedback # noqa: F401
|
||||
from app.models.meituan_coupon import MeituanCoupon # noqa: F401
|
||||
from app.models.ops_marquee_seed import OpsMarqueeSeed # noqa: F401
|
||||
from app.models.ops_stat_config import OpsStatConfig # noqa: F401
|
||||
from app.models.price_observation import PriceObservation # noqa: F401
|
||||
|
||||
@@ -1,130 +0,0 @@
|
||||
"""领券今日状态(弹窗频控 + 领券记录)两张表。
|
||||
|
||||
- `coupon_prompt_engagement`:按 (device, 自然日) 记"今天是否对领券引导窗表达过意向"
|
||||
——点「一键领取」(claim_started) 或 点拒绝/关闭 (dismissed) 都算。切到外卖 App 时
|
||||
据此决定弹不弹:今天 engage 过就不再弹。判断维度是 **device_id**——券发到的是设备上
|
||||
登录的那个外卖账号,device 比 user 更贴近"哪个登录环境",且 device_id 全链路现成、
|
||||
不依赖领券鉴权。
|
||||
|
||||
- `coupon_claim_record`:按 (device, 券, 自然日) 记每张券的领取结果,纯沉淀(资产/画像/
|
||||
排查)。当前**不参与**"要不要领"的判断(MVP 先不去重:今天 engage 过就不弹,A 路径
|
||||
主动领则全跑)。留作以后做按券去重 / CPS 归因 / 用户画像的数据源。
|
||||
|
||||
口径:
|
||||
- 日期 = Asia/Shanghai 的自然日(claim_date / engage_date)。
|
||||
- user_id 领券登录态有就记(资产),可空、不进唯一键、不阻塞判断。
|
||||
- device_id 客户端生成存 SP,卸载重装会变 → 重装当新设备重新弹一次(产品预期)。
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import date, datetime
|
||||
|
||||
from sqlalchemy import (
|
||||
JSON,
|
||||
Date,
|
||||
DateTime,
|
||||
Index,
|
||||
Integer,
|
||||
String,
|
||||
UniqueConstraint,
|
||||
func,
|
||||
)
|
||||
from sqlalchemy.dialects.postgresql import JSONB
|
||||
from sqlalchemy.orm import Mapped, mapped_column
|
||||
|
||||
from app.db.base import Base
|
||||
|
||||
# PG 上用 JSONB,SQLite(本地/测试)退化为通用 JSON(同 price_observation / comparison)。
|
||||
_JSON = JSON().with_variant(JSONB(), "postgresql")
|
||||
|
||||
|
||||
class CouponClaimRecord(Base):
|
||||
"""单张券一天一条领取记录(资产层,当前不做去重判断)。"""
|
||||
|
||||
__tablename__ = "coupon_claim_record"
|
||||
__table_args__ = (
|
||||
# 同设备、同券、同一天只一条:领券每帧 last_coupon_result + done 帧 coupon_results
|
||||
# 会重复上报同一张券,靠它幂等 upsert。
|
||||
UniqueConstraint(
|
||||
"device_id", "coupon_id", "claim_date",
|
||||
name="uq_coupon_claim_device_coupon_date",
|
||||
),
|
||||
# 去重/统计查询按 (device, 日) 取一天所有券。
|
||||
Index("ix_coupon_claim_device_date", "device_id", "claim_date"),
|
||||
)
|
||||
|
||||
id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
|
||||
|
||||
# 判断/聚合维度。client getOrCreateDeviceId 生成,重装会变。
|
||||
device_id: Mapped[str] = mapped_column(String(64), nullable=False)
|
||||
# 登录态有就记(资产/画像),可空、不进唯一键。
|
||||
user_id: Mapped[int | None] = mapped_column(Integer, index=True, nullable=True)
|
||||
|
||||
coupon_id: Mapped[str] = mapped_column(String(64), nullable=False)
|
||||
# Asia/Shanghai 自然日。每日可领的券(签到/天天红包)靠这天然每天一条。
|
||||
claim_date: Mapped[date] = mapped_column(Date, nullable=False)
|
||||
|
||||
# success / already_claimed / failed / skipped(原样取 pricebot coupon 结果)
|
||||
status: Mapped[str] = mapped_column(String(24), nullable=False)
|
||||
vendor: Mapped[str | None] = mapped_column(String(48), nullable=True)
|
||||
coupon_name: Mapped[str | None] = mapped_column(String(128), nullable=True)
|
||||
# 这张领到几张(pricebot display_count;给不出时为 None)
|
||||
claimed_count: Mapped[int | None] = mapped_column(Integer, nullable=True)
|
||||
# 哪次任务领的,回指 pricebot work_logs / 排查
|
||||
trace_id: Mapped[str | None] = mapped_column(String(64), index=True, nullable=True)
|
||||
# failed / skipped 原因
|
||||
reason: Mapped[str | None] = mapped_column(String(255), nullable=True)
|
||||
# 杂项兜底:券的结构化信息(面额/入口/关键节点摘要等),免得加字段就迁移。
|
||||
# ⚠️ 别塞原始无障碍树(几十 KB → 行膨胀);原始大树看 trace_id 指过去的 work_logs。
|
||||
extra: Mapped[dict | None] = mapped_column(_JSON, nullable=True)
|
||||
|
||||
created_at: Mapped[datetime] = mapped_column(
|
||||
DateTime(timezone=True), server_default=func.now(), nullable=False
|
||||
)
|
||||
updated_at: Mapped[datetime] = mapped_column(
|
||||
DateTime(timezone=True), server_default=func.now(), onupdate=func.now(),
|
||||
nullable=False,
|
||||
)
|
||||
|
||||
def __repr__(self) -> str: # pragma: no cover
|
||||
return (
|
||||
f"<CouponClaimRecord device={self.device_id} coupon={self.coupon_id} "
|
||||
f"date={self.claim_date} status={self.status}>"
|
||||
)
|
||||
|
||||
|
||||
class CouponPromptEngagement(Base):
|
||||
"""按 (device, 自然日) 记"今天是否对领券引导窗表达过意向"——弹窗频控源。"""
|
||||
|
||||
__tablename__ = "coupon_prompt_engagement"
|
||||
__table_args__ = (
|
||||
# 一台设备一天一条:今天 engage 过(领或拒)就不再弹。
|
||||
UniqueConstraint(
|
||||
"device_id", "engage_date",
|
||||
name="uq_coupon_engage_device_date",
|
||||
),
|
||||
)
|
||||
|
||||
id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
|
||||
|
||||
device_id: Mapped[str] = mapped_column(String(64), nullable=False)
|
||||
user_id: Mapped[int | None] = mapped_column(Integer, index=True, nullable=True)
|
||||
# Asia/Shanghai 自然日。
|
||||
engage_date: Mapped[date] = mapped_column(Date, nullable=False)
|
||||
# claim_started(点了一键领取)/ dismissed(点了拒绝/关闭)。仅记录区分,
|
||||
# 判断只看"今天有没有这条",type 不影响弹不弹。
|
||||
engage_type: Mapped[str] = mapped_column(String(16), nullable=False)
|
||||
|
||||
created_at: Mapped[datetime] = mapped_column(
|
||||
DateTime(timezone=True), server_default=func.now(), nullable=False
|
||||
)
|
||||
updated_at: Mapped[datetime] = mapped_column(
|
||||
DateTime(timezone=True), server_default=func.now(), onupdate=func.now(),
|
||||
nullable=False,
|
||||
)
|
||||
|
||||
def __repr__(self) -> str: # pragma: no cover
|
||||
return (
|
||||
f"<CouponPromptEngagement device={self.device_id} "
|
||||
f"date={self.engage_date} type={self.engage_type}>"
|
||||
)
|
||||
@@ -1,93 +0,0 @@
|
||||
"""美团 CPS 券本地缓存表(meituan_coupon)。
|
||||
|
||||
把美团联盟 CPS 的券定时抓进本地库,供「销量 / 佣金排序」等查询**直接从库里捞、本地排序**,
|
||||
不再每次实时打美团接口——美团搜索/供给对销量排序支持差(实测乱序)、且有 402 限流和召回上限。
|
||||
|
||||
数据来源(source,北京试点):
|
||||
- search_waimai : 到家/外卖, 搜索关键词「外卖」, 翻到尽头
|
||||
- search_meishi : 到家/外卖, 搜索关键词「美食」, 翻到尽头
|
||||
- store_supply : 到店, 多业务线供给(到餐+到综+酒店+门票), 翻到尽头
|
||||
|
||||
每 5min(本地测试)/ 1h(线上)全量抓一次,按 (source, product_view_sign) upsert 存最新态,
|
||||
last_seen 每轮刷新(可据此清理长期未再出现的陈旧券)。
|
||||
|
||||
去重说明:`product_view_sign` / `sku_view_id` 都是「按召回渠道生成」的,跨渠道(搜索 vs 供给)
|
||||
会变,**不能当商品全局唯一 id**。因此:
|
||||
- 存储:按 (source, product_view_sign) upsert(同源短周期内 sign 稳定)。
|
||||
- 查询:用 `dedup_key = md5(brand|name|price)` 跨源去重,再按 `sale_volume_num` / `commission_percent` 排序。
|
||||
|
||||
价格统一存「分」(cents),与 price_report / comparison_record 一致。`raw` 保留整条原始返回
|
||||
(字段越详细越好,避免后续漏字段还要重抓)。
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import datetime
|
||||
|
||||
from sqlalchemy import JSON, DateTime, Float, Integer, String, UniqueConstraint, func
|
||||
from sqlalchemy.dialects.postgresql import JSONB
|
||||
from sqlalchemy.orm import Mapped, mapped_column
|
||||
|
||||
from app.db.base import Base
|
||||
|
||||
|
||||
class MeituanCoupon(Base):
|
||||
__tablename__ = "meituan_coupon"
|
||||
__table_args__ = (
|
||||
UniqueConstraint("source", "product_view_sign", name="uq_meituan_coupon_source_sign"),
|
||||
)
|
||||
|
||||
id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
|
||||
|
||||
# ===== 来源 / 分类 =====
|
||||
# source: search_waimai | search_meishi | store_supply
|
||||
source: Mapped[str] = mapped_column(String(16), index=True, nullable=False)
|
||||
platform: Mapped[int] = mapped_column(Integer, nullable=False) # 1 到家/外卖, 2 到店
|
||||
biz_line: Mapped[int | None] = mapped_column(Integer, nullable=True) # 到店: 1到餐 2到综 3酒店 4门票
|
||||
city_id: Mapped[str] = mapped_column(String(32), index=True, nullable=False)
|
||||
|
||||
# ===== 召回 id(按渠道生成,跨渠道会变;product_view_sign 用于换推广链) =====
|
||||
product_view_sign: Mapped[str] = mapped_column(String(128), nullable=False)
|
||||
sku_view_id: Mapped[str | None] = mapped_column(String(128), nullable=True)
|
||||
|
||||
# ===== 商品本体 =====
|
||||
name: Mapped[str | None] = mapped_column(String(256), nullable=True)
|
||||
brand_name: Mapped[str | None] = mapped_column(String(128), index=True, nullable=True)
|
||||
sell_price_cents: Mapped[int | None] = mapped_column(Integer, nullable=True)
|
||||
original_price_cents: Mapped[int | None] = mapped_column(Integer, nullable=True)
|
||||
head_url: Mapped[str | None] = mapped_column(String(512), nullable=True)
|
||||
|
||||
# ===== 销量(美团只给粗档位:热销1w+;num=排序用的下界数值,如 1w+ → 10000) =====
|
||||
sale_volume: Mapped[str | None] = mapped_column(String(32), nullable=True)
|
||||
sale_volume_num: Mapped[int | None] = mapped_column(Integer, index=True, nullable=True)
|
||||
|
||||
# ===== 佣金(percent 如 1.4 表示 1.4%;amount 为分) =====
|
||||
commission_percent: Mapped[float | None] = mapped_column(Float, index=True, nullable=True)
|
||||
commission_amount_cents: Mapped[int | None] = mapped_column(Integer, nullable=True)
|
||||
|
||||
# ===== 门店 / 距离 =====
|
||||
poi_name: Mapped[str | None] = mapped_column(String(128), nullable=True)
|
||||
available_poi_num: Mapped[int | None] = mapped_column(Integer, nullable=True)
|
||||
delivery_distance_m: Mapped[float | None] = mapped_column(Float, nullable=True)
|
||||
|
||||
# ===== 跨源去重键 + 原始返回 =====
|
||||
dedup_key: Mapped[str] = mapped_column(String(64), index=True, nullable=False)
|
||||
raw: Mapped[dict] = mapped_column(
|
||||
JSON().with_variant(JSONB(), "postgresql"), nullable=False, default=dict
|
||||
)
|
||||
|
||||
# ===== 抓取时间窗 =====
|
||||
first_seen: Mapped[datetime] = mapped_column(
|
||||
DateTime(timezone=True), server_default=func.now(), nullable=False
|
||||
)
|
||||
last_seen: Mapped[datetime] = mapped_column(
|
||||
DateTime(timezone=True), server_default=func.now(), index=True, nullable=False
|
||||
)
|
||||
updated_at: Mapped[datetime] = mapped_column(
|
||||
DateTime(timezone=True), server_default=func.now(), onupdate=func.now(), nullable=False
|
||||
)
|
||||
|
||||
def __repr__(self) -> str: # pragma: no cover
|
||||
return (
|
||||
f"<MeituanCoupon id={self.id} source={self.source} "
|
||||
f"name={self.name!r} sale={self.sale_volume} comm={self.commission_percent}>"
|
||||
)
|
||||
@@ -1,129 +0,0 @@
|
||||
"""领券今日状态读写:弹窗频控(engagement)+ 领券记录(claim)。
|
||||
|
||||
写操作按唯一键幂等 upsert,自带 commit + 并发 IntegrityError 兜底(对齐 price_observation)。
|
||||
日期口径 = Asia/Shanghai 的自然日。
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from datetime import date, datetime
|
||||
from zoneinfo import ZoneInfo
|
||||
|
||||
from sqlalchemy import select
|
||||
from sqlalchemy.exc import IntegrityError
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from app.models.coupon_state import CouponClaimRecord, CouponPromptEngagement
|
||||
|
||||
logger = logging.getLogger("shagua.coupon_state")
|
||||
|
||||
_CN_TZ = ZoneInfo("Asia/Shanghai")
|
||||
|
||||
|
||||
def today_cn() -> date:
|
||||
"""Asia/Shanghai 的自然日(领券判断的"今天")。"""
|
||||
return datetime.now(_CN_TZ).date()
|
||||
|
||||
|
||||
# ===== 弹窗频控(coupon_prompt_engagement)=====
|
||||
|
||||
def has_engaged_today(db: Session, device_id: str) -> bool:
|
||||
"""这台设备今天是否已对领券引导窗表达过意向(领或拒)。有 = 不再弹。"""
|
||||
row = db.execute(
|
||||
select(CouponPromptEngagement.id).where(
|
||||
CouponPromptEngagement.device_id == device_id,
|
||||
CouponPromptEngagement.engage_date == today_cn(),
|
||||
)
|
||||
).first()
|
||||
return row is not None
|
||||
|
||||
|
||||
def mark_engagement(
|
||||
db: Session, device_id: str, user_id: int | None, engage_type: str
|
||||
) -> None:
|
||||
"""记今日意向(claim_started / dismissed)。(device, 今天) 唯一,幂等 upsert。"""
|
||||
today = today_cn()
|
||||
row = db.execute(
|
||||
select(CouponPromptEngagement).where(
|
||||
CouponPromptEngagement.device_id == device_id,
|
||||
CouponPromptEngagement.engage_date == today,
|
||||
)
|
||||
).scalar_one_or_none()
|
||||
if row is not None:
|
||||
row.engage_type = engage_type
|
||||
if user_id is not None:
|
||||
row.user_id = user_id
|
||||
else:
|
||||
db.add(CouponPromptEngagement(
|
||||
device_id=device_id, user_id=user_id,
|
||||
engage_date=today, engage_type=engage_type,
|
||||
))
|
||||
try:
|
||||
db.commit()
|
||||
except IntegrityError:
|
||||
# 并发下另一请求刚插了同 (device, 日) → 唯一约束撞,回滚忽略(本就幂等)。
|
||||
db.rollback()
|
||||
|
||||
|
||||
# ===== 领券记录(coupon_claim_record)=====
|
||||
|
||||
def record_claims(
|
||||
db: Session,
|
||||
device_id: str,
|
||||
user_id: int | None,
|
||||
trace_id: str | None,
|
||||
results: list[dict],
|
||||
) -> int:
|
||||
"""一批券领取结果幂等写入,返回写入(新增 + 更新)条数。
|
||||
|
||||
results 单项取自 pricebot 的 last_coupon_result / done.coupon_results,识别字段:
|
||||
coupon_id(必需)/ status(必需)/ name / vendor / reason /(display_count)。
|
||||
(device, coupon_id, 今天) 唯一:重复上报同张券走更新(status 以最后一次为准)。
|
||||
"""
|
||||
today = today_cn()
|
||||
written = 0
|
||||
seen: set[str] = set() # 同批去重防御:autoflush=False 下同 coupon_id 重复会两次 add → 撞唯一约束回滚整批
|
||||
for r in results:
|
||||
coupon_id = r.get("coupon_id")
|
||||
status = r.get("status")
|
||||
if not coupon_id or not status or coupon_id in seen:
|
||||
continue # 脏数据 / 同批重复跳过
|
||||
seen.add(coupon_id)
|
||||
count = r.get("display_count")
|
||||
if count is None:
|
||||
count = r.get("claimed_count")
|
||||
row = db.execute(
|
||||
select(CouponClaimRecord).where(
|
||||
CouponClaimRecord.device_id == device_id,
|
||||
CouponClaimRecord.coupon_id == coupon_id,
|
||||
CouponClaimRecord.claim_date == today,
|
||||
)
|
||||
).scalar_one_or_none()
|
||||
if row is not None:
|
||||
row.status = status
|
||||
row.reason = r.get("reason")
|
||||
if user_id is not None:
|
||||
row.user_id = user_id
|
||||
if count is not None:
|
||||
row.claimed_count = count
|
||||
row.extra = r
|
||||
else:
|
||||
db.add(CouponClaimRecord(
|
||||
device_id=device_id, user_id=user_id,
|
||||
coupon_id=coupon_id, claim_date=today,
|
||||
status=status, vendor=r.get("vendor"), coupon_name=r.get("name"),
|
||||
claimed_count=count, trace_id=trace_id, reason=r.get("reason"),
|
||||
extra=r,
|
||||
))
|
||||
written += 1
|
||||
if written == 0:
|
||||
return 0
|
||||
try:
|
||||
db.commit()
|
||||
except IntegrityError:
|
||||
db.rollback()
|
||||
logger.warning(
|
||||
"coupon_claim 并发幂等冲突 device=%s trace=%s,回滚", device_id, trace_id
|
||||
)
|
||||
return 0
|
||||
return written
|
||||
@@ -1,21 +0,0 @@
|
||||
"""领券今日状态端点的收发模型。"""
|
||||
from __future__ import annotations
|
||||
|
||||
from pydantic import BaseModel
|
||||
|
||||
|
||||
class CouponPromptDismissIn(BaseModel):
|
||||
"""客户端拒绝/关闭领券引导窗的通知体。
|
||||
|
||||
server 据此记一条今日 engagement(dismissed)→ 今天这台设备不再弹引导窗。
|
||||
MVP 不鉴权,按 device_id 判断;user_id 登录态带上就一并记(资产),可空。
|
||||
"""
|
||||
|
||||
device_id: str
|
||||
user_id: int | None = None
|
||||
|
||||
|
||||
class CouponPromptShouldShowOut(BaseModel):
|
||||
"""切到外卖 App 时是否还应弹领券引导窗。今天已 engage(领或拒)过 → false。"""
|
||||
|
||||
should_show: bool
|
||||
@@ -134,12 +134,6 @@ class FeedRequest(BaseModel):
|
||||
latitude: float = Field(..., description="纬度")
|
||||
page: int = Field(1, ge=1)
|
||||
page_size: int = Field(20, ge=1, le=20)
|
||||
# 筛选/排序口径,后端据此处理后返回(前端不再自己筛/排):
|
||||
# rec = 智能推荐:榜单混合 feed 去掉佣金率 < 3%(分页)
|
||||
# distance = 距离最近:拉齐全部轮次后全局按距离由近及远(一次性返回, has_next=False)
|
||||
# 留空/其它 = 原混合 feed 不筛(老客户端兼容,新 app 会显式传 tab)
|
||||
# (销量最高 sales 走 /coupons 同城热销,不在本接口)
|
||||
tab: str = Field("", description="rec 智能推荐 / distance 距离最近 / 空=不筛(兼容)")
|
||||
|
||||
|
||||
class FeedResponse(BaseModel):
|
||||
@@ -148,13 +142,6 @@ class FeedResponse(BaseModel):
|
||||
page: int = 1
|
||||
|
||||
|
||||
class TopSalesRequest(BaseModel):
|
||||
"""销量最高 tab:从离线库 meituan_coupon 按销量降序取(不实时打美团)。"""
|
||||
page: int = Field(1, ge=1)
|
||||
page_size: int = Field(20, ge=1, le=50)
|
||||
platform: int | None = Field(None, description="可选: 1只外卖 / 2只到店; 不填=全部(全城销量)")
|
||||
|
||||
|
||||
# ───────────────── 换链 请求 / 响应 ─────────────────
|
||||
|
||||
class ReferralLinkRequest(BaseModel):
|
||||
|
||||
@@ -1,331 +0,0 @@
|
||||
"""美团 CPS 券定时抓取入库(北京试点)。
|
||||
|
||||
把 3 路券抓进 meituan_coupon 表,供「销量/佣金排序」从库里捞、本地排序,不再实时打美团:
|
||||
1. search_waimai : 到家/外卖, 搜「外卖」 翻到尽头
|
||||
2. search_meishi : 到家/外卖, 搜「美食」 翻到尽头
|
||||
3. store_supply : 到店, 多业务线供给(到餐+到综+酒店+门票) 翻到尽头
|
||||
|
||||
按 (source, product_view_sign) upsert 存最新态;last_seen 每轮刷新。带文件锁,防止
|
||||
上一轮没跑完下一轮又起(本地 5~10min、跨进程 cron 都安全)。
|
||||
|
||||
用法:
|
||||
# 单轮(打通验证 / 给 cron 用,线上每 1h 一次)
|
||||
python -m scripts.pull_meituan_coupons --once
|
||||
|
||||
# 本地循环(默认每 10min 一轮)
|
||||
python -m scripts.pull_meituan_coupons --loop --interval 600
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import hashlib
|
||||
import os
|
||||
import re
|
||||
import sys
|
||||
import time
|
||||
from datetime import datetime, timedelta, timezone
|
||||
|
||||
# Windows 控制台按 UTF-8 输出中文/¥
|
||||
try:
|
||||
sys.stdout.reconfigure(encoding="utf-8") # type: ignore[attr-defined]
|
||||
except Exception: # noqa: BLE001
|
||||
pass
|
||||
|
||||
from sqlalchemy import delete, func, select
|
||||
from sqlalchemy.dialects.postgresql import insert as pg_insert
|
||||
|
||||
from app.db.session import SessionLocal
|
||||
from app.integrations.meituan import MeituanCpsError, _call
|
||||
from app.models.meituan_coupon import MeituanCoupon
|
||||
|
||||
CITY_BEIJING = "WKV2HMXUEK634WP64CUCUQGM64"
|
||||
QUERY_PATH = "/cps_open/common/api/v1/query_coupon"
|
||||
PAGE_SIZE = 20
|
||||
MAX_PAGES = 80 # 单路安全上限(搜索 ~52 页、供给 ~70 页)
|
||||
PAGE_SLEEP = 0.35 # 页间配速,缓解 402
|
||||
RETRY = 7
|
||||
LOCK_FILE = "data/.meituan_etl.lock"
|
||||
LOCK_STALE_SEC = 30 * 60 # 锁超过 30min 视为陈旧(进程异常退出残留),自动接管
|
||||
|
||||
SOURCES = [
|
||||
{"code": "search_waimai", "label": "外卖·搜外卖", "kind": "search", "platform": 1, "keyword": "外卖"},
|
||||
{"code": "search_meishi", "label": "外卖·搜美食", "kind": "search", "platform": 1, "keyword": "美食"},
|
||||
{"code": "store_supply", "label": "到店·多业务线供给", "kind": "supply",
|
||||
"platform": 2, "biz_lines": [1, 2, 3, 4]},
|
||||
]
|
||||
|
||||
|
||||
# ───────────────────────── 美团调用 ─────────────────────────
|
||||
|
||||
def _call_retry(body: dict) -> dict | None:
|
||||
"""打美团,402/频繁退避重试;其它错误打印并放弃本页。"""
|
||||
for a in range(RETRY):
|
||||
try:
|
||||
return _call(QUERY_PATH, body)
|
||||
except MeituanCpsError as e:
|
||||
msg = str(e)
|
||||
if "402" in msg or "频繁" in msg:
|
||||
time.sleep(2.5 * (a + 1))
|
||||
continue
|
||||
print(f" [warn] meituan: {msg[:80]}")
|
||||
return None
|
||||
except Exception as e: # noqa: BLE001
|
||||
print(f" [warn] {type(e).__name__}: {str(e)[:60]}")
|
||||
time.sleep(2.0 * (a + 1))
|
||||
return None
|
||||
|
||||
|
||||
def _pull_search(platform: int, keyword: str) -> list[dict]:
|
||||
rows: list[dict] = []
|
||||
sid = None
|
||||
pg = 1
|
||||
while pg <= MAX_PAGES:
|
||||
body = {"platform": platform, "searchText": keyword, "cityId": CITY_BEIJING, "pageSize": PAGE_SIZE}
|
||||
if sid:
|
||||
body["searchId"] = sid
|
||||
else:
|
||||
body["pageNo"] = pg
|
||||
r = _call_retry(body)
|
||||
if not r:
|
||||
break
|
||||
sid = r.get("searchId")
|
||||
data = r.get("data") or []
|
||||
rows.extend(data)
|
||||
if not r.get("hasNext") or not data:
|
||||
break
|
||||
pg += 1
|
||||
time.sleep(PAGE_SLEEP)
|
||||
return rows
|
||||
|
||||
|
||||
def _pull_supply(platform: int, biz_lines: list[int]) -> list[dict]:
|
||||
rows: list[dict] = []
|
||||
sid = None
|
||||
biz_param = [{"bizLine": b} for b in biz_lines]
|
||||
for _ in range(MAX_PAGES):
|
||||
body = {
|
||||
"multipleSupplyList": [{"platform": platform, "bizLineParamList": biz_param}],
|
||||
"cityId": CITY_BEIJING,
|
||||
"sortField": 2, # 供给查询 sortField 必填;我们入库后本地再排,这里给个默认
|
||||
"pageSize": PAGE_SIZE,
|
||||
}
|
||||
if sid:
|
||||
body["searchId"] = sid
|
||||
r = _call_retry(body)
|
||||
if not r:
|
||||
break
|
||||
sid = r.get("searchId")
|
||||
data = r.get("data") or []
|
||||
rows.extend(data)
|
||||
if not r.get("hasNext") or not data:
|
||||
break
|
||||
time.sleep(PAGE_SLEEP)
|
||||
return rows
|
||||
|
||||
|
||||
# ───────────────────────── 解析 ─────────────────────────
|
||||
|
||||
_SV_RE = re.compile(r"(\d+(?:\.\d+)?)\s*(w|万|k)?", re.I)
|
||||
|
||||
|
||||
def _sale_volume_num(s: str | None) -> int | None:
|
||||
"""'热销1w+' → 10000, '热销500+' → 500(取下界,排序用)。"""
|
||||
if not s:
|
||||
return None
|
||||
m = _SV_RE.search(str(s))
|
||||
if not m:
|
||||
return None
|
||||
v = float(m.group(1))
|
||||
u = (m.group(2) or "").lower()
|
||||
if u in ("w", "万"):
|
||||
v *= 10000
|
||||
elif u == "k":
|
||||
v *= 1000
|
||||
return int(v)
|
||||
|
||||
|
||||
def _to_cents(yuan) -> int | None:
|
||||
if yuan in (None, "", "null"):
|
||||
return None
|
||||
try:
|
||||
return round(float(yuan) * 100)
|
||||
except (TypeError, ValueError):
|
||||
return None
|
||||
|
||||
|
||||
def _parse_item(item: dict, source: dict) -> dict | None:
|
||||
cpd = item.get("couponPackDetail") or {}
|
||||
br = item.get("brandInfo") or {}
|
||||
ci = item.get("commissionInfo") or {}
|
||||
poi = item.get("availablePoiInfo") or {}
|
||||
dp = item.get("deliverablePoiInfo") or {}
|
||||
|
||||
sign = cpd.get("productViewSign") or cpd.get("skuViewId")
|
||||
if not sign:
|
||||
return None
|
||||
|
||||
name = cpd.get("name") or ""
|
||||
brand = br.get("brandName") or ""
|
||||
price_cents = _to_cents(cpd.get("sellPrice"))
|
||||
|
||||
comm_pct = None
|
||||
if ci.get("commissionPercent") is not None:
|
||||
try:
|
||||
comm_pct = float(ci["commissionPercent"]) / 100.0 # 140 → 1.4(%)
|
||||
except (TypeError, ValueError):
|
||||
comm_pct = None
|
||||
|
||||
dist = dp.get("deliveryDistance")
|
||||
try:
|
||||
dist = float(dist) if dist not in (None, "", "null") else None
|
||||
except (TypeError, ValueError):
|
||||
dist = None
|
||||
|
||||
dedup_raw = f"{brand}|{name}|{price_cents}"
|
||||
return {
|
||||
"source": source["code"],
|
||||
"platform": source["platform"],
|
||||
"biz_line": item.get("bizLine") or cpd.get("bizLine"),
|
||||
"city_id": CITY_BEIJING,
|
||||
"product_view_sign": str(sign)[:128],
|
||||
"sku_view_id": cpd.get("skuViewId"),
|
||||
"name": (name[:256] or None),
|
||||
"brand_name": (brand[:128] or None),
|
||||
"sell_price_cents": price_cents,
|
||||
"original_price_cents": _to_cents(cpd.get("originalPrice")),
|
||||
"head_url": ((cpd.get("headUrl") or "").split("@")[0][:512] or None),
|
||||
"sale_volume": cpd.get("saleVolume"),
|
||||
"sale_volume_num": _sale_volume_num(cpd.get("saleVolume")),
|
||||
"commission_percent": comm_pct,
|
||||
"commission_amount_cents": _to_cents(ci.get("commission")),
|
||||
"poi_name": ((dp.get("poiName") or "")[:128] or None),
|
||||
"available_poi_num": poi.get("availablePoiNum"),
|
||||
"delivery_distance_m": dist,
|
||||
"dedup_key": hashlib.md5(dedup_raw.encode("utf-8")).hexdigest(),
|
||||
"raw": item,
|
||||
}
|
||||
|
||||
|
||||
# ───────────────────────── 入库(upsert) ─────────────────────────
|
||||
|
||||
def _upsert(db, rows: list[dict], now: datetime) -> tuple[int, int]:
|
||||
"""按 (source, product_view_sign) upsert。返回 (入库去重后条数, 本源内重复条数)。"""
|
||||
if not rows:
|
||||
return 0, 0
|
||||
dedup: dict[tuple, dict] = {}
|
||||
for r in rows:
|
||||
dedup[(r["source"], r["product_view_sign"])] = r # 同轮内同键保留最后一条
|
||||
payload = list(dedup.values())
|
||||
for r in payload:
|
||||
r["last_seen"] = now
|
||||
r["updated_at"] = now
|
||||
|
||||
chunk = 500
|
||||
for i in range(0, len(payload), chunk):
|
||||
part = payload[i:i + chunk]
|
||||
stmt = pg_insert(MeituanCoupon).values(part)
|
||||
update_cols = {
|
||||
c: getattr(stmt.excluded, c)
|
||||
for c in part[0]
|
||||
if c not in ("source", "product_view_sign", "first_seen")
|
||||
}
|
||||
stmt = stmt.on_conflict_do_update(
|
||||
constraint="uq_meituan_coupon_source_sign", set_=update_cols
|
||||
)
|
||||
db.execute(stmt)
|
||||
db.commit()
|
||||
return len(payload), len(rows) - len(payload)
|
||||
|
||||
|
||||
# ───────────────────────── 运行锁 ─────────────────────────
|
||||
|
||||
def _acquire_lock() -> bool:
|
||||
os.makedirs(os.path.dirname(LOCK_FILE) or ".", exist_ok=True)
|
||||
try:
|
||||
fd = os.open(LOCK_FILE, os.O_CREAT | os.O_EXCL | os.O_WRONLY)
|
||||
os.write(fd, f"{os.getpid()} {time.time()}".encode())
|
||||
os.close(fd)
|
||||
return True
|
||||
except FileExistsError:
|
||||
try:
|
||||
with open(LOCK_FILE) as f:
|
||||
parts = f.read().split()
|
||||
ts = float(parts[1]) if len(parts) > 1 else 0.0
|
||||
if time.time() - ts > LOCK_STALE_SEC:
|
||||
os.remove(LOCK_FILE)
|
||||
return _acquire_lock()
|
||||
except Exception: # noqa: BLE001
|
||||
pass
|
||||
return False
|
||||
|
||||
|
||||
def _release_lock() -> None:
|
||||
try:
|
||||
os.remove(LOCK_FILE)
|
||||
except OSError:
|
||||
pass
|
||||
|
||||
|
||||
# ───────────────────────── 主流程 ─────────────────────────
|
||||
|
||||
def run_once(prune_hours: int = 24) -> None:
|
||||
if not _acquire_lock():
|
||||
print(f"[{datetime.now():%H:%M:%S}] 上一轮还在跑(锁占用),跳过本轮")
|
||||
return
|
||||
t0 = time.time()
|
||||
now = datetime.now(timezone.utc)
|
||||
db = SessionLocal()
|
||||
try:
|
||||
total = 0
|
||||
for src in SOURCES:
|
||||
ts = time.time()
|
||||
if src["kind"] == "search":
|
||||
items = _pull_search(src["platform"], src["keyword"])
|
||||
else:
|
||||
items = _pull_supply(src["platform"], src["biz_lines"])
|
||||
parsed = [p for p in (_parse_item(it, src) for it in items) if p]
|
||||
up, dup = _upsert(db, parsed, now)
|
||||
total += up
|
||||
print(f" {src['label']:18} 抓{len(items):5} 解析{len(parsed):5} "
|
||||
f"入库{up:5} (本源去重{dup:3}) {time.time() - ts:4.0f}s")
|
||||
# 清理长期未再出现的陈旧券(美团 sign 轮换 / 券下架后的残留),默认 24h 宽限
|
||||
if prune_hours and prune_hours > 0:
|
||||
cutoff = now - timedelta(hours=prune_hours)
|
||||
pruned = db.execute(
|
||||
delete(MeituanCoupon).where(MeituanCoupon.last_seen < cutoff)
|
||||
).rowcount
|
||||
db.commit()
|
||||
if pruned:
|
||||
print(f" 清理陈旧券(>{prune_hours}h 未再出现): {pruned} 条")
|
||||
cnt = db.execute(select(func.count()).select_from(MeituanCoupon)).scalar()
|
||||
print(f"[{datetime.now():%H:%M:%S}] 本轮完成: 入库 {total} 条, 表总计 {cnt} 行, "
|
||||
f"用时 {time.time() - t0:.0f}s")
|
||||
finally:
|
||||
db.close()
|
||||
_release_lock()
|
||||
|
||||
|
||||
def main() -> None:
|
||||
ap = argparse.ArgumentParser(description="美团 CPS 券定时抓取入库")
|
||||
ap.add_argument("--once", action="store_true", help="只跑一轮(默认)")
|
||||
ap.add_argument("--loop", action="store_true", help="循环跑")
|
||||
ap.add_argument("--interval", type=int, default=600, help="循环间隔秒(默认 600=10min)")
|
||||
ap.add_argument("--prune-hours", type=int, default=24,
|
||||
help="清理超过 N 小时未再出现的陈旧券(默认 24;0=不清理)")
|
||||
args = ap.parse_args()
|
||||
|
||||
if args.loop:
|
||||
print(f"循环模式: 每 {args.interval}s 一轮 (Ctrl-C 退出)")
|
||||
while True:
|
||||
try:
|
||||
run_once(args.prune_hours)
|
||||
except Exception as e: # noqa: BLE001
|
||||
print(f"[{datetime.now():%H:%M:%S}] 本轮异常: {type(e).__name__}: {e}")
|
||||
_release_lock()
|
||||
time.sleep(args.interval)
|
||||
else:
|
||||
run_once(args.prune_hours)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
Reference in New Issue
Block a user