| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201 |
- # -*- 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()
|