_backfill_gap.py 9.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240
  1. # -*- coding: utf-8 -*-
  2. # Author : Charley
  3. # 历史回补脚本:按 pid 直取 /detail 与 /report/list,补齐 vod/list 滚动窗口已滚出的历史数据。
  4. # 可续跑:产品走 INSERT IGNORE,报告靠 report_state 门控(1/2=已完成,0/3=待处理)。
  5. #
  6. # 用法:
  7. # python _backfill_gap.py # 用默认区间 [1638, 2812]
  8. # python _backfill_gap.py 2800 3200 # 指定回补 pid 区间 [2800, 3200]
  9. # 说明:区间内已存在的 pid 会被自动跳过,无效/不存在的 pid 会被识别并计入"无效",不写库。
  10. import os
  11. import sys
  12. import time
  13. import requests
  14. os.chdir(os.path.dirname(os.path.abspath(__file__)))
  15. from mysql_pool import MySQLConnectionPool
  16. from loguru import logger
  17. logger.remove()
  18. logger.add("./_backfill_gap.log", encoding="utf-8",
  19. format="[{time:YYYY-MM-DD HH:mm:ss}] {message}", level="INFO")
  20. logger.add(sys.stdout, format="[{time:HH:mm:ss}] {message}", level="INFO")
  21. # 回补区间默认值(覆盖 2026 年 3-8 月缺口);可用命令行参数覆盖
  22. DEFAULT_PID_MIN, DEFAULT_PID_MAX = 1638, 2812
  23. DETAIL = "https://cxx.cardsvault.net/app/teamup/detail"
  24. REPORT = "https://cxx.cardsvault.net/app/teamup/report/list"
  25. def build_headers(token):
  26. """构造请求头(与线上爬虫一致的精简头)。
  27. Args:
  28. token (str): 鉴权 token,来自 super_vault_token 表。
  29. Returns:
  30. dict: requests 可用的 headers。
  31. """
  32. return {
  33. "User-Agent": "okhttp/4.9.0",
  34. "Authorization": token,
  35. "CXX-APP-API-VERSION": "V2",
  36. "Content-Type": "application/json; charset=UTF-8",
  37. }
  38. def req(method, url, headers, retries=4, **kw):
  39. """带重试的请求,返回解析后的 JSON。
  40. Args:
  41. method (str): "GET" 或 "POST"。
  42. url (str): 请求地址。
  43. headers (dict): 请求头。
  44. retries (int, optional): 重试次数。Defaults to 4。
  45. Returns:
  46. dict | None: 响应 JSON;多次失败返回 None。
  47. """
  48. for i in range(retries):
  49. try:
  50. r = requests.request(method, url, headers=headers, timeout=22, **kw)
  51. r.raise_for_status()
  52. return r.json()
  53. except Exception as e:
  54. if i == retries - 1:
  55. logger.warning(f"请求失败 {url} {kw.get('params') or kw.get('json')}: {e}")
  56. return None
  57. time.sleep(1)
  58. return None
  59. def map_detail_to_row(d):
  60. """把 /detail 的 data 映射为 super_vault_product_record 的一行。
  61. Args:
  62. d (dict): detail 接口返回的 data 字段。
  63. Returns:
  64. dict: 产品行;从 liveInfo 直接取 live_id / vod_url,report_state 置 0。
  65. """
  66. li = d.get("liveInfo") or {}
  67. total_price = d.get("totalPrice")
  68. sign_price = d.get("signPrice")
  69. return {
  70. "pid": d.get("id"),
  71. "title": d.get("title"),
  72. "serial": d.get("serial"),
  73. "type_name": d.get("typeName"),
  74. "is_pre": d.get("isPre"),
  75. "count": d.get("count"),
  76. "total_price": total_price / 100 if total_price else 0,
  77. "sign_price": sign_price / 100 if sign_price else 0,
  78. "sell_time": d.get("sellTime"),
  79. "sell_days": d.get("sellDays"),
  80. "status": d.get("status"),
  81. "group_num": d.get("groupNum"),
  82. "description": d.get("description"),
  83. "create_time": d.get("createTime"),
  84. "completion_time": d.get("completionTime"),
  85. "cover_url": (d.get("cover") or {}).get("url"),
  86. "anchor_id": (d.get("anchor") or {}).get("id"),
  87. "anchor_username": (d.get("anchor") or {}).get("userName"),
  88. "sold_count": d.get("soldCount"),
  89. "detail_url": d.get("detailUrl"),
  90. "goods_url": d.get("goodsUrl"),
  91. "standard_name": d.get("standardName"),
  92. "live_task_time": d.get("liveTaskTime"),
  93. "live_id": li.get("id"),
  94. "vod_url": (li.get("vod_info") or {}).get("vodUrl"),
  95. "report_state": 0,
  96. }
  97. def backfill_reports(pool, headers, pid):
  98. """抓取单个 pid 的全部拆卡报告并入库,同时更新 report_state。
  99. Args:
  100. pool (MySQLConnectionPool): 连接池。
  101. headers (dict): 请求头。
  102. pid (int): 商品 id。
  103. Returns:
  104. int: 本次入库尝试的报告条数(0 表示该商品无报告)。
  105. """
  106. page_num, total_pages, page_size = 1, 1, 20
  107. inserted = 0
  108. while page_num <= total_pages:
  109. j = req("POST", REPORT, headers,
  110. json={"pageSize": page_size, "my": 0, "pageNum": page_num, "tid": pid})
  111. if not j or j.get("status") != 200:
  112. raise RuntimeError(f"report/list 返回异常 pid={pid} page={page_num}")
  113. data = j.get("data") or {}
  114. total = data.get("total", 0)
  115. if page_num == 1:
  116. if total == 0:
  117. pool.update_one_or_dict(table="super_vault_product_record",
  118. data={"report_state": 2}, condition={"pid": pid})
  119. return 0
  120. total_pages = (total + page_size - 1) // page_size
  121. items = data.get("data", []) or []
  122. rows = [{
  123. "pid": pid,
  124. "user_name": it.get("userName"),
  125. "level": it.get("level"),
  126. "team_name_cn": it.get("teamNameCn"),
  127. "team_name_en": it.get("teamNameEn"),
  128. "count": it.get("count"),
  129. "picture_url": (it.get("picture") or {}).get("url"),
  130. "alias": it.get("alias"),
  131. "create_time": it.get("createTime"),
  132. } for it in items]
  133. if rows:
  134. pool.insert_many(table="super_vault_report_record", data_list=rows, ignore=True)
  135. inserted += len(rows)
  136. page_num += 1
  137. pool.update_one_or_dict(table="super_vault_product_record",
  138. data={"report_state": 1}, condition={"pid": pid})
  139. return inserted
  140. def main(pid_min, pid_max):
  141. """回补主流程:先补产品详情,再补拆卡报告,全程可续跑。
  142. Args:
  143. pid_min (int): 回补区间起始 pid(含)。
  144. pid_max (int): 回补区间结束 pid(含)。
  145. """
  146. pool = MySQLConnectionPool(log=logger)
  147. if not pool.check_pool_health():
  148. logger.error("数据库连接池异常")
  149. return
  150. token = pool.select_one("SELECT token FROM super_vault_token")[0]
  151. headers = build_headers(token)
  152. # ---------- Phase A: 补产品详情 ----------
  153. have = set(x[0] for x in pool.select_all("SELECT pid FROM super_vault_product_record"))
  154. missing = [p for p in range(pid_min, pid_max + 1) if p not in have]
  155. logger.info(f"Phase A 开始:区间[{pid_min},{pid_max}] 缺失 {len(missing)} 个 pid 待探测")
  156. inserted_p = invalid = 0
  157. for i, pid in enumerate(missing, 1):
  158. j = req("GET", DETAIL, headers, params={"id": str(pid)})
  159. d = (j or {}).get("data") or {}
  160. if j and j.get("status") == 200 and d.get("id"):
  161. pool.insert_many(table="super_vault_product_record",
  162. data_list=[map_detail_to_row(d)], ignore=True)
  163. inserted_p += 1
  164. else:
  165. invalid += 1
  166. if i % 50 == 0 or i == len(missing):
  167. logger.info(f" Phase A 进度 {i}/{len(missing)} 新增产品{inserted_p} 无效{invalid}")
  168. time.sleep(0.1)
  169. logger.info(f"Phase A 完成:新增产品 {inserted_p} 条,无效 pid {invalid} 个")
  170. # ---------- Phase B: 补拆卡报告(含之前超时遗留 + 本次新增) ----------
  171. pending = [x[0] for x in pool.select_all(
  172. "SELECT pid FROM super_vault_product_record WHERE report_state NOT IN (1,2) ORDER BY pid")]
  173. logger.info(f"Phase B 开始:待抓报告产品 {len(pending)} 个")
  174. ok = err = total_reports = 0
  175. for i, pid in enumerate(pending, 1):
  176. try:
  177. total_reports += backfill_reports(pool, headers, pid)
  178. ok += 1
  179. except Exception as e:
  180. err += 1
  181. pool.update_one_or_dict(table="super_vault_product_record",
  182. data={"report_state": 3}, condition={"pid": pid})
  183. logger.warning(f" pid={pid} 报告回补失败: {e}")
  184. if i % 25 == 0 or i == len(pending):
  185. logger.info(f" Phase B 进度 {i}/{len(pending)} 成功{ok} 失败{err} 累计报告{total_reports}")
  186. time.sleep(0.05)
  187. logger.info(f"Phase B 完成:成功{ok} 失败{err},新增拆卡报告约 {total_reports} 条")
  188. p_cnt = pool.select_one("SELECT COUNT(*) FROM super_vault_product_record")[0]
  189. r_cnt = pool.select_one("SELECT COUNT(*) FROM super_vault_report_record")[0]
  190. logger.info(f"===== 回补结束:产品表共 {p_cnt} 条,报告表共 {r_cnt} 条 =====")
  191. def parse_range_args(argv):
  192. """从命令行参数解析回补 pid 区间。
  193. Args:
  194. argv (list[str]): sys.argv[1:],可为空或 [起始pid, 结束pid]。
  195. Returns:
  196. tuple[int, int]: (pid_min, pid_max);未传参时返回默认区间。
  197. Raises:
  198. ValueError: 参数无法转为整数或 pid_min > pid_max 时抛出。
  199. """
  200. if len(argv) >= 2:
  201. pid_min, pid_max = int(argv[0]), int(argv[1])
  202. if pid_min > pid_max:
  203. raise ValueError(f"起始pid({pid_min}) 不能大于 结束pid({pid_max})")
  204. return pid_min, pid_max
  205. return DEFAULT_PID_MIN, DEFAULT_PID_MAX
  206. if __name__ == "__main__":
  207. _pid_min, _pid_max = parse_range_args(sys.argv[1:])
  208. main(_pid_min, _pid_max)