# -*- coding: utf-8 -*- # Author : Charley # Python : 3.12.10 # Date : 2026/08/19 """集物星球已售抓取:商家列表→逐商家已售 corp/history,并对全站已售商品抓拆卡报告与购买记录。 接口对应(以抓包为准): - 商家列表 hotRecommend(需登录)→ jw_shop_record - 已售历史 corp/history(免登录)→ jw_sold_product_record(只增) - 拆卡报告 /goods/gift/report/query/search/pager(需登录, sbt=6)→ jw_report_record(只增, giftReportId 去重) - 购买记录 /order/merchant/app/query/gift/publicity/user/group/pager(赠品公示·玩家维度, 需登录)→ jw_player_record 每条=一个买家 userId/userNick/picId/count,翻页拿全量;**不 DB 去重**,靠 jw_sold_product_record.buy_fetched 状态位控制每商品抓一次(翻页取全→批量入库成功→再置 buy_fetched=1;失败不置位、下轮重试)。 拆卡报告「更新不及时」:对结束 REPORT_REFETCH_DAYS 天内(或未结束)的已售商品每轮重查(INSERT IGNORE 自然补齐), 老车只补抓从未抓过的。 """ import sys import time import json import schedule from loguru import logger from tenacity import retry, stop_after_attempt, wait_fixed from mysql_pool import MySQLConnectionPool import jiwu_core as core import jw_detail TARGET_CORPS = [] # 空=全站;如 [100716, 100715] 只抓 Jake/九叔 SYSTEM_BUSINESS_TYPE = 5 PAGE_LIMIT = 20 MAX_SHOP_PAGES = 30 # 商家列表翻页上限 MAX_SOLD_PAGES = 300 # 单商家已售翻页上限 REPORT_PAGE_LIMIT = 20 # 拆卡报告翻页每页 MAX_REPORT_PAGES = 50 # 单商品拆卡报告翻页上限 REPORT_REFETCH_DAYS = 3 # 拆卡报告:结束 N 天内每轮重查(补迟到更新),老车只补抓未抓过的 BUY_PAGE_LIMIT = 20 # 购买记录翻页每页 MAX_BUY_PAGES = 50 # 单商品购买记录翻页上限 logger.remove() logger.add("./logs/sold_{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") _SOLD_COLS = ["goods_id", "program_id", "goods_name", "product_type", "business_type", "goods_ip_id", "goods_ip_name", "random_type", "corp_info_id", "corp_info_name", "amount", "highest_price", "lowest_price", "stock_amount", "residue_stock_amount", "specification_name", "report_status", "live_replay_url", "sold_time", "soldout_time", "finish_time"] def parse_shop(rec: dict) -> dict: """把 hotRecommend 一条商家记录规范化为入库字典。 hotRecommend 商家字段实测:`corpInfoId`(商家id) / `corpInfoName`(名) / `number`(**粉丝数**) / `saleNum`(**在售商品数**)。响应里的 `userId` 是查询者本人账号id(每行恒等)、非商家id,不入库; 好评/新品/简介接口体系无此字段,不入库。 Args: rec (dict): 商家列表返回的一条。 Returns: dict: 与 jw_shop_record 列对齐的字典(corp_info_id/corp_name/fans_amount/onsale_amount)。 """ return { "corp_info_id": rec.get("corpInfoId"), "corp_name": rec.get("corpInfoName"), "fans_amount": rec.get("number"), # number = 粉丝数(实测确认) "onsale_amount": rec.get("saleNum"), # saleNum = 在售商品数 } LIVE_BASE_URL = "https://play.jiwustar.com" # 直播回放域名(BASE_PULL_STREAM_URL, 逆向 EnvironmentManager RELEASE) def full_live_url(path) -> str | None: """把 liveReplayUrl 相对路径(形如 /live/xxx.mp4)拼成完整直播回放链接。 接口返回的 liveReplayUrl 只含 `/live/...`,需拼直播域名 https://play.jiwustar.com(路径已带 /live/)。 Args: path (str | None): 接口返回的 liveReplayUrl。 Returns: str | None: 完整 URL;空值返回 None;已是 http(s) 开头则原样返回。 """ if not path: return None if str(path).startswith("http"): return path return LIVE_BASE_URL + (path if str(path).startswith("/") else "/" + path) def parse_sold(rec: dict) -> dict: """把 corp/history 一条已售记录规范化为入库字典。 Args: rec (dict): 已售历史返回的一条。 Returns: dict: 与 jw_sold_product_record 列对齐的字典。 """ return { "goods_id": rec.get("goodsId"), "program_id": rec.get("programId"), "goods_name": rec.get("goodsName"), "product_type": rec.get("productType"), "business_type": rec.get("businessType"), "goods_ip_id": rec.get("goodsIPId"), "goods_ip_name": rec.get("goodsIPName"), "random_type": rec.get("randomType"), "corp_info_id": rec.get("corpInfoId"), "corp_info_name": rec.get("corpInfoName"), "amount": core.to_yuan(rec.get("amount")), "highest_price": core.to_yuan(rec.get("highestPrice")), "lowest_price": core.to_yuan(rec.get("lowestPrice")), "stock_amount": rec.get("stockAmount"), "residue_stock_amount": rec.get("residueStockAmount"), "specification_name": rec.get("specificationName"), "report_status": rec.get("reportStatus"), "live_replay_url": full_live_url(rec.get("liveReplayUrl")), "sold_time": rec.get("soldTime"), "soldout_time": rec.get("soldOutTime"), "finish_time": rec.get("finishTime"), } def get_shops(log, pool) -> list: """抓商家列表 hotRecommend 落库 jw_shop_record,返回 corp_info_id 列表。 Args: log: 日志对象。 pool: 数据库连接池。 Returns: list: 商家 corp_info_id 列表(受 TARGET_CORPS 过滤)。 """ corp_ids = [] for page in range(1, MAX_SHOP_PAGES + 1): j = core.do_request(log, "/search/app/index/corp/hotRecommend", {"currentPage": str(page), "limit": str(PAGE_LIMIT), "systemBusinessType": 6}, need_auth=True) # 商家列表实测需登录 if not j: break data = j.get("data") or {} recs = data.get("records") or [] if not recs: break shops = [parse_shop(r) for r in recs] for s in shops: pool.update_one( "INSERT INTO jw_shop_record (corp_info_id,corp_name,fans_amount,onsale_amount) " "VALUES (%s,%s,%s,%s) " "ON DUPLICATE KEY UPDATE corp_name=VALUES(corp_name)," "fans_amount=VALUES(fans_amount),onsale_amount=VALUES(onsale_amount)", (s["corp_info_id"], s["corp_name"], s["fans_amount"], s["onsale_amount"])) corp_ids.append(s["corp_info_id"]) # 末页判定用响应里的实际每页大小 data.limit:hotRecommend 服务端把每页压到 10 条(≠请求的 PAGE_LIMIT), # 若用 PAGE_LIMIT 判会第 1 页就误判停(之前只翻到 10 个商家的 bug)。实测全量约 82 个商家、9 页。 srv_limit = data.get("limit") or PAGE_LIMIT if len(recs) < srv_limit: break time.sleep(0.3) if TARGET_CORPS: corp_ids = [c for c in corp_ids if c in TARGET_CORPS] log.info(f"商家列表入库完成,待抓已售商家 {len(corp_ids)} 家") return corp_ids def fetch_sold_for_corp(log, pool, corp_id: int) -> list: """翻页抓某商家的已售历史,INSERT IGNORE 只增落库。 Args: log: 日志对象。 pool: 数据库连接池。 corp_id (int): 商家 corpInfoId。 Returns: list: 本商家本轮抓到的已售商品 goods_id 列表(供抓拆卡报告/购买记录)。 """ total, goods_ids = 0, [] for page in range(1, MAX_SOLD_PAGES + 1): j = core.do_request(log, "/search/app/corp/history", { "corpInfoId": str(corp_id), "currentPage": str(page), "limit": str(PAGE_LIMIT), "systemBusinessType": SYSTEM_BUSINESS_TYPE}) # 免登录 if not j: break recs = (j.get("data") or {}).get("records") or [] if not recs: break rows = [parse_sold(r) for r in recs] pool.insert_many(table="jw_sold_product_record", data_list=rows, ignore=True) goods_ids.extend(r["goods_id"] for r in rows if r.get("goods_id")) total += len(rows) if len(recs) < PAGE_LIMIT: break time.sleep(0.3) log.info(f"商家 {corp_id} 已售抓取 {total} 条") return goods_ids # ---------------- 拆卡报告(需登录, sbt=6, INSERT IGNORE 去重, 结束N天内重查) ---------------- def parse_report(rec: dict, goods_id: int) -> dict: """把拆卡报告一条记录规范化为入库字典(goods_id 由调用方补,响应体不含)。 Args: rec (dict): gift/report/query/search/pager 返回 records 里的一条。 goods_id (int): 所属商品ID(查询参数)。 Returns: dict: 与 jw_report_record 列对齐的字典。 """ res = rec.get("resAddrList") return { "gift_report_id": rec.get("giftReportId"), "goods_id": goods_id, "user_id": rec.get("userId"), "username": rec.get("username"), "corp_id": rec.get("corpId"), "corp_name": rec.get("corpName"), "goods_name": rec.get("goodsName"), "serial_item_name": rec.get("serialItemName"), "res_addr_list": json.dumps(res, ensure_ascii=False) if res is not None else None, "winner_status": rec.get("winnerStatus"), "anonymous_status": rec.get("anonymousStatus"), "my_gift_status": rec.get("myGiftStatus"), "report_time": rec.get("createTime"), } def fetch_reports_for_goods(log, pool, goods_id: int) -> None: """翻页抓某商品的拆卡报告(需登录),INSERT IGNORE 按 giftReportId 只增落库。 Args: log: 日志对象。 pool: 数据库连接池。 goods_id (int): 商品 goodsId。 """ total = 0 for page in range(1, MAX_REPORT_PAGES + 1): j = core.do_request(log, "/goods/gift/report/query/search/pager", { "currentPage": str(page), "goodsId": str(goods_id), "limit": str(REPORT_PAGE_LIMIT), "systemBusinessType": 6}, need_auth=True) # 拆卡报告实测需登录, sbt=6 if not j: break recs = (j.get("data") or {}).get("records") or [] if not recs: break rows = [parse_report(r, goods_id) for r in recs] pool.insert_many(table="jw_report_record", data_list=rows, ignore=True) total += len(rows) if len(recs) < REPORT_PAGE_LIMIT: break time.sleep(0.2) if total: log.info(f"商品 {goods_id} 拆卡报告 {total} 条") def fetch_reports_for_sold(log, pool, goods_ids: list) -> None: """对已售商品抓拆卡报告:结束 REPORT_REFETCH_DAYS 天内(或未结束)每轮重查,老车只补抓从未抓过的。 拆卡报告有唯一 giftReportId,重查靠 INSERT IGNORE 去重,故重查安全、能补齐迟到的报告。 Args: log: 日志对象。 pool: 数据库连接池。 goods_ids (list): 本商家本轮的已售商品 goods_id 列表。 """ ids = list({g for g in goods_ids if g}) if not ids: return ph = ",".join(["%s"] * len(ids)) have = {r[0] for r in pool.select_all( f"SELECT DISTINCT goods_id FROM jw_report_record WHERE goods_id IN ({ph})", tuple(ids))} recent = {r[0] for r in pool.select_all( f"SELECT goods_id FROM jw_sold_product_record WHERE goods_id IN ({ph}) " "AND (finish_time IS NULL OR finish_time >= NOW() - INTERVAL %s DAY)", tuple(ids) + (REPORT_REFETCH_DAYS,))} todo = recent | (set(ids) - have) # 近N天(或未结束)重查 + 从未抓过补一次 for gid in todo: try: fetch_reports_for_goods(log, pool, gid) except Exception as e: log.error(f"商品 {gid} 拆卡报告抓取异常: {e}") # ---------------- 购买记录(赠品公示·玩家维度, 需登录, 翻页拿全, buy_fetched 状态位控制抓一次) ---------------- def parse_buy(rec: dict, goods_id: int) -> dict: """把玩家维度购买记录一条买家规范化为入库字典。 Args: rec (dict): publicity/user/group/pager 返回 records 里的一条。 goods_id (int): 所属商品ID(兜底,响应体一般也带 goodsId)。 Returns: dict: 与 jw_player_record 列对齐的字典(goods_id/user_id/user_nick/pic_id/buy_count)。 """ return { "goods_id": rec.get("goodsId") or goods_id, "user_id": rec.get("userId"), "user_nick": rec.get("userNick"), "pic_id": rec.get("picId"), "buy_count": rec.get("count"), } def fetch_buy_for_goods(log, pool, goods_id: int) -> bool: """翻页取全某商品的购买记录(玩家维度,需登录)→ 批量入库(不去重)。 先把所有页收齐再一次性入库;只要有一页请求失败就判为未取全、返回 False(不入库、不置状态位、下轮重试)。 Args: log: 日志对象。 pool: 数据库连接池。 goods_id (int): 商品 goodsId。 Returns: bool: 取全并入库成功返回 True(供上层置 buy_fetched=1);请求失败返回 False。 """ rows = [] for page in range(1, MAX_BUY_PAGES + 1): j = core.do_request(log, "/order/merchant/app/query/gift/publicity/user/group/pager", { "currentPage": str(page), "giftBusinessName": "", "goodsId": str(goods_id), "limit": str(BUY_PAGE_LIMIT), "systemBusinessType": SYSTEM_BUSINESS_TYPE}, need_auth=True) # 需登录 if not j: return False # 请求失败=未取全,返回 False 让上层不置状态位、下轮重试 recs = (j.get("data") or {}).get("records") or [] if not recs: break rows.extend(parse_buy(r, goods_id) for r in recs if r.get("userId") is not None) if len(recs) < BUY_PAGE_LIMIT: break time.sleep(0.2) if rows: pool.insert_many(table="jw_player_record", data_list=rows, ignore=False) # 不去重,批量保存 log.info(f"商品 {goods_id} 购买记录 {len(rows)} 条") return True def fetch_buy_for_pending(log, pool, goods_ids: list) -> None: """对未抓过购买记录(buy_fetched=0)的已售商品抓全入库,成功后置 buy_fetched=1。 买家列表售罄后即固定,用已售表状态位控制「每个商品抓一次」;取全并入库成功才置位,失败下轮重试。 Args: log: 日志对象。 pool: 数据库连接池。 goods_ids (list): 本商家本轮的已售商品 goods_id 列表。 """ ids = list({g for g in goods_ids if g}) if not ids: return ph = ",".join(["%s"] * len(ids)) pending = [r[0] for r in pool.select_all( f"SELECT goods_id FROM jw_sold_product_record WHERE goods_id IN ({ph}) AND buy_fetched=0", tuple(ids))] for gid in pending: try: if fetch_buy_for_goods(log, pool, gid): pool.update_one("UPDATE jw_sold_product_record SET buy_fetched=1 WHERE goods_id=%s", (gid,)) except Exception as e: log.error(f"商品 {gid} 购买记录抓取异常: {e}") @retry(stop=stop_after_attempt(100), wait=wait_fixed(3600), after=core.after_log) def main_task(log) -> None: """已售抓取主流程(挂了每小时重试):已售历史 + 拆卡报告 + 购买记录。 Args: log: 日志对象。 Raises: RuntimeError: 数据库连接池异常时抛出以触发重试。 """ log.info("开始已售抓取" + "." * 40) pool = MySQLConnectionPool(log=log) if not pool.check_pool_health(): log.error("数据库连接池异常") raise RuntimeError("数据库连接池异常") try: get_shops(log, pool) # 先刷新商家表(hotRecommend 每日轮换, INSERT/UPDATE 累积) # 已售对库中【全量商家】循环查询——hotRecommend 每天只返回轮换的一批热门商家, # jw_shop_record 随天数累积覆盖更全;故从表里取全量(而非只取本轮 get_shops 返回的)。 corp_ids = [r[0] for r in pool.select_all("SELECT corp_info_id FROM jw_shop_record")] if TARGET_CORPS: # 如只想抓指定商家(如 Jake/九叔)在此过滤 corp_ids = [c for c in corp_ids if c in TARGET_CORPS] log.info(f"待抓已售商家 {len(corp_ids)} 家(取自 jw_shop_record 全量)") for cid in corp_ids: try: goods_ids = fetch_sold_for_corp(log, pool, cid) # 已售商品(免登录) jw_detail.enrich_detail(log, pool, "jw_sold_product_record", goods_ids) # 详情补全(免登录) fetch_reports_for_sold(log, pool, goods_ids) # 拆卡报告(需登录, 按状态重查) fetch_buy_for_pending(log, pool, goods_ids) # 购买记录(需登录, buy_fetched 状态位控制) except Exception as e: log.error(f"商家 {cid} 已售/拆卡报告/购买记录抓取异常: {e}") except Exception as e: log.error(f"已售抓取异常: {e}") finally: log.info("已售抓取结束,等待下一轮" + "." * 20) def schedule_task(): """定时入口:每天 08:00 抓一次已售(含拆卡报告、购买记录)。 08:00 配合已售日报的业务日窗口 [昨17:00, 今06:00]:窗口 06:00 关闭后再采, 报告 09:10 前全站已售数据齐全,不漏窗口尾部(00:30~06:00)结束的团。 """ # main_task(log=logger) # 立即跑一次(调试时取消注释) schedule.every().day.at("08:00").do(main_task, log=logger) while True: schedule.run_pending() time.sleep(1) if __name__ == "__main__": schedule_task()