# -*- coding: utf-8 -*- # Author : Charley # 历史回补脚本:按 pid 直取 /detail 与 /report/list,补齐 vod/list 滚动窗口已滚出的历史数据。 # 可续跑:产品走 INSERT IGNORE,报告靠 report_state 门控(1/2=已完成,0/3=待处理)。 # # 用法: # python _backfill_gap.py # 用默认区间 [1638, 2812] # python _backfill_gap.py 2800 3200 # 指定回补 pid 区间 [2800, 3200] # 说明:区间内已存在的 pid 会被自动跳过,无效/不存在的 pid 会被识别并计入"无效",不写库。 import os import sys import time import requests os.chdir(os.path.dirname(os.path.abspath(__file__))) from mysql_pool import MySQLConnectionPool from loguru import logger logger.remove() logger.add("./_backfill_gap.log", encoding="utf-8", format="[{time:YYYY-MM-DD HH:mm:ss}] {message}", level="INFO") logger.add(sys.stdout, format="[{time:HH:mm:ss}] {message}", level="INFO") # 回补区间默认值(覆盖 2026 年 3-8 月缺口);可用命令行参数覆盖 DEFAULT_PID_MIN, DEFAULT_PID_MAX = 1638, 2812 DETAIL = "https://cxx.cardsvault.net/app/teamup/detail" REPORT = "https://cxx.cardsvault.net/app/teamup/report/list" def build_headers(token): """构造请求头(与线上爬虫一致的精简头)。 Args: token (str): 鉴权 token,来自 super_vault_token 表。 Returns: dict: requests 可用的 headers。 """ return { "User-Agent": "okhttp/4.9.0", "Authorization": token, "CXX-APP-API-VERSION": "V2", "Content-Type": "application/json; charset=UTF-8", } def req(method, url, headers, retries=4, **kw): """带重试的请求,返回解析后的 JSON。 Args: method (str): "GET" 或 "POST"。 url (str): 请求地址。 headers (dict): 请求头。 retries (int, optional): 重试次数。Defaults to 4。 Returns: dict | None: 响应 JSON;多次失败返回 None。 """ for i in range(retries): try: r = requests.request(method, url, headers=headers, timeout=22, **kw) r.raise_for_status() return r.json() except Exception as e: if i == retries - 1: logger.warning(f"请求失败 {url} {kw.get('params') or kw.get('json')}: {e}") return None time.sleep(1) return None def map_detail_to_row(d): """把 /detail 的 data 映射为 super_vault_product_record 的一行。 Args: d (dict): detail 接口返回的 data 字段。 Returns: dict: 产品行;从 liveInfo 直接取 live_id / vod_url,report_state 置 0。 """ li = d.get("liveInfo") or {} total_price = d.get("totalPrice") sign_price = d.get("signPrice") return { "pid": d.get("id"), "title": d.get("title"), "serial": d.get("serial"), "type_name": d.get("typeName"), "is_pre": d.get("isPre"), "count": d.get("count"), "total_price": total_price / 100 if total_price else 0, "sign_price": sign_price / 100 if sign_price else 0, "sell_time": d.get("sellTime"), "sell_days": d.get("sellDays"), "status": d.get("status"), "group_num": d.get("groupNum"), "description": d.get("description"), "create_time": d.get("createTime"), "completion_time": d.get("completionTime"), "cover_url": (d.get("cover") or {}).get("url"), "anchor_id": (d.get("anchor") or {}).get("id"), "anchor_username": (d.get("anchor") or {}).get("userName"), "sold_count": d.get("soldCount"), "detail_url": d.get("detailUrl"), "goods_url": d.get("goodsUrl"), "standard_name": d.get("standardName"), "live_task_time": d.get("liveTaskTime"), "live_id": li.get("id"), "vod_url": (li.get("vod_info") or {}).get("vodUrl"), "report_state": 0, } def backfill_reports(pool, headers, pid): """抓取单个 pid 的全部拆卡报告并入库,同时更新 report_state。 Args: pool (MySQLConnectionPool): 连接池。 headers (dict): 请求头。 pid (int): 商品 id。 Returns: int: 本次入库尝试的报告条数(0 表示该商品无报告)。 """ page_num, total_pages, page_size = 1, 1, 20 inserted = 0 while page_num <= total_pages: j = req("POST", REPORT, headers, json={"pageSize": page_size, "my": 0, "pageNum": page_num, "tid": pid}) if not j or j.get("status") != 200: raise RuntimeError(f"report/list 返回异常 pid={pid} page={page_num}") data = j.get("data") or {} total = data.get("total", 0) if page_num == 1: if total == 0: pool.update_one_or_dict(table="super_vault_product_record", data={"report_state": 2}, condition={"pid": pid}) return 0 total_pages = (total + page_size - 1) // page_size items = data.get("data", []) or [] rows = [{ "pid": pid, "user_name": it.get("userName"), "level": it.get("level"), "team_name_cn": it.get("teamNameCn"), "team_name_en": it.get("teamNameEn"), "count": it.get("count"), "picture_url": (it.get("picture") or {}).get("url"), "alias": it.get("alias"), "create_time": it.get("createTime"), } for it in items] if rows: pool.insert_many(table="super_vault_report_record", data_list=rows, ignore=True) inserted += len(rows) page_num += 1 pool.update_one_or_dict(table="super_vault_product_record", data={"report_state": 1}, condition={"pid": pid}) return inserted def main(pid_min, pid_max): """回补主流程:先补产品详情,再补拆卡报告,全程可续跑。 Args: pid_min (int): 回补区间起始 pid(含)。 pid_max (int): 回补区间结束 pid(含)。 """ pool = MySQLConnectionPool(log=logger) if not pool.check_pool_health(): logger.error("数据库连接池异常") return token = pool.select_one("SELECT token FROM super_vault_token")[0] headers = build_headers(token) # ---------- Phase A: 补产品详情 ---------- have = set(x[0] for x in pool.select_all("SELECT pid FROM super_vault_product_record")) missing = [p for p in range(pid_min, pid_max + 1) if p not in have] logger.info(f"Phase A 开始:区间[{pid_min},{pid_max}] 缺失 {len(missing)} 个 pid 待探测") inserted_p = invalid = 0 for i, pid in enumerate(missing, 1): j = req("GET", DETAIL, headers, params={"id": str(pid)}) d = (j or {}).get("data") or {} if j and j.get("status") == 200 and d.get("id"): pool.insert_many(table="super_vault_product_record", data_list=[map_detail_to_row(d)], ignore=True) inserted_p += 1 else: invalid += 1 if i % 50 == 0 or i == len(missing): logger.info(f" Phase A 进度 {i}/{len(missing)} 新增产品{inserted_p} 无效{invalid}") time.sleep(0.1) logger.info(f"Phase A 完成:新增产品 {inserted_p} 条,无效 pid {invalid} 个") # ---------- Phase B: 补拆卡报告(含之前超时遗留 + 本次新增) ---------- pending = [x[0] for x in pool.select_all( "SELECT pid FROM super_vault_product_record WHERE report_state NOT IN (1,2) ORDER BY pid")] logger.info(f"Phase B 开始:待抓报告产品 {len(pending)} 个") ok = err = total_reports = 0 for i, pid in enumerate(pending, 1): try: total_reports += backfill_reports(pool, headers, pid) ok += 1 except Exception as e: err += 1 pool.update_one_or_dict(table="super_vault_product_record", data={"report_state": 3}, condition={"pid": pid}) logger.warning(f" pid={pid} 报告回补失败: {e}") if i % 25 == 0 or i == len(pending): logger.info(f" Phase B 进度 {i}/{len(pending)} 成功{ok} 失败{err} 累计报告{total_reports}") time.sleep(0.05) logger.info(f"Phase B 完成:成功{ok} 失败{err},新增拆卡报告约 {total_reports} 条") p_cnt = pool.select_one("SELECT COUNT(*) FROM super_vault_product_record")[0] r_cnt = pool.select_one("SELECT COUNT(*) FROM super_vault_report_record")[0] logger.info(f"===== 回补结束:产品表共 {p_cnt} 条,报告表共 {r_cnt} 条 =====") def parse_range_args(argv): """从命令行参数解析回补 pid 区间。 Args: argv (list[str]): sys.argv[1:],可为空或 [起始pid, 结束pid]。 Returns: tuple[int, int]: (pid_min, pid_max);未传参时返回默认区间。 Raises: ValueError: 参数无法转为整数或 pid_min > pid_max 时抛出。 """ if len(argv) >= 2: pid_min, pid_max = int(argv[0]), int(argv[1]) if pid_min > pid_max: raise ValueError(f"起始pid({pid_min}) 不能大于 结束pid({pid_max})") return pid_min, pid_max return DEFAULT_PID_MIN, DEFAULT_PID_MAX if __name__ == "__main__": _pid_min, _pid_max = parse_range_args(sys.argv[1:]) main(_pid_min, _pid_max)