Przeglądaj źródła

feat(super_vault): 实现历史数据缺口回补与陈旧数据刷新

- 新增 _backfill_gap.py 脚本,实现按 pid 区间回补产品详情和拆卡报告
- 支持续跑与自动跳过已存在或无效 pid
- 添加 _run_update2.py 脚本,用于调用更新陈旧产品信息接口
- 调整 start_cxx_spider.py,暂停执行 cxx_daily_main 任务,保留 cxx_sale_main
- super_vault_daily_spider.py 新增刷新陈旧拼团字段功能,修复 INSERT IGNORE 导致的字段固化问题
- 实现针对已售罄团的 sold_time 精准补齐逻辑
- 通过 detail 接口批量更新状态、价格、销量、直播信息等动态字段,确保数据及时同步
- 增加日志和异常处理,支持重试机制,保障数据补齐流程稳定运行
charley 1 miesiąc temu
rodzic
commit
dd8d30f9ef

+ 240 - 0
cxx_spider/_backfill_gap.py

@@ -0,0 +1,240 @@
+# -*- 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)

+ 12 - 0
cxx_spider/_run_update2.py

@@ -0,0 +1,12 @@
+# -*- coding: utf-8 -*-
+import sys, io
+sys.stdout = io.TextIOWrapper(sys.stdout.buffer, encoding='utf-8')
+from loguru import logger
+logger.remove()
+logger.add(sys.stderr, format="[{time:HH:mm:ss}] {message}", level="INFO")
+import super_vault_daily_spider as s
+from mysql_pool import MySQLConnectionPool
+pool = MySQLConnectionPool(log=logger)
+token = pool.select_one("SELECT token FROM super_vault_token")[0]
+s.update_stale_products(logger, pool, token)
+print("UPDATE_DONE")

+ 1 - 1
cxx_spider/start_cxx_spider.py

