# -*- coding: utf-8 -*- # Author : Charley # Python : 3.12.10 # Date : 2026/08/19 """集物星球在售抓取:免登录遍历各商品类型翻页拉 index/top,落库 jw_onsale_product_record。 每天 09/15/20/01 四档各跑一次(对齐 deca:上午/下午/晚上/凌晨场),每档同时把当刻在售写入每日快照表 jw_onsale_daily_record(商品+日期唯一,同日多档刷成最新值),供在售日报按天差分出「在售趋势」。 抓完做下架对账:本轮未出现的商品置 is_on_sale=0。index/top 实测免登录,故 do_request 用默认 need_auth=False。 """ import sys import time from datetime import date 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 # 商品类型:1福袋 2变风盒 3错版卡 4原盒;均属业务线 5(集卡) PRODUCT_TYPES = ["1", "2", "3", "4"] SYSTEM_BUSINESS_TYPE = 5 PAGE_LIMIT = 20 MAX_PAGES = 200 # 单类型翻页保护上限 ENRICH_DETAIL = True # 是否逐商品补拉详情字段(免登录);量大可关 logger.remove() logger.add("./logs/onsale_{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") # 入库列(与 jw_onsale_product_record 对齐) _COLS = ["goods_id", "program_id", "goods_name", "product_type", "business_type", "block_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", "sold_count", "specification_name", "report_status", "live_status", "sold_time", "product_pic", "is_on_sale"] def parse_product(rec: dict) -> dict: """把接口一条在售记录规范化为入库字典。 Args: rec (dict): index/top 返回 records 里的一条。 Returns: dict: 键与 jw_onsale_product_record 列对齐的字典。 """ stock = rec.get("stockAmount") residue = rec.get("residueStockAmount") sold = (stock - residue) if (isinstance(stock, int) and isinstance(residue, int)) else None 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"), "block_type": rec.get("blockType"), "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": stock, "residue_stock_amount": residue, "sold_count": sold, "specification_name": rec.get("specificationName"), "report_status": rec.get("reportStatus"), "live_status": rec.get("liveStatus"), "sold_time": rec.get("soldTime"), "product_pic": rec.get("productPic"), "is_on_sale": 1, } def upsert_products(log, pool, rows: list) -> None: """把在售商品批量 upsert 到库(存在则更新库存/价格等最新状态)。 Args: log: 日志对象。 pool: 数据库连接池。 rows (list): parse_product 结果列表。 """ if not rows: return cols_sql = ",".join(f"`{c}`" for c in _COLS) ph = ",".join(["%s"] * len(_COLS)) upd = ",".join(f"`{c}`=VALUES(`{c}`)" for c in _COLS if c != "goods_id") sql = f"INSERT INTO jw_onsale_product_record ({cols_sql}) VALUES ({ph}) ON DUPLICATE KEY UPDATE {upd}" args_list = [tuple(r[c] for c in _COLS) for r in rows] pool.insert_many(query=sql, args_list=args_list) log.info(f"upsert 在售商品 {len(rows)} 条") # 每日快照入库列(与 jw_onsale_daily_record 对齐);唯一键命中(同日多档)时只刷 _SNAP_UPD 里的量价状态 _SNAP_COLS = ["goods_id", "corp_info_id", "corp_info_name", "snapshot_date", "product_type", "stock_amount", "sold_count", "residue_stock_amount", "live_status", "amount"] _SNAP_UPD = ["stock_amount", "sold_count", "residue_stock_amount", "live_status", "amount"] def upsert_daily_snapshot(log, pool, rows: list) -> None: """把本页在售商品写入每日快照表(goods_id+snapshot_date 唯一,当天多档跑刷成最新值)。 对齐参考项目 deca 的每日快照:唯一键命中(同日多档)时只更新量价状态、保留商品/商家/日期, 保证每商品每天一行、恒为当天最后一次采集值,供在售日报按天差分出「在售趋势」。 Args: log: 日志对象。 pool: 数据库连接池。 rows (list): parse_product 结果列表(含 amount 等,金额已换算为元)。 """ if not rows: return today = date.today().isoformat() # 快照日期;01:00 凌晨场归入新自然日(与 deca 一致) cols_sql = ",".join(f"`{c}`" for c in _SNAP_COLS) ph = ",".join(["%s"] * len(_SNAP_COLS)) upd = ",".join(f"`{c}`=VALUES(`{c}`)" for c in _SNAP_UPD) sql = (f"INSERT INTO jw_onsale_daily_record ({cols_sql}) VALUES ({ph}) " f"ON DUPLICATE KEY UPDATE {upd}") args_list = [(r["goods_id"], r["corp_info_id"], r["corp_info_name"], today, r["product_type"], r["stock_amount"], r["sold_count"], r["residue_stock_amount"], r["live_status"], r["amount"]) for r in rows] pool.insert_many(query=sql, args_list=args_list) def sweep_offsale(log, pool, seen_ids: set) -> None: """下架对账:库中标记在售、但本轮未出现的商品置 is_on_sale=0。 Args: log: 日志对象。 pool: 数据库连接池。 seen_ids (set): 本轮抓到的全部 goods_id。 """ rows = pool.select_all("SELECT goods_id FROM jw_onsale_product_record WHERE is_on_sale=1") db_ids = {r[0] for r in rows} gone = db_ids - seen_ids for gid in gone: pool.update_one("UPDATE jw_onsale_product_record SET is_on_sale=0 WHERE goods_id=%s", (gid,)) if gone: log.info(f"下架对账:{len(gone)} 个商品置为已下架") def fetch_onsale(log, pool) -> None: """免登录遍历各商品类型翻页抓在售并落库,最后做下架对账。 Args: log: 日志对象。 pool: 数据库连接池。 """ seen_ids = set() 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: break recs = (j.get("data") or {}).get("records") or [] if not recs: break rows = [parse_product(r) for r in recs] seen_ids.update(r["goods_id"] for r in rows) upsert_products(log, pool, rows) upsert_daily_snapshot(log, pool, rows) # 同步写每日快照(供在售趋势差分) if len(recs) < PAGE_LIMIT: # 末页 break time.sleep(0.3) sweep_offsale(log, pool, seen_ids) if ENRICH_DETAIL: # 逐商品补拉详情字段(免登录, detail_fetched 控制每商品补一次) jw_detail.enrich_detail(log, pool, "jw_onsale_product_record", list(seen_ids)) log.success(f"本轮在售抓取完成,共 {len(seen_ids)} 个在售商品") @retry(stop=stop_after_attempt(100), wait=wait_fixed(3600), after=core.after_log) def main_task(log) -> None: """在售抓取主流程(挂了每小时重试,最多100次)。 Args: log: 日志对象。 Raises: RuntimeError: 数据库连接池异常时抛出以触发重试。 """ log.info("开始在售抓取" + "." * 40) pool = MySQLConnectionPool(log=log) if not pool.check_pool_health(): log.error("数据库连接池异常") raise RuntimeError("数据库连接池异常") try: fetch_onsale(log, pool) except Exception as e: log.error(f"在售抓取异常: {e}") finally: log.info("在售抓取结束,等待下一轮" + "." * 20) def schedule_task(): """定时入口:每天 09:00 / 15:00 / 20:00 / 01:00 各抓一次在售(对齐 deca 四档:上午/下午/晚上/凌晨场)。""" # main_task(log=logger) # 调试时取消注释立即跑一次 for _hhmm in ("09:00", "15:00", "20:00", "01:00"): schedule.every().day.at(_hhmm).do(main_task, log=logger) while True: schedule.run_pending() time.sleep(1) if __name__ == "__main__": schedule_task()