| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398 |
- # -*- coding: utf-8 -*-
- # Author : Charley
- # Python : 3.12.10
- # Date : 2026/08/11
- """得卡 DECA · 随机团 teams 明细高频采集(选队随机 + 剩余随机)。
- 背景:
- 卡牌拼团分两种"随机团"玩法,同一团会经历两阶段:
- · 选队随机(playTypeName 以「选队随机」开头):每支球队一价,team-options 接口(带 token)
- 返回逐队 cardCount / availableStock / unitPrice;已售数 = cardCount − availableStock。
- · 剩余随机(playTypeName == 剩余随机):卖不动后转成的兜底阶段。团购价变成"剩余池加权
- 均价"、买家随机开一张剩余卡。team-options **对剩余随机永远返 29000**、详情 snapshot
- 冻结在转换那一刻(每队只有 unitPrice/availableStock,无原始 cardCount)。
- 关键结论:**一旦转成剩余随机,原始 cardCount 就再也没接口能补回**。所以在选队阶段必须持续
- 抓 team-options 存表(deca_groupbuy_team_record),转剩余后我们才能用「存表的原始总价 − 当前
- 剩余货值」算出实时已售额。转剩余前没被我们抓到的老团,team_total_amount 留 NULL、报告端
- 标注"原始数据缺失"。
- 一轮做什么:
- 1) 从 deca_onsale_product_record 拉当前在售·选队随机团 → team-options →
- upsert 到 teams 表(data_source=team_options,同 (code, team_id) 覆盖成最新,captured_at
- 记本轮时刻) → 算 Σ 单价×(cardCount − availableStock) 写回 team_total_amount。
- 2) 从 deca_onsale_product_record 拉当前在售·剩余随机团 → 详情 snapshot 首次存表
- (data_source=snapshot;冻结不变,重复见到直接跳过) → 若库里有该团选队阶段的原始
- cardCount,则算 Σ 单价×cardCount − detail.unitPrice × detail.availableStock 写回
- team_total_amount;否则留 NULL。
- 登录态:
- - team-options 强制带 token(不带返 10002);
- - 详情接口 groupbuy/detail 免登录。
- - 与 alert / sold_daily 走同一账号池(每批随机换号 + 各号专属 IP),无共享 token.json;
- 续签在 `AccountPool.ensure_access` 里按号独立轮换 refreshToken,由 `report_failure` 兜底判死。
- 运行:项目根目录 `python deca_team_spider.py`(2026/08/20 由 on_sale/ 迁回根目录)
- """
- import os
- import sys
- import time
- from datetime import datetime
- from decimal import Decimal, ROUND_HALF_UP
- import schedule
- from loguru import logger
- from tenacity import retry, stop_after_attempt, wait_fixed
- # 挂靠新项目根:sys.path 指向 common、CWD 对齐新根(application.yml / ./logs / 账号池 DB 均在此)
- _ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
- sys.path.insert(0, os.path.join(_ROOT, "common"))
- os.chdir(_ROOT)
- import deca_sold_core as core # noqa: E402
- from mysql_pool import MySQLConnectionPool # noqa: E402
- # ==================== 配置 ====================
- INTERVAL_SEC = 300 # 采集间隔:5 分钟一轮(选队随机变化较快、又不至于压接口)
- BETWEEN_ITEM_SEC = 0.3 # 单商品之间的最小间隔,避免瞬时限流
- MAX_ITEM_PER_ROUND = 500 # 单轮采集商品数上限(保护;正常远小于此)
- ONSALE_TABLE = "deca_onsale_product_record"
- TEAM_TABLE = "deca_groupbuy_team_record"
- # 只监控重点商家的选队随机(2026/09/10):原扫全站 is_on_sale=1(实测 449 个/99 家含大量陈旧未下架
- # 死团,每 5 分钟对每个打带 token 的 team-options,白烧 token 扩风控面)。改为只扫这 5 家(与购买记录
- # buy_record_spider.MERCHANT_IDS 一致),请求量骤降、省 token 降风控面。
- WATCH_MERCHANT_IDS = [
- "881226408", # 魔都兄弟球星卡
- "274584650", # 卡皇拆卡
- "538252487", # 尼卡拆卡
- "591544726", # 文泰卡屋
- "606370597", # 魔都兄弟综合体育
- ]
- # ==================== 日志 ====================
- logger.remove()
- logger.add("./logs/team_spider_{time:YYYYMMDD}.log", encoding="utf-8", rotation="00:00",
- format="[{time:YYYY-MM-DD HH:mm:ss.SSS}] {level} {message}",
- level="DEBUG", retention="7 day")
- # logger.add(sys.stderr, level="INFO",
- # format="[{time:HH:mm:ss}] {level} {message}")
- def after_log(retry_state):
- """tenacity 重试回调,业务函数首参约定为 log。
- Args:
- retry_state: tenacity 传入的 RetryCallState,含调用参数与结果。
- """
- log = retry_state.args[0] if retry_state.args else logger
- if retry_state.outcome.failed:
- log.warning(f"Function '{retry_state.fn.__name__}', Attempt {retry_state.attempt_number} Times")
- else:
- log.info(f"Function '{retry_state.fn.__name__}', Attempt {retry_state.attempt_number} succeeded")
- def _round2(x) -> Decimal | None:
- """Decimal 化并四舍五入到 2 位小数。
- Args:
- x: 数值(int/float/str/Decimal),None/异常直接返回 None。
- Returns:
- Decimal | None: 保留 2 位;None 表示不可算。
- """
- if x is None:
- return None
- try:
- return Decimal(str(x)).quantize(Decimal("0.01"), rounding=ROUND_HALF_UP)
- except Exception:
- return None
- def fetch_onsale_random(log, pool) -> tuple[list, list]:
- """从 deca_onsale_product_record 查当前在售·选队随机 / 剩余随机团 code 列表。
- Args:
- log: 日志对象。
- pool: MySQL 连接池。
- Returns:
- tuple[list[str], list[str]]: (选队随机 codes, 剩余随机 codes)。
- """
- rows = pool.select_all(
- f"SELECT product_code, play_type_name FROM {ONSALE_TABLE} "
- f"WHERE is_on_sale=1 AND play_type_name IS NOT NULL "
- f"AND merchant_user_id IN ({','.join(['%s'] * len(WATCH_MERCHANT_IDS))}) " # 只扫重点 5 家(省 token 降风控面)
- f"LIMIT {MAX_ITEM_PER_ROUND}", tuple(WATCH_MERCHANT_IDS)) or []
- xd, sy = [], []
- for code, ptn in rows:
- if not ptn:
- continue
- if "选队随机" in ptn:
- xd.append(code)
- elif "剩余随机" in ptn:
- sy.append(code)
- log.info(f"[发现] 在售·选队随机 {len(xd)} / 剩余随机 {len(sy)}")
- return xd, sy
- @retry(stop=stop_after_attempt(3), wait=wait_fixed(2), after=after_log)
- def fetch_team_options(log, code: str) -> list | None:
- """打 team-options 接口拿一个团的逐队 teams(带 token)。
- Args:
- log: 日志对象。
- code (str): 商品编码。
- Returns:
- list | None: teams 列表;接口报 29000(已转剩余随机)或失败时返回 None(None 是业务态,
- 不视为需要重试的错误——只有网络级异常才会被 tenacity 拦下重试)。
- """
- r = core.do_request(log, "/api/v1/app/groupbuy/team-options", {"code": code}, need_auth=True)
- if not r or r.get("code") != 0:
- return None # 29000 = 商家开启剩余随机中,正常业务态
- return (r.get("data") or {}).get("list") or []
- @retry(stop=stop_after_attempt(3), wait=wait_fixed(2), after=after_log)
- def fetch_detail(log, code: str) -> dict | None:
- """打 detail 接口拿详情(免 token)。
- Args:
- log: 日志对象。
- code (str): 商品编码。
- Returns:
- dict | None: data 字典;失败返回 None。
- """
- d = core.do_request(log, "/api/v1/app/groupbuy/detail", {"code": code}, need_auth=False)
- return (d or {}).get("data") or None
- def upsert_team_row(pool, code: str, ptn: str, source: str, team: dict,
- captured_at: str, snap_total: int | None = None) -> None:
- """把一条 team 明细 upsert 进 deca_groupbuy_team_record。
- 唯一键 (product_code, team_id, data_source) 冲突时按最新覆盖——
- 这就把「在售团每轮覆盖成最新」自动做掉了。
- Args:
- pool: MySQL 连接池。
- code (str): 商品编码。
- ptn (str): 采集时的 playTypeName。
- source (str): 数据来源 'team_options' / 'snapshot'。
- team (dict): 单条 team 数据(team-options 项 或 snapshot.teams 项)。
- captured_at (str): 采集时刻 YYYY-MM-DD HH:MM:SS。
- snap_total (int | None, optional): 剩余随机 snapshot.totalQuantity。Defaults to None。
- """
- tid = team.get("teamId")
- if tid is None:
- return
- cc = team.get("cardCount") # 选队阶段有;剩余快照无 → NULL
- av = team.get("availableStock")
- up = _round2(team.get("unitPrice"))
- sold = (cc - av) if (cc is not None and av is not None) else None
- pool._execute(
- f"INSERT INTO {TEAM_TABLE} "
- f"(product_code, play_type_name, data_source, team_id, team_name_en, team_name_zh, "
- f" team_logo_image_url, unit_price, card_count, available_stock, sold_count, "
- f" snapshot_total_quantity, captured_at) "
- f"VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s) "
- f"ON DUPLICATE KEY UPDATE "
- f" play_type_name=VALUES(play_type_name), team_name_en=VALUES(team_name_en), "
- f" team_name_zh=VALUES(team_name_zh), team_logo_image_url=VALUES(team_logo_image_url), "
- f" unit_price=VALUES(unit_price), card_count=VALUES(card_count), "
- f" available_stock=VALUES(available_stock), sold_count=VALUES(sold_count), "
- f" snapshot_total_quantity=VALUES(snapshot_total_quantity), "
- f" captured_at=VALUES(captured_at)",
- (code, ptn, source, tid, team.get("teamNameEn"), team.get("teamNameZh"),
- team.get("teamLogoImageUrl"), up, cc, av, sold, snap_total, captured_at),
- commit=True)
- def has_snapshot(pool, code: str) -> bool:
- """判断某团是否已存过剩余随机 snapshot(冻结值,只需存一次)。
- Args:
- pool: MySQL 连接池。
- code (str): 商品编码。
- Returns:
- bool: True 已存过、可跳过;False 未存。
- """
- row = pool.select_one(
- f"SELECT 1 FROM {TEAM_TABLE} WHERE product_code=%s AND data_source='snapshot' LIMIT 1",
- (code,))
- return bool(row)
- def get_original_total(pool, code: str) -> Decimal | None:
- """从 teams 表查某团选队阶段抓过的「原始总价」= Σ 单价×cardCount。
- Args:
- pool: MySQL 连接池。
- code (str): 商品编码。
- Returns:
- Decimal | None: 原始总价;缺 team-options 记录时返回 None。
- """
- row = pool.select_one(
- f"SELECT SUM(unit_price * card_count) FROM {TEAM_TABLE} "
- f"WHERE product_code=%s AND data_source='team_options' AND card_count IS NOT NULL",
- (code,))
- if not row or row[0] is None:
- return None
- return _round2(row[0])
- def update_team_total(pool, code: str, amount: Decimal | None) -> None:
- """把算好的团总价写回 deca_onsale_product_record.team_total_amount。
- Args:
- pool: MySQL 连接池。
- code (str): 商品编码。
- amount (Decimal | None): 总价;None 会显式清空(表示"数据缺失")。
- """
- pool.update_one(
- f"UPDATE {ONSALE_TABLE} SET team_total_amount=%s WHERE product_code=%s",
- (amount, code))
- def process_xuandui(log, pool, code: str, captured_at: str) -> bool:
- """处理一个在售·选队随机团:抓 team-options → 存表 → 算总价 → 回写商品。
- Args:
- log: 日志对象。
- pool: MySQL 连接池。
- code (str): 商品编码。
- captured_at (str): 本轮采集时刻。
- Returns:
- bool: True 成功;False 接口不可用(如刚转成剩余随机)。
- """
- teams = fetch_team_options(log, code)
- if not teams:
- return False # 可能这一瞬间刚转成剩余随机,本轮跳过;下轮以剩余随机身份进另一分支
- total = Decimal("0")
- ptn = "选队随机" # 精确 name 已在商品行 play_type_name;此处只作 teams 表内的存档
- for t in teams:
- upsert_team_row(pool, code, ptn, "team_options", t, captured_at)
- cc = t.get("cardCount") or 0
- av = t.get("availableStock") or 0
- up = Decimal(str(t.get("unitPrice") or "0"))
- total += up * (cc - av)
- update_team_total(pool, code, _round2(total))
- return True
- def process_shengyu(log, pool, code: str, captured_at: str) -> bool:
- """处理一个在售·剩余随机团:首次存 snapshot;每轮从 detail 取实时剩余算总价回写。
- 公式:team_total_amount = 原始总价(存表 Σ 单价×cardCount) − detail.unitPrice × detail.availableStock。
- 若库里没有该团选队阶段的 cardCount 记录,team_total_amount 留 NULL 表示"数据缺失"。
- Args:
- log: 日志对象。
- pool: MySQL 连接池。
- code (str): 商品编码。
- captured_at (str): 本轮采集时刻(仅首次存 snapshot 时用)。
- Returns:
- bool: True 处理成功;False 详情拉失败。
- """
- dd = fetch_detail(log, code)
- if not dd:
- return False
- # 1) 首次存 snapshot(冻结不变,只存一次)
- if not has_snapshot(pool, code):
- snap = dd.get("remainingRandomTeamSnapshot") or {}
- for t in (snap.get("teams") or []):
- upsert_team_row(pool, code, "剩余随机", "snapshot", t, captured_at,
- snap_total=snap.get("totalQuantity"))
- # 2) 每轮重算实时总价
- original = get_original_total(pool, code)
- if original is None:
- # 转剩余随机前没被我们抓到过 team-options → 拿不到原始总价,显式留 NULL
- update_team_total(pool, code, None)
- log.warning(f"[数据缺失] {code} 无选队阶段 team-options 记录,team_total_amount=NULL")
- return True
- up_now = Decimal(str(dd.get("unitPrice") or "0"))
- av_now = int(dd.get("availableStock") or 0)
- remaining_value = up_now * av_now
- total = original - remaining_value
- if total < 0:
- # 边界:原始总价是"名义价",理论上不会 <0;出现即为数据异常,记日志、置 None
- log.warning(f"[异常] {code} original={original} - remaining={remaining_value} <0,置 NULL")
- update_team_total(pool, code, None)
- else:
- update_team_total(pool, code, _round2(total))
- return True
- @retry(stop=stop_after_attempt(100), wait=wait_fixed(3600), after=after_log)
- def main_task(log):
- """采集主函数:一轮遍历在售随机团、更新 teams 表与 team_total_amount。
- 挂了每小时重试(无人值守);单轮内各商品独立 try/except,单个失败不拖垮整轮。
- Args:
- log: 日志对象。
- Raises:
- RuntimeError: 数据库连接池异常时抛出以触发重试。
- """
- log.info(f"开始运行 {sys._getframe().f_code.co_name}" + "." * 40)
- pool = MySQLConnectionPool(log=log)
- core.init_account_pool(pool, task_tag="team") # 启用账号池:need_auth 请求走 20 号池 + 各号专属 IP
- if not pool.check_pool_health():
- log.error("数据库连接池异常")
- raise RuntimeError("数据库连接池异常")
- try:
- xd_codes, sy_codes = fetch_onsale_random(log, pool)
- captured_at = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
- # 选队随机:每轮全量刷新 team-options
- ok_xd = fail_xd = 0
- for c in xd_codes:
- try:
- if process_xuandui(log, pool, c, captured_at):
- ok_xd += 1
- else:
- fail_xd += 1
- except Exception as e:
- fail_xd += 1
- log.error(f"[选队随机·失败] {c}: {e}")
- time.sleep(BETWEEN_ITEM_SEC)
- log.info(f"[选队随机] 处理完成: 成功 {ok_xd} 失败 {fail_xd}")
- # 剩余随机:每轮从 detail 取实时剩余、更新总价
- ok_sy = fail_sy = 0
- for c in sy_codes:
- try:
- if process_shengyu(log, pool, c, captured_at):
- ok_sy += 1
- else:
- fail_sy += 1
- except Exception as e:
- fail_sy += 1
- log.error(f"[剩余随机·失败] {c}: {e}")
- time.sleep(BETWEEN_ITEM_SEC)
- log.info(f"[剩余随机] 处理完成: 成功 {ok_sy} 失败 {fail_sy}")
- except Exception as e:
- log.error(f"{sys._getframe().f_code.co_name} error: {e}")
- finally:
- log.info(f"{sys._getframe().f_code.co_name} 运行结束,等待下一轮" + "." * 20)
- def schedule_task():
- """定时入口:每 INTERVAL_SEC 秒跑一次 main_task。启动时立即跑一次。"""
- main_task(log=logger) # 启动即刻跑一轮,避免等 5 分钟才开始
- schedule.every(INTERVAL_SEC).seconds.do(main_task, log=logger)
- while True:
- schedule.run_pending()
- time.sleep(1)
- if __name__ == "__main__":
- schedule_task()
|