@@ -31,7 +31,7 @@ def schedule_task():
     爬虫模块的启动文件
     """
     # 立即运行一次任务
-    run_threaded(cxx_daily_main, log=logger)
+    # run_threaded(cxx_daily_main, log=logger)
     # time.sleep(5)
 
     run_threaded(cxx_sale_main, log=logger)

+ 157 - 2
cxx_spider/super_vault_daily_spider.py

@@ -111,7 +111,7 @@ def get_vod_single_page(log, page_num=1, token=""):
     response.raise_for_status()
 
     result = response.json()
-    # print(result)
+    print(result)
     if result.get("status") == 200:
         data = result.get("data", {})
         total = data.get("total", 0)
@@ -364,6 +364,149 @@ def get_report_list(log, detail_id, token, sql_pool):
         page_num += 1
 
 
+# ----------------------------------------------------------------------------------------------------------------------
+# 补充/更新任务:修正 INSERT IGNORE「只增不改」导致的字段固化
+#   列表接口靠 pid 唯一键 INSERT IGNORE 去重,商品首次入库常在预售阶段(sign_price/sold_count/
+#   status 尚为 0/初始值,且未开播故 live_id/vod_url/completion_time 为空),此后即便平台更新,
+#   库里也不再回写——这正是「大量团缺 completion_time / 直播回放」的根因(实测 status!=9 且已售出
+#   的团里约 92% 其实已完成、detail 里已有 completionTime)。本段对 status!=9 或缺直播信息的团
+#   逐个查 detail 接口刷新会变字段,补齐 status/completion_time/live_task_time/live_id/vod_url 等。
+
+@retry(stop=stop_after_attempt(5), wait=wait_fixed(1), after=after_log)
+def get_teamup_detail(log, pid, token):
+    """请求单个拼团的详情,返回其最新完整字段。
+
+    Args:
+        log: 日志对象。
+        pid (int): 拼团商品 id。
+        token (str): 鉴权 token。
+
+    Returns:
+        dict | None: 详情 data 字典;接口业务码非 200 或无数据时返回 None。
+
+    Raises:
+        requests.HTTPError: HTTP 状态码非 2xx 时抛出以触发重试。
+    """
+    HEADERS["Authorization"] = token
+    response = requests.get("https://cxx.cardsvault.net/app/teamup/detail",
+                            headers=HEADERS, params={"id": str(pid)}, timeout=22)
+    response.raise_for_status()
+    result = response.json()
+    if result.get("status") != 200:
+        log.warning(f"detail 接口业务码异常 pid={pid}: {result.get('msg', '未知错误')}")
+        return None
+    return result.get("data") or None
+
+
+def build_update_fields(data):
+    """从详情 data 里挑出「会随售卖进程变化」的字段,组装成 UPDATE 用的 dict。
+
+    signPrice/totalPrice 接口以「分」返回,统一 /100 转「元」入库(与列表解析口径一致)。
+    直播相关字段(live_task_time/live_id/vod_url)为「直播后才产生」,首次预售入库常为空,一并补。
+
+    Args:
+        data (dict): get_teamup_detail 返回的详情字典。
+
+    Returns:
+        dict: 待更新字段 {列名: 值},含 status/sold_count/sign_price/total_price/count/
+              completion_time/sold_time/sell_time/sell_days/live_task_time/live_id/vod_url。
+    """
+    total_price = data.get("totalPrice")
+    sign_price = data.get("signPrice")
+    live_info = data.get("liveInfo") or {}               # 未开播可能为 None
+    vod_info = live_info.get("vod_info") or {}           # 无回放可能为 None
+    return {
+        "status": data.get("status"),                          # 状态(9:完成 8:待发货 ...)
+        "sold_count": data.get("soldCount"),                   # 已售份数
+        "sign_price": sign_price / 100 if sign_price else 0,   # 单价(分→元)
+        "total_price": total_price / 100 if total_price else 0,  # 团总价(分→元)
+        "count": data.get("count"),                            # 总份数
+        "completion_time": data.get("completionTime"),         # 完成时间(发货完成态才有)
+        "sold_time": data.get("soldTime"),                     # 售罄成交时刻(补 completion_time 缺口,近期团几乎都有)
+        "sell_time": data.get("sellTime"),                     # 开售时间
+        "sell_days": data.get("sellDays"),                     # 售卖天数
+        "live_task_time": data.get("liveTaskTime"),            # 直播时间
+        "live_id": live_info.get("id"),                        # 直播 id
+        "vod_url": vod_info.get("vodUrl"),                     # 直播回放地址
+    }
+
+
+def update_stale_products(log, sql_pool, token):
+    """刷新陈旧拼团的会变字段,回写被 INSERT IGNORE 固化的旧值。
+
+    候选 = status!=9(字段仍会变、常有已完成却未回写的团)或缺直播信息(live_task_time/
+    live_id/vod_url 之一为空,直播后才产生、首次入库时常缺)的商品;逐个查详情并 UPDATE,
+    使 sign_price/sold_count/status/total_price/completion_time/直播回放等保持最新。
+    已完成(status=9)且直播信息齐全的字段已定型,不在候选内、省请求。
+
+    Args:
+        log: 日志对象。
+        sql_pool (MySQLConnectionPool): 数据库连接池。
+        token (str): 鉴权 token。
+    """
+    rows = sql_pool.select_all(
+        "SELECT pid FROM super_vault_product_record "
+        "WHERE status <> 9 OR status IS NULL "
+        "OR vod_url IS NULL OR live_id IS NULL OR live_task_time IS NULL "
+        "OR (sold_time IS NULL AND sold_count > 0)")
+    pids = [r[0] for r in rows] if rows else []
+    log.info(f"待刷新未完成拼团 {len(pids)} 个")
+    updated = 0
+    for pid in pids:
+        try:
+            data = get_teamup_detail(log, pid, token)
+            if not data:
+                continue
+            sql_pool.update_one_or_dict(
+                table="super_vault_product_record",
+                data=build_update_fields(data),
+                condition={"pid": pid})
+            updated += 1
+        except Exception as e:
+            log.error(f"刷新 pid={pid} 失败: {e}")
+        time.sleep(0.3)  # 限速,避免请求过密触发风控
+    log.info(f"未完成拼团刷新完成,成功更新 {updated}/{len(pids)} 个")
+
+
+def backfill_sold_time(log, sql_pool, token):
+    """精准补齐已售罄团缺失的 sold_time(售罄成交时刻)。
+
+    只挑「已进入售后阶段(status 3/8)、已有人购买(sold_count>0)、但 sold_time 仍为空」的团:
+    这类团常在列表首次入库(INSERT IGNORE)后才售罄,soldTime 由平台事后生成、旧值不再回写,故需
+    逐个查 detail 补。刻意不扫 status=9 的历史团——早期平台无 soldTime 字段、源头即为空,查了也
+    补不到、纯浪费请求;未售罄的(status 0/10 等)源头同样无值,一并排除。仅当 detail 真返回
+    soldTime 时才 UPDATE,且只写 sold_time 单字段,避免误改其它数据。
+
+    Args:
+        log: 日志对象。
+        sql_pool (MySQLConnectionPool): 数据库连接池。
+        token (str): 鉴权 token。
+    """
+    rows = sql_pool.select_all(
+        "SELECT pid FROM super_vault_product_record "
+        "WHERE sold_time IS NULL AND sold_count > 0 AND status IN (3, 8)")
+    pids = [r[0] for r in rows] if rows else []
+    log.info(f"待补 sold_time 的已售罄团 {len(pids)} 个")
+    filled = 0
+    for pid in pids:
+        try:
+            data = get_teamup_detail(log, pid, token)
+            if not data:
+                continue
+            sold_time = data.get("soldTime")
+            if not sold_time:  # 源头也无(团未真正售罄),跳过不误写
+                continue
+            sql_pool.update_one_or_dict(
+                table="super_vault_product_record",
+                data={"sold_time": sold_time},
+                condition={"pid": pid})
+            filled += 1
+        except Exception as e:
+            log.error(f"补 sold_time pid={pid} 失败: {e}")
+        time.sleep(0.3)  # 限速,避免请求过密触发风控
+    log.info(f"sold_time 补齐完成,成功回写 {filled}/{len(pids)} 个")
+
+
 @retry(stop=stop_after_attempt(100), wait=wait_fixed(3600), after=after_log)
 def cxx_daily_main(log):
     """
@@ -388,6 +531,18 @@ def cxx_daily_main(log):
         except Exception as e:
             log.error(f"Error fetching last_product_id: {e}")
 
+        # 精准补齐已售罄团(status 3/8)的 sold_time:轻量、只改单字段,最新入库的团先补上
+        try:
+            backfill_sold_time(log, sql_pool, token[0])
+        except Exception as e:
+            log.error(f"Error backfilling sold_time: {e}")
+
+        # 刷新未完成拼团(status!=9)的会变字段,修正 INSERT IGNORE 固化的旧 sign_price/sold_count/status 等
+        try:
+            update_stale_products(log, sql_pool, token[0])
+        except Exception as e:
+            log.error(f"Error updating stale products: {e}")
+
         time.sleep(5)
 
         # 获取所有 report_state = 0 的 pid
@@ -417,7 +572,7 @@ def cxx_daily_main(log):
 
 if __name__ == '__main__':
     # get_vod_list(logger, None, '')
-    # get_vod_single_page(logger, 1)
+    # get_vod_single_page(logger, 1,'eyJ0eXAiOiJKV1QiLCJhbGciOiJIUzI1NiJ9.eyJzdWIiOiJDSEFPWElOWElORyNBUFAiLCJhdWQiOiJDSEFPWElOWElORyIsIm5iZiI6MTc4NjEwMjA5MywiZGF0YSI6Ijk1MjMiLCJpc3MiOiI3ViNweHlQZSIsImV4cCI6MTc4NzMwMjA5MywiaWF0IjoxNzg2MTAyMDkzLCJqdGkiOiI4MWE0YTIxNi1mNTc0LTQ2ZDUtYmIxOS05ZjMzNWMxYTczM2UifQ.ieTyqCqFGsDBt1WkO1aHI58hngXP9jcViNYzGRdkOYA')
     # get_report_single_page(logger,1, '','')
 
     cxx_daily_main(logger)