# -*- coding: utf-8 -*- # Author : Charley # Python : 3.12.10 # Date : 2026/08/19 """集物星球在售监控 / 自动上架提醒:轮询全站在售、过滤关注商家(Jake/九叔),三类告警推企微。 逻辑与参考项目 deca_spider 的 onsale_alert_spider 一致,只把取数层换成集物的接口: - 取数:/search/app/index/top 无按商家过滤参数,需对 productType∈{1,2,3,4} 各翻全部页, 合并去重后按 corp_info_id 客户端过滤出关注商家的在售商品(返回 (list, ok),取数不全整轮跳过)。 - 三类告警:新品上架(new) / 进度过半(half) / 一车结束(ended)。 - 去重靠 jw_onsale_alert_record 的三个标记位 + 进程内内存集合 _seen_onsale: · new / ended:send-then-mark(发送成功才置位,失败下轮补发); · half:先置位再发(容忍偶发漏发,简化逻辑); · _seen_onsale 只让「本进程运行后见过在售」的商品参与结束对账,避免冷启动补发历史老车。 - 新品判断用列表自带 soldTime(上架时间):soldTime ≥ 本场窗口起点(最近一个 RUN_START)才算「本场刚上架」 (对齐 deca 最新逻辑);更早上架的老货静默建档(new_notified=1)不发消息,防冷启动刷屏。 - 结束判定:详情 surplusStock≤0(售罄) 或 已过 offShelfTime(下架/销售结束时间),二者任一即结束(对齐 deca)。 """ import sys import time import random import argparse from collections import defaultdict from datetime import datetime, timedelta from loguru import logger from tenacity import retry, stop_after_attempt, wait_fixed from mysql_pool import MySQLConnectionPool import jiwu_core as core import auto_send_wx_msg as wx # ---------- 业务配置 ---------- WATCH_CORPS = {100716: "Jake球星卡", 100715: "九叔的喷火龙"} # 只盯这些商家(corpInfoId→名称) PRODUCT_TYPES = ["1", "2", "3", "4"] # 商品类型:1福袋 2变风盒 3错版卡 4原盒 SYSTEM_BUSINESS_TYPE = 5 # 业务线:集卡 PAGE_LIMIT = 20 # 翻页每页条数 MAX_PAGES = 200 # 单类型翻页保护上限 HALF_THRESHOLD = 0.5 # 进度过半阈值:已售/总份数 ≥ 0.5 # ---------- 调度配置 ---------- MIN_INTERVAL_SEC = 60 # 每轮跑完最小间隔(秒) MAX_INTERVAL_SEC = 90 # 每轮跑完最大间隔(秒),随机打散降风控 RUN_ALL_DAY = False # 对齐 deca:只在 [RUN_START, RUN_END] 直播窗口内轮询(可跨午夜);设 True 则全天 RUN_START = "20:30" # 运行窗口起点,同时是「本场新上架」判定的时间门槛(soldTime≥此才算新品) RUN_END = "06:00" # 运行窗口终点(跨午夜,2026/08/15 得卡由 03:00 延到 06:00) # 进程内已见在售集合:只有本进程运行后见过在售的 goods_id 才参与结束对账(重启即清空) _seen_onsale = set() logger.remove() logger.add("./logs/onsale_alert_{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") def fmt_money(raw) -> str | None: """把原始金额(×1000000)格式化为「元」字符串:整数去小数,非整去尾零。 Args: raw (int | float | None): 接口原始金额,÷1000000 为元。 Returns: str | None: 如 "199"/"12.5";raw 为空返回 None。 """ if raw is None: return None y = raw / 1000000 return str(int(y)) if y == int(y) else f"{y:.2f}".rstrip("0").rstrip(".") def fmt_price(low, high) -> str: """把最低价/最高价拼成展示价格文案:有区间显区间,否则显单价。 Args: low (int | None): 最低价原始值(×1000000)。 high (int | None): 最高价原始值(×1000000)。 Returns: str: 如 "¥199~299"/"¥199";均空返回 "¥-"。 """ lm, hm = fmt_money(low), fmt_money(high) if lm and hm and lm != hm: return f"¥{lm}~{hm}" if hm: return f"¥{hm}" if lm: return f"¥{lm}" return "¥-" def parse_dt(v) -> datetime | None: """把接口的上架时间字段解析为 datetime(兼容字符串日期与毫秒/秒级时间戳)。 Args: v (str | int | None): soldTime 原始值。 Returns: datetime | None: 解析成功返回 datetime;空值/无法解析返回 None(新品判断按保守处理,视为非新品)。 """ if not v: return None if isinstance(v, (int, float)): # 时间戳:毫秒>1e12 转秒 ts = v / 1000 if v > 1e12 else v try: return datetime.fromtimestamp(ts) except (ValueError, OSError): return None for fmt in ("%Y-%m-%d %H:%M:%S", "%Y-%m-%dT%H:%M:%S", "%Y/%m/%d %H:%M:%S", "%Y-%m-%d"): try: return datetime.strptime(str(v), fmt) except ValueError: continue return None def parse_onsale(rec: dict) -> dict: """把 index/top 一条在售记录抽取为监控所需字段字典。 Args: rec (dict): index/top 返回 records 里的一条。 Returns: dict: 含 goods_id/corp_info_id/corp_info_name/goods_name/价格/库存/规格/上架时间。 """ return { "goods_id": rec.get("goodsId"), "corp_info_id": rec.get("corpInfoId"), "corp_info_name": rec.get("corpInfoName"), "goods_name": rec.get("goodsName"), "amount": rec.get("amount"), "highest_price": rec.get("highestPrice"), "lowest_price": rec.get("lowestPrice"), "stock_amount": rec.get("stockAmount"), "residue_stock_amount": rec.get("residueStockAmount"), "specification_name": rec.get("specificationName"), "sold_time": rec.get("soldTime"), } def compute_ratio(item: dict) -> float: """计算某在售商品的售出进度比例。 Args: item (dict): parse_onsale 结果。 Returns: float: 已售/总份数;总份数缺失或 ≤0 返回 0.0。 """ total = item.get("stock_amount") residue = item.get("residue_stock_amount") if not isinstance(total, int) or total <= 0 or not isinstance(residue, int): return 0.0 sold = total - residue return sold / total if sold > 0 else 0.0 def _window_start(now: datetime) -> datetime: """本场监控窗口起点:最近一个已过去的 RUN_START(新上架判定的时间门槛)。 对齐 deca:当前已过今天 RUN_START 用今天的,否则用昨天的 RUN_START(跨午夜场)。 Args: now (datetime): 当前时间。 Returns: datetime: 本场窗口起点时刻。 """ sh, sm = _parse_hhmm(RUN_START) start = now.replace(hour=sh, minute=sm, second=0, microsecond=0) if now < start: # 今天 RUN_START 还没到 → 用昨天的 start -= timedelta(days=1) return start def is_new_arrival(item: dict, now: datetime) -> bool: """判断某商品是否为「本场新上架的新品」:soldTime ≥ 本场窗口起点(最近的 RUN_START)。 对齐 deca 最新逻辑:只把「本场监控窗口起点之后上架」的当新品;更早上架的老货静默建档、不提醒; 空值/解析失败一律视为非新品(保守,不误报老货)。 Args: item (dict): parse_onsale 结果。 now (datetime): 当前时间。 Returns: bool: soldTime ≥ 本场窗口起点返回 True;空值/解析失败/老货返回 False。 """ dt = parse_dt(item.get("sold_time")) if dt is None: return False return dt >= _window_start(now) def fetch_watch_onsale(log) -> tuple[list, bool]: """轮询全站在售(各 productType 翻全部页),过滤出关注商家的在售商品并按 goods_id 去重。 Args: log: 日志对象。 Returns: tuple[list, bool]: (关注商家在售列表, ok)。ok=False 表示取数中途失败、数据不全,本轮应整轮跳过。 """ watched, seen = [], set() watch_ids = set(WATCH_CORPS.keys()) for pt in PRODUCT_TYPES: for page in range(1, MAX_PAGES + 1): j = core.do_request(log, "/search/app/index/top", { "currentPage": str(page), "limit": str(PAGE_LIMIT), "productType": pt, "systemBusinessType": SYSTEM_BUSINESS_TYPE}) if not j: # 请求失败:数据不全,标记本轮无效,避免误判下架/漏报 log.warning(f"在售翻页失败 productType={pt} page={page},本轮取数不全") return watched, False recs = (j.get("data") or {}).get("records") or [] if not recs: break for r in recs: it = parse_onsale(r) if it["corp_info_id"] in watch_ids and it["goods_id"] not in seen: seen.add(it["goods_id"]) watched.append(it) if len(recs) < PAGE_LIMIT: # 末页 break time.sleep(0.3) return watched, True def load_existing(log, pool) -> dict: """读取关注商家在监控表里的已建档商品及其三个提醒标记位。 Args: log: 日志对象。 pool: 数据库连接池。 Returns: dict: {goods_id: {"corp_info_id","corp_info_name","goods_name","new","half","ended"}}。 """ ids = list(WATCH_CORPS.keys()) ph = ",".join(["%s"] * len(ids)) rows = pool.select_all( f"SELECT goods_id, corp_info_id, corp_info_name, goods_name, new_notified, half_notified, ended_notified " f"FROM jw_onsale_alert_record WHERE corp_info_id IN ({ph})", tuple(ids)) return {r[0]: {"corp_info_id": r[1], "corp_info_name": r[2], "goods_name": r[3], "new": r[4], "half": r[5], "ended": r[6]} for r in rows} def archive_item(pool, item: dict, new_notified: int) -> None: """把新出现在监控里的商品建档到 jw_onsale_alert_record(已存在则只刷新名称,不动标记位)。 Args: pool: 数据库连接池。 item (dict): parse_onsale 结果。 new_notified (int): 建档时的新品提醒标记;1=老货静默建档不发,0=待发新品提醒。 """ pool.update_one( "INSERT INTO jw_onsale_alert_record (goods_id, corp_info_id, corp_info_name, goods_name, new_notified) " "VALUES (%s,%s,%s,%s,%s) " "ON DUPLICATE KEY UPDATE corp_info_name=VALUES(corp_info_name), goods_name=VALUES(goods_name)", (item["goods_id"], item["corp_info_id"], item["corp_info_name"], item["goods_name"], new_notified)) def mark_flag(pool, goods_id: int, column: str) -> None: """把某商品的某个提醒标记位置 1。 Args: pool: 数据库连接池。 goods_id (int): 商品ID。 column (str): 标记列名,取值 new_notified / half_notified / ended_notified。 """ pool.update_one(f"UPDATE jw_onsale_alert_record SET {column}=1 WHERE goods_id=%s", (goods_id,)) def confirm_ended(log, goods_id: int) -> dict | None: """二次确认某商品是否真的结束:打详情,售罄(surplusStock≤0) 或 已过下架时间(offShelfTime) 才算结束。 对齐 deca 最新逻辑:结束 = 剩余库存售罄 OR 当前已过销售结束/下架时间(offShelfTime),二者任一命中即结束; 详情取不到 / 两条件都不满足 → 视为本轮抖动缺席,返回 None、下轮再看(不误判下架)。 Args: log: 日志对象。 goods_id (int): 商品ID。 Returns: dict | None: 确认结束返回详情 data(供拼战报);仍在售/详情取不到返回 None。 """ j = core.do_request(log, "/search/app/merchantGoodsId", {"goodsId": str(goods_id), "systemBusinessType": SYSTEM_BUSINESS_TYPE}) if not j: return None data = j.get("data") or {} residue = data.get("surplusStock") # 详情用 surplusStock(剩余库存),非 index/top 的 residueStockAmount if isinstance(residue, int) and residue <= 0: # 售罄 return data end_dt = parse_dt(data.get("offShelfTime")) # 下架/销售结束时间 if end_dt is not None and datetime.now() >= end_dt: # 已过下架时间 return data return None def send_new(log, corp_id: int, items: list) -> bool: """发送某商家的新品上架提醒(一条 markdown,多个新品逐条列出)。 Args: log: 日志对象。 corp_id (int): 商家ID。 items (list): 待提醒的新品 item 列表。 Returns: bool: 发送成功返回 True(供 send-then-mark 置位)。 """ corp_name = WATCH_CORPS.get(corp_id, items[0].get("corp_info_name") or str(corp_id)) lines = [] for it in items: price = fmt_price(it.get("lowest_price"), it.get("highest_price")) stock = it.get("stock_amount") residue = it.get("residue_stock_amount") lines.append(f"**{it.get('goods_name')}**\n💰 {price} | 📦 {stock}份 | 🎯 余{residue}/{stock}") title = f"🆕 集物星球 · {corp_name} 新商品上架 {len(items)} 个" return wx.send_wechat_group_msg(log=log, items=lines, title=title) is not None def send_half(log, corp_id: int, items: list) -> bool: """发送某商家的进度过半提醒(一条 markdown,多个过半商品逐条列出)。 Args: log: 日志对象。 corp_id (int): 商家ID。 items (list): 元素为 (item, ratio) 的列表。 Returns: bool: 发送成功返回 True。 """ corp_name = WATCH_CORPS.get(corp_id, items[0][0].get("corp_info_name") or str(corp_id)) lines = [] for it, ratio in items: price = fmt_price(it.get("lowest_price"), it.get("highest_price")) stock = it.get("stock_amount") residue = it.get("residue_stock_amount") lines.append(f"**{it.get('goods_name')}**\n💰 {price} | 📈 进度{int(ratio * 100)}% | 🎯 余{residue}/{stock}") title = f"🔥 集物星球 · {corp_name} 拼团进度过半 {len(items)} 个" return wx.send_wechat_group_msg(log=log, items=lines, title=title) is not None def get_buyer_count(log, goods_id: int) -> int: """取某商品的购买人数(赠品公示·玩家维度接口的 totalCount,需登录)。 Args: log: 日志对象。 goods_id (int): 商品ID。 Returns: int: 购买人数;取不到返回 0。 """ j = core.do_request(log, "/order/merchant/app/query/gift/publicity/user/group/pager", {"currentPage": "1", "giftBusinessName": "", "goodsId": str(goods_id), "limit": "1", "systemBusinessType": SYSTEM_BUSINESS_TYPE}, need_auth=True) if not j: return 0 return (j.get("data") or {}).get("totalCount") or 0 def send_ended(log, goods_id: int, ex: dict, detail: dict) -> bool: """发送某商品的一车结束战报(标题/价格/售出件数/购买人数)。 Args: log: 日志对象。 goods_id (int): 商品ID。 ex (dict): 库内已建档信息(商家名、商品名兜底)。 detail (dict): confirm_ended 返回的商品详情 data(详情字段:price/highestPrice/stock)。 Returns: bool: 发送成功返回 True(供 send-then-mark 置位)。 """ corp_name = WATCH_CORPS.get(ex.get("corp_info_id"), ex.get("corp_info_name") or "") goods_name = detail.get("goodsName") or ex.get("goods_name") price = fmt_price(detail.get("price"), detail.get("highestPrice")) # 详情用 price/highestPrice total = detail.get("stock") # 详情总份数=stock;已售罄即全部售出 buyers = get_buyer_count(log, goods_id) # 玩家维度准确购买人数 line = f"**{goods_name}**\n💰 {price} | 🎯 售出 {total} 件\n👥 {buyers} 人购买" title = f"🏁 集物星球 · {corp_name} 一车结束" return wx.send_wechat_group_msg(log=log, items=[line], title=title) is not None def detect_and_report_ended(log, pool, existing: dict, onsale_codes: set) -> None: """结束对账:本进程见过在售、本轮已消失、且未播报过的商品,二次确认后发结束战报。 Args: log: 日志对象。 pool: 数据库连接池。 existing (dict): load_existing 结果(本轮读库快照)。 onsale_codes (set): 本轮抓到的关注商家在售 goods_id 集合。 """ candidates = [gid for gid, ex in existing.items() if gid in _seen_onsale and gid not in onsale_codes and ex["ended"] == 0] for gid in candidates: detail = confirm_ended(log, gid) if detail is None: # 未确认结束(可能只是本轮抖动缺席),下轮再看 continue if send_ended(log, gid, existing[gid], detail): mark_flag(pool, gid, "ended_notified") log.success(f"结束战报已发 goods_id={gid}") def run_once(log, pool) -> None: """跑一轮监控:取数→结束对账→遍历在售判新品/过半→建档并推送。 Args: log: 日志对象。 pool: 数据库连接池。 """ items, ok = fetch_watch_onsale(log) if not ok: log.warning("本轮取数不完整,跳过判断") return onsale_codes = {it["goods_id"] for it in items} log.info(f"本轮关注商家在售 {len(items)} 个") existing = load_existing(log, pool) # 1) 先做结束对账(用「旧 existing + 本轮在售」判断谁消失了) detect_and_report_ended(log, pool, existing, onsale_codes) # 2) 标记本进程已见过这些在售(此后它们消失才纳入结束对账) _seen_onsale.update(onsale_codes) # 3) 遍历在售,收集新品/过半 now = datetime.now() new_alerts, half_alerts = [], [] for it in items: gid = it["goods_id"] ex = existing.get(gid) if ex is None: # 监控里首次出现的商品:建档 new = is_new_arrival(it, now) archive_item(pool, it, new_notified=0 if new else 1) ex = {"corp_info_id": it["corp_info_id"], "new": 0 if new else 1, "half": 0, "ended": 0} existing[gid] = ex if new: new_alerts.append(it) elif ex["new"] == 0 and ex["ended"] == 0: # 上轮发送失败的新品,补发 new_alerts.append(it) # 过半判断(未过半、未结束才评估;首次达标即发) if ex["half"] == 0 and ex["ended"] == 0: ratio = compute_ratio(it) if ratio >= HALF_THRESHOLD: mark_flag(pool, gid, "half_notified") # 先置位再发 half_alerts.append((it, ratio)) # 4) 推送:新品按商家分组,send-then-mark(发成功才置 new_notified) if new_alerts: groups = defaultdict(list) for it in new_alerts: groups[it["corp_info_id"]].append(it) for cid, its in groups.items(): if send_new(log, cid, its): for it in its: mark_flag(pool, it["goods_id"], "new_notified") log.info(f"新品上架提醒 {len(new_alerts)} 个") # 5) 推送:过半按商家分组(已先置位,直接发) if half_alerts: hgroups = defaultdict(list) for it, ratio in half_alerts: hgroups[it["corp_info_id"]].append((it, ratio)) for cid, its in hgroups.items(): send_half(log, cid, its) log.info(f"进度过半提醒 {len(half_alerts)} 个") @retry(stop=stop_after_attempt(100), wait=wait_fixed(600), after=core.after_log) def main_task(log) -> None: """监控主流程(挂了每 10 分钟重试,最多 100 次):按运行窗口手写轮询循环。 Args: log: 日志对象。 Raises: RuntimeError: 数据库连接池异常时抛出以触发重试。 """ log.info("启动在售监控" + "." * 40) pool = MySQLConnectionPool(log=log) if not pool.check_pool_health(): log.error("数据库连接池异常") raise RuntimeError("数据库连接池异常") while True: now = datetime.now() if not _in_run_window(now): sleep_s = _seconds_to_next_window(now) log.info(f"当前不在运行窗口,休眠 {sleep_s} 秒到下次窗口") time.sleep(sleep_s) continue try: run_once(log, pool) except Exception as e: log.error(f"监控单轮异常: {e}") gap = random.randint(MIN_INTERVAL_SEC, MAX_INTERVAL_SEC) log.info(f"本轮结束,{gap} 秒后再轮询") time.sleep(gap) def _parse_hhmm(s: str) -> tuple[int, int]: """把 "HH:MM" 解析为 (时, 分)。 Args: s (str): 形如 "20:30" 的时间串。 Returns: tuple[int, int]: (小时, 分钟)。 """ h, m = s.split(":") return int(h), int(m) def _in_run_window(now: datetime) -> bool: """判断当前是否在运行窗口内(支持跨午夜窗口)。 Args: now (datetime): 当前时间。 Returns: bool: RUN_ALL_DAY 恒 True;否则按 [RUN_START, RUN_END] 判断。 """ if RUN_ALL_DAY: return True sh, sm = _parse_hhmm(RUN_START) eh, em = _parse_hhmm(RUN_END) start = now.replace(hour=sh, minute=sm, second=0, microsecond=0) end = now.replace(hour=eh, minute=em, second=0, microsecond=0) if start <= end: # 同日窗口 return start <= now <= end return now >= start or now <= end # 跨午夜:[start,次日end] def _seconds_to_next_window(now: datetime) -> int: """计算距下一次进入运行窗口还有多少秒。 Args: now (datetime): 当前时间。 Returns: int: 休眠秒数(至少 1)。 """ sh, sm = _parse_hhmm(RUN_START) start = now.replace(hour=sh, minute=sm, second=0, microsecond=0) if now >= start: # 今天窗口起点已过,等明天 start += timedelta(days=1) return max(1, int((start - now).total_seconds())) def _parse_start_time(text: str) -> str: """校验并规范化命令行传入的窗口起点时间,返回 "HH:MM" 字符串。 Args: text (str): 形如 "20:30" 或 "20:30:00" 的时间串。 Returns: str: 规范化的 "HH:MM"(丢掉秒,供 _parse_hhmm 使用)。 Raises: argparse.ArgumentTypeError: 格式非法(非 HH:MM / HH:MM:SS)时抛出。 """ text = text.strip() for fmt in ("%H:%M:%S", "%H:%M"): try: return datetime.strptime(text, fmt).strftime("%H:%M") except ValueError: continue raise argparse.ArgumentTypeError(f"开始时间格式非法:{text!r},应为 HH:MM 或 HH:MM:SS,如 20:30") def _parse_args() -> argparse.Namespace: """解析命令行参数:可指定运行窗口起点(同时作为「新上架」判定门槛)。 对齐 deca:位置参数与 --start 等价、位置优先;不传则用默认 RUN_START。 Returns: argparse.Namespace: 含 start(位置) / start_opt(--start) 两个可选时间(均为规范化 "HH:MM" 或 None)。 """ parser = argparse.ArgumentParser( description="集物星球在售提醒:可指定运行窗口起点(该时间后的新上架才提醒)") parser.add_argument("start", nargs="?", type=_parse_start_time, default=None, help="运行窗口起点 HH:MM[:SS],默认 20:30;位置参数写法,如 17:00") parser.add_argument("--start", dest="start_opt", type=_parse_start_time, default=None, help="运行窗口起点 HH:MM[:SS],与位置参数等价,如 --start 17:00") return parser.parse_args() def schedule_task(): """监控入口:直接进入 main_task 的手写轮询循环。""" main_task(log=logger) if __name__ == "__main__": _args = _parse_args() _start = _args.start or _args.start_opt # 位置参数优先,其次 --start;都未传则保持默认 RUN_START if _start is not None: RUN_START = _start # 覆盖模块级默认;窗口判定与「新上架」门槛均随之改变(两处均裸读该全局) logger.info(f"运行窗口起点由命令行指定为 {RUN_START}") schedule_task()