jw_onsale_spider.py 8.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201
  1. # -*- coding: utf-8 -*-
  2. # Author : Charley
  3. # Python : 3.12.10
  4. # Date : 2026/08/19
  5. """集物星球在售抓取:免登录遍历各商品类型翻页拉 index/top,落库 jw_onsale_product_record。
  6. 每天 09/15/20/01 四档各跑一次(对齐 deca:上午/下午/晚上/凌晨场),每档同时把当刻在售写入每日快照表
  7. jw_onsale_daily_record(商品+日期唯一,同日多档刷成最新值),供在售日报按天差分出「在售趋势」。
  8. 抓完做下架对账:本轮未出现的商品置 is_on_sale=0。index/top 实测免登录,故 do_request 用默认 need_auth=False。
  9. """
  10. import sys
  11. import time
  12. from datetime import date
  13. import schedule
  14. from loguru import logger
  15. from tenacity import retry, stop_after_attempt, wait_fixed
  16. from mysql_pool import MySQLConnectionPool
  17. import jiwu_core as core
  18. import jw_detail
  19. # 商品类型:1福袋 2变风盒 3错版卡 4原盒;均属业务线 5(集卡)
  20. PRODUCT_TYPES = ["1", "2", "3", "4"]
  21. SYSTEM_BUSINESS_TYPE = 5
  22. PAGE_LIMIT = 20
  23. MAX_PAGES = 200 # 单类型翻页保护上限
  24. ENRICH_DETAIL = True # 是否逐商品补拉详情字段(免登录);量大可关
  25. logger.remove()
  26. logger.add("./logs/onsale_{time:YYYYMMDD}.log", encoding="utf-8", rotation="00:00",
  27. format="[{time:YYYY-MM-DD HH:mm:ss.SSS}] {level} {message}", level="DEBUG", retention="7 day")
  28. # 入库列(与 jw_onsale_product_record 对齐)
  29. _COLS = ["goods_id", "program_id", "goods_name", "product_type", "business_type", "block_type",
  30. "goods_ip_id", "goods_ip_name", "random_type", "corp_info_id", "corp_info_name",
  31. "amount", "highest_price", "lowest_price", "stock_amount", "residue_stock_amount",
  32. "sold_count", "specification_name", "report_status", "live_status", "sold_time",
  33. "product_pic", "is_on_sale"]
  34. def parse_product(rec: dict) -> dict:
  35. """把接口一条在售记录规范化为入库字典。
  36. Args:
  37. rec (dict): index/top 返回 records 里的一条。
  38. Returns:
  39. dict: 键与 jw_onsale_product_record 列对齐的字典。
  40. """
  41. stock = rec.get("stockAmount")
  42. residue = rec.get("residueStockAmount")
  43. sold = (stock - residue) if (isinstance(stock, int) and isinstance(residue, int)) else None
  44. return {
  45. "goods_id": rec.get("goodsId"), "program_id": rec.get("programId"),
  46. "goods_name": rec.get("goodsName"), "product_type": rec.get("productType"),
  47. "business_type": rec.get("businessType"), "block_type": rec.get("blockType"),
  48. "goods_ip_id": rec.get("goodsIPId"), "goods_ip_name": rec.get("goodsIPName"),
  49. "random_type": rec.get("randomType"), "corp_info_id": rec.get("corpInfoId"),
  50. "corp_info_name": rec.get("corpInfoName"), "amount": core.to_yuan(rec.get("amount")),
  51. "highest_price": core.to_yuan(rec.get("highestPrice")), "lowest_price": core.to_yuan(rec.get("lowestPrice")),
  52. "stock_amount": stock, "residue_stock_amount": residue, "sold_count": sold,
  53. "specification_name": rec.get("specificationName"), "report_status": rec.get("reportStatus"),
  54. "live_status": rec.get("liveStatus"), "sold_time": rec.get("soldTime"),
  55. "product_pic": rec.get("productPic"), "is_on_sale": 1,
  56. }
  57. def upsert_products(log, pool, rows: list) -> None:
  58. """把在售商品批量 upsert 到库(存在则更新库存/价格等最新状态)。
  59. Args:
  60. log: 日志对象。
  61. pool: 数据库连接池。
  62. rows (list): parse_product 结果列表。
  63. """
  64. if not rows:
  65. return
  66. cols_sql = ",".join(f"`{c}`" for c in _COLS)
  67. ph = ",".join(["%s"] * len(_COLS))
  68. upd = ",".join(f"`{c}`=VALUES(`{c}`)" for c in _COLS if c != "goods_id")
  69. sql = f"INSERT INTO jw_onsale_product_record ({cols_sql}) VALUES ({ph}) ON DUPLICATE KEY UPDATE {upd}"
  70. args_list = [tuple(r[c] for c in _COLS) for r in rows]
  71. pool.insert_many(query=sql, args_list=args_list)
  72. log.info(f"upsert 在售商品 {len(rows)} 条")
  73. # 每日快照入库列(与 jw_onsale_daily_record 对齐);唯一键命中(同日多档)时只刷 _SNAP_UPD 里的量价状态
  74. _SNAP_COLS = ["goods_id", "corp_info_id", "corp_info_name", "snapshot_date", "product_type",
  75. "stock_amount", "sold_count", "residue_stock_amount", "live_status", "amount"]
  76. _SNAP_UPD = ["stock_amount", "sold_count", "residue_stock_amount", "live_status", "amount"]
  77. def upsert_daily_snapshot(log, pool, rows: list) -> None:
  78. """把本页在售商品写入每日快照表(goods_id+snapshot_date 唯一,当天多档跑刷成最新值)。
  79. 对齐参考项目 deca 的每日快照:唯一键命中(同日多档)时只更新量价状态、保留商品/商家/日期,
  80. 保证每商品每天一行、恒为当天最后一次采集值,供在售日报按天差分出「在售趋势」。
  81. Args:
  82. log: 日志对象。
  83. pool: 数据库连接池。
  84. rows (list): parse_product 结果列表(含 amount 等,金额已换算为元)。
  85. """
  86. if not rows:
  87. return
  88. today = date.today().isoformat() # 快照日期;01:00 凌晨场归入新自然日(与 deca 一致)
  89. cols_sql = ",".join(f"`{c}`" for c in _SNAP_COLS)
  90. ph = ",".join(["%s"] * len(_SNAP_COLS))
  91. upd = ",".join(f"`{c}`=VALUES(`{c}`)" for c in _SNAP_UPD)
  92. sql = (f"INSERT INTO jw_onsale_daily_record ({cols_sql}) VALUES ({ph}) "
  93. f"ON DUPLICATE KEY UPDATE {upd}")
  94. args_list = [(r["goods_id"], r["corp_info_id"], r["corp_info_name"], today, r["product_type"],
  95. r["stock_amount"], r["sold_count"], r["residue_stock_amount"],
  96. r["live_status"], r["amount"]) for r in rows]
  97. pool.insert_many(query=sql, args_list=args_list)
  98. def sweep_offsale(log, pool, seen_ids: set) -> None:
  99. """下架对账:库中标记在售、但本轮未出现的商品置 is_on_sale=0。
  100. Args:
  101. log: 日志对象。
  102. pool: 数据库连接池。
  103. seen_ids (set): 本轮抓到的全部 goods_id。
  104. """
  105. rows = pool.select_all("SELECT goods_id FROM jw_onsale_product_record WHERE is_on_sale=1")
  106. db_ids = {r[0] for r in rows}
  107. gone = db_ids - seen_ids
  108. for gid in gone:
  109. pool.update_one("UPDATE jw_onsale_product_record SET is_on_sale=0 WHERE goods_id=%s", (gid,))
  110. if gone:
  111. log.info(f"下架对账:{len(gone)} 个商品置为已下架")
  112. def fetch_onsale(log, pool) -> None:
  113. """免登录遍历各商品类型翻页抓在售并落库,最后做下架对账。
  114. Args:
  115. log: 日志对象。
  116. pool: 数据库连接池。
  117. """
  118. seen_ids = set()
  119. for pt in PRODUCT_TYPES:
  120. for page in range(1, MAX_PAGES + 1):
  121. j = core.do_request(log, "/search/app/index/top", {
  122. "currentPage": str(page), "limit": str(PAGE_LIMIT),
  123. "productType": pt, "systemBusinessType": SYSTEM_BUSINESS_TYPE}) # 免登录
  124. if not j:
  125. break
  126. recs = (j.get("data") or {}).get("records") or []
  127. if not recs:
  128. break
  129. rows = [parse_product(r) for r in recs]
  130. seen_ids.update(r["goods_id"] for r in rows)
  131. upsert_products(log, pool, rows)
  132. upsert_daily_snapshot(log, pool, rows) # 同步写每日快照(供在售趋势差分)
  133. if len(recs) < PAGE_LIMIT: # 末页
  134. break
  135. time.sleep(0.3)
  136. sweep_offsale(log, pool, seen_ids)
  137. if ENRICH_DETAIL: # 逐商品补拉详情字段(免登录, detail_fetched 控制每商品补一次)
  138. jw_detail.enrich_detail(log, pool, "jw_onsale_product_record", list(seen_ids))
  139. log.success(f"本轮在售抓取完成,共 {len(seen_ids)} 个在售商品")
  140. @retry(stop=stop_after_attempt(100), wait=wait_fixed(3600), after=core.after_log)
  141. def main_task(log) -> None:
  142. """在售抓取主流程(挂了每小时重试,最多100次)。
  143. Args:
  144. log: 日志对象。
  145. Raises:
  146. RuntimeError: 数据库连接池异常时抛出以触发重试。
  147. """
  148. log.info("开始在售抓取" + "." * 40)
  149. pool = MySQLConnectionPool(log=log)
  150. if not pool.check_pool_health():
  151. log.error("数据库连接池异常")
  152. raise RuntimeError("数据库连接池异常")
  153. try:
  154. fetch_onsale(log, pool)
  155. except Exception as e:
  156. log.error(f"在售抓取异常: {e}")
  157. finally:
  158. log.info("在售抓取结束,等待下一轮" + "." * 20)
  159. def schedule_task():
  160. """定时入口:每天 09:00 / 15:00 / 20:00 / 01:00 各抓一次在售(对齐 deca 四档:上午/下午/晚上/凌晨场)。"""
  161. # main_task(log=logger) # 调试时取消注释立即跑一次
  162. for _hhmm in ("09:00", "15:00", "20:00", "01:00"):
  163. schedule.every().day.at(_hhmm).do(main_task, log=logger)
  164. while True:
  165. schedule.run_pending()
  166. time.sleep(1)
  167. if __name__ == "__main__":
  168. schedule_task()