kj_daily_spider.py 30 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747
  1. # -*- coding: utf-8 -*-
  2. # Author : Charley
  3. # Python : 3.12.10
  4. # Date : 2026/07/02
  5. """卡集(card.kaji.com)每日采集爬虫。
  6. 三级采集管道(结构对齐 zc_new_daily_spider):
  7. 1. 商户发现:遍历「在售主列表」forsale/main,从每个在售商品里提取商户,写入 kj_shop_record。
  8. 2. 商品采集:遍历库内每个商户的「已售/已完成列表」merchant?tp=2,逐商品补 detail 后写入 kj_product_record。
  9. 3. 玩家采集:遍历 kj_product_record 中未采集过玩家的商品,抓 result/merge 中奖/参与名单写入 kj_player_record。
  10. 鉴权:卡集接口的 Authorization 采用 Open-Auth-Sig 请求签名(见 kj_auth.py),
  11. 每请求用当前时间实时生成,无登录 token、无过期问题,适合长期无人值守。
  12. 其中商户已售列表 merchant?tp=2 为公开接口,不需要签名。
  13. """
  14. import re
  15. import sys
  16. import time
  17. from datetime import datetime
  18. from urllib.parse import unquote
  19. import requests
  20. import schedule
  21. from loguru import logger
  22. from tenacity import retry, stop_after_attempt, wait_fixed
  23. from mysql_pool import MySQLConnectionPool
  24. import kj_auth
  25. logger.remove()
  26. logger.add("./logs/{time:YYYYMMDD}.log", encoding="utf-8", rotation="00:00",
  27. format="[{time:YYYY-MM-DD HH:mm:ss.SSS}] {level} {message}",
  28. level="DEBUG", retention="7 day")
  29. # ==================== 基础配置 ====================
  30. # 两个业务域名(来源:抓包):主列表 / videoPlay 走 server 域,商户列表 / 详情 / 玩家走 page 域
  31. SERVER_BASE = "https://server.ssl1.kaji6.com"
  32. PAGE_BASE = "https://page.ssl1.kaji6.com"
  33. PAGE_SIZE = 10 # 列表接口每页条数
  34. PLAYER_PAGE_SIZE = 30 # 玩家名单每页条数(抓包实测为 30)
  35. MAX_PAGES = 100 # 单个列表翻页保护上限,防异常时无限翻页
  36. # 是否为每个商品补拉 detail(拿开始/结束时间、商户销量/粉丝、规格)。
  37. # 关闭可大幅减少请求量,但 start_date/end_date/规格/商户 sold_number/fans 会缺失。
  38. FETCH_DETAIL = True
  39. # 是否使用代理。卡集实测可直连;若遇 IP 风控再置 True 并配置 get_proxys。
  40. USE_PROXY = False
  41. # 固定 UA(来源:抓包,卡集为 uni-app WebView)
  42. UA = ("Mozilla/5.0 (Linux; Android 11; Pixel 5 Build/RQ3A.211001.001; wv) "
  43. "AppleWebKit/537.36 (KHTML, like Gecko) Version/4.0 Chrome/148.0.7778.120 "
  44. "Mobile Safari/537.36 uni-app Html5Plus/1.0 (Immersed/52.727272)")
  45. # 基础请求头(Authorization 每请求单独生成,不放这里)
  46. BASE_HEADERS = {
  47. "deviceType": "phone",
  48. "Accept": "application/json, text/plain, */*",
  49. "plat": "android",
  50. "version": "2.5.39", # 来源:抓包 App 版本
  51. "appVersionCode": "10003", # 来源:抓包 App versionCode
  52. "user-agent": UA,
  53. }
  54. def after_log(retry_state):
  55. """tenacity 重试回调,记录每次尝试的结果。
  56. Args:
  57. retry_state: tenacity 传入的 RetryCallState 对象,含调用参数与结果。
  58. """
  59. # 约定业务函数首个位置参数为 log;取不到时回退全局 logger
  60. if retry_state.args and len(retry_state.args) > 0:
  61. log = retry_state.args[0]
  62. else:
  63. log = logger
  64. if retry_state.outcome.failed:
  65. log.warning(f"Function '{retry_state.fn.__name__}', Attempt {retry_state.attempt_number} Times")
  66. else:
  67. log.info(f"Function '{retry_state.fn.__name__}', Attempt {retry_state.attempt_number} succeeded")
  68. @retry(stop=stop_after_attempt(5), wait=wait_fixed(1), after=after_log)
  69. def get_proxys(log):
  70. """获取隧道代理配置(默认不启用,见 USE_PROXY)。
  71. Args:
  72. log: 日志对象。
  73. Returns:
  74. dict: requests 可用的 proxies 字典。
  75. Raises:
  76. Exception: 组装代理配置异常时向上抛出以触发重试。
  77. """
  78. tunnel = "x371.kdltps.com:15818"
  79. kdl_username = "t13753103189895"
  80. kdl_password = "o0yefv6z"
  81. try:
  82. proxies = {
  83. "http": "http://%(user)s:%(pwd)s@%(proxy)s/" % {"user": kdl_username, "pwd": kdl_password, "proxy": tunnel},
  84. "https": "http://%(user)s:%(pwd)s@%(proxy)s/" % {"user": kdl_username, "pwd": kdl_password,
  85. "proxy": tunnel},
  86. }
  87. return proxies
  88. except Exception as e:
  89. log.error(f"Error getting proxy: {e}")
  90. raise e
  91. def ts_to_dt(ts) -> str | None:
  92. """将秒级 Unix 时间戳转成 'YYYY-MM-DD HH:MM:SS' 字符串。
  93. Args:
  94. ts (int | None): 秒级时间戳;为 None / 0 时返回 None。
  95. Returns:
  96. str | None: 格式化时间字符串;无有效时间戳时返回 None。
  97. """
  98. if not ts:
  99. return None
  100. return datetime.fromtimestamp(ts).strftime("%Y-%m-%d %H:%M:%S")
  101. def parse_desc(desc: str) -> tuple[str | None, int | None]:
  102. """从 detail.good.desc(URL 编码)中提取拼团系列与拼团天数。
  103. desc 解码后各字段以回车(\\r)分隔,形如:
  104. 拼团系列:27-28
  105. 拼团规格:1张/包,1包/盒,1盒/箱,共3包 3张
  106. 拼团份数:1205份
  107. 拼团时间:15天
  108. Args:
  109. desc (str): detail.good.desc 原始 URL 编码串;为空时返回 (None, None)。
  110. Returns:
  111. tuple[str | None, int | None]: (拼团系列, 拼团天数)。系列为文本(如 "27-28"),
  112. 天数为整数(如 15);对应项提取失败为 None。
  113. """
  114. if not desc:
  115. return None, None
  116. text = unquote(desc) # URL 解码:%EF%BC%9A→: %0D→\r
  117. m_series = re.search(r"拼团系列[::]\s*(.*?)\s*(?:[\r\n]|$)", text)
  118. m_days = re.search(r"拼团时间[::]\s*(\d+)", text) # 只取数字,"15天"→15
  119. series = m_series.group(1).strip() if m_series else None
  120. days = int(m_days.group(1)) if m_days else None
  121. return series, days
  122. @retry(stop=stop_after_attempt(5), wait=wait_fixed(1), after=after_log)
  123. def kj_request(log, path: str, method: str = "GET", base: str = PAGE_BASE,
  124. params: dict = None, json_body: dict = None,
  125. need_auth: bool = True, extra_headers: dict = None):
  126. """卡集通用请求函数(带重试)。
  127. 根据 need_auth 决定是否为该接口路径实时生成 Open-Auth-Sig 签名并注入 Authorization。
  128. Args:
  129. log: 日志对象。
  130. path (str): 接口相对路径,不含域名 / '/api/v4/' 前缀 / query,如 "goodlist/forsale/main"。
  131. method (str, optional): 请求方法,"GET" 或 "POST"。Defaults to "GET"。
  132. base (str, optional): 域名前缀,SERVER_BASE 或 PAGE_BASE。Defaults to PAGE_BASE。
  133. params (dict, optional): URL query 参数。Defaults to None。
  134. json_body (dict, optional): POST 的 JSON body。Defaults to None。
  135. need_auth (bool, optional): 是否注入 Authorization 签名。Defaults to True。
  136. extra_headers (dict, optional): 追加 / 覆盖的请求头。Defaults to None。
  137. Returns:
  138. dict | None: 响应 JSON;非 200 时抛异常触发重试。
  139. Raises:
  140. RuntimeError: HTTP 状态码非 200 时抛出。
  141. """
  142. url = f"{base}/api/v4/{path}"
  143. req_headers = BASE_HEADERS.copy()
  144. if need_auth:
  145. # 签名基于不含 query 的 path,kj_auth 内部会补 '/api/v4/' 前缀并裁掉 query
  146. req_headers["Authorization"] = kj_auth.gen_authorization(path)
  147. if extra_headers:
  148. req_headers.update(extra_headers)
  149. proxies = get_proxys(log) if USE_PROXY else None
  150. if method.upper() == "POST":
  151. resp = requests.post(url, headers=req_headers, json=json_body, params=params, timeout=(5, 30), proxies=proxies)
  152. else:
  153. resp = requests.get(url, headers=req_headers, params=params, timeout=(5, 30), proxies=proxies)
  154. if resp.status_code != 200:
  155. log.error(f"请求失败 {resp.status_code}: {url}")
  156. raise RuntimeError(f"HTTP {resp.status_code}")
  157. return resp.json()
  158. # ==================== 一、商户发现(在售主列表) ====================
  159. def get_forsale_page(log, fetch_from: int, fetch_size: int = PAGE_SIZE):
  160. """获取一页在售主列表 goodlist/forsale/main。
  161. Args:
  162. log: 日志对象。
  163. fetch_from (int): 游标起始位置(首页为 1,逐页按 fetch_size 递增)。
  164. fetch_size (int, optional): 每页条数。Defaults to PAGE_SIZE。
  165. Returns:
  166. dict | None: 响应 JSON(含 goodList / isFetchEnd);失败返回 None。
  167. """
  168. return kj_request(
  169. log, "goodlist/forsale/main", base=SERVER_BASE,
  170. params={"fetchFrom": str(fetch_from), "fetchSize": str(fetch_size)},
  171. need_auth=True,
  172. )
  173. def save_shops(log, good_list: list, sql_pool, seen: set) -> int:
  174. """从在售商品列表提取商户并写入 kj_shop_record(存在则更新店名)。
  175. 在售商品项只带 merchantAlias / merchantName,故本步只写 shop_id + shop_name;
  176. sold_number / fans 由商品采集阶段的 detail.publisher 回填。
  177. Args:
  178. log: 日志对象。
  179. good_list (list[dict]): forsale 返回的 goodList,每项含 merchantAlias / merchantName。
  180. sql_pool (MySQLConnectionPool): MySQL 连接池。
  181. seen (set): 跨页去重的 shop_id 集合,避免重复写库。
  182. Returns:
  183. int: 本次实际写库的商户数(已去重)。
  184. """
  185. info_list = []
  186. # print(good_list)
  187. for item in good_list:
  188. shop_id = item.get("merchantAlias")
  189. shop_name = item.get("merchantName")
  190. data_dict = {"shop_id": shop_id, "shop_name": shop_name}
  191. if not shop_id or shop_id in seen: continue
  192. seen.add(shop_id)
  193. info_list.append(data_dict)
  194. # print(data_dict)
  195. if not info_list:
  196. log.info("无新商户,不写库")
  197. return 0
  198. # 存在则更新店名并把 is_deleted 刷回 0:能出现在在售列表 = 商户在营业
  199. # (曾被标记注销的商户在此复活,重新纳入采集)
  200. sql = ("INSERT INTO kj_shop_record (shop_id, shop_name) VALUES (%s, %s) "
  201. "ON DUPLICATE KEY UPDATE shop_name = VALUES(shop_name), is_deleted = 0")
  202. args_list = [(d["shop_id"], d["shop_name"]) for d in info_list]
  203. sql_pool.insert_many(query=sql, args_list=args_list)
  204. return len(info_list)
  205. def get_shop_list(log, sql_pool) -> int:
  206. """遍历在售主列表翻页,发现并写入全部在售商户。
  207. 翻页靠响应的 isFetchEnd 与空页判断,shop_id 唯一键去重兜底翻页游标的不确定性。
  208. Args:
  209. log: 日志对象。
  210. sql_pool (MySQLConnectionPool): MySQL 连接池。
  211. Returns:
  212. int: 去重后发现的商户总数。
  213. """
  214. seen = set()
  215. fetch_from = 1
  216. while fetch_from <= MAX_PAGES * PAGE_SIZE:
  217. try:
  218. data = get_forsale_page(log, fetch_from, PAGE_SIZE)
  219. except Exception as e:
  220. log.error(f"在售主列表 fetch_from={fetch_from} 请求失败: {e}")
  221. break
  222. if not data:
  223. break
  224. good_list = data.get("goodList", [])
  225. if not good_list:
  226. log.info(f"在售主列表 fetch_from={fetch_from} 无数据,停止翻页")
  227. break
  228. save_shops(log, good_list, sql_pool, seen)
  229. log.info(f"在售主列表 fetch_from={fetch_from} 完成,本页 {len(good_list)} 条,累计商户 {len(seen)}")
  230. if data.get("isFetchEnd") or len(good_list) < PAGE_SIZE:
  231. break
  232. fetch_from += PAGE_SIZE
  233. return len(seen)
  234. # ==================== 二、商品采集(商户已售列表) ====================
  235. def get_merchant_sold_page(log, merchant_code: str, page_index: int, page_size: int = PAGE_SIZE):
  236. """获取某商户「已售/已完成」列表的一页(merchant?tp=2,公开接口免签名)。
  237. Args:
  238. log: 日志对象。
  239. merchant_code (str): 商户编码(kj_shop_record.shop_id,如 MCT7550827)。
  240. page_index (int): 页码,从 1 开始。
  241. page_size (int, optional): 每页条数。Defaults to PAGE_SIZE。
  242. Returns:
  243. dict | None: 响应 JSON(含 list / totalPage);失败返回 None。
  244. """
  245. return kj_request(
  246. log, f"merchant/1/goodlist/{merchant_code}", base=PAGE_BASE,
  247. params={"pageIndex": str(page_index), "pageSize": str(page_size), "tp": "2"},
  248. need_auth=False,
  249. )
  250. def get_good_detail(log, good_code: str) -> dict:
  251. """获取商品详情,返回 detail.good(含时间 / 商户 publisher / 规格)。
  252. Args:
  253. log: 日志对象。
  254. good_code (str): 商品编码,如 ZC7653941。
  255. Returns:
  256. dict: detail 中的 good 字典;无数据时返回空字典。
  257. """
  258. j = kj_request(
  259. log, f"good/{good_code}/1/detail", base=PAGE_BASE,
  260. params={"referer": "MerchantList"}, need_auth=True,
  261. )
  262. return (j or {}).get("good", {}) or {}
  263. def get_video(log, good_code: str, play_code: str) -> str | None:
  264. """获取商品「拆卡回放」视频地址。
  265. play_code 取自 detail.good.broadcast.playCode(为空表示回放未就绪);body 的 sign 由
  266. kj_auth.video_sign 生成,videoPlay 接口本身不需要 Authorization。响应的 media_url 即视频地址。
  267. Args:
  268. log: 日志对象。
  269. good_code (str): 商品编码。
  270. play_code (str): 回放播放码,来自 detail.good.broadcast.playCode。
  271. Returns:
  272. str | None: 回放视频 URL;play_code 为空或接口无视频时返回 None。
  273. """
  274. if not play_code:
  275. return None
  276. ts = int(time.time())
  277. body = {"playCode": play_code, "sign": kj_auth.video_sign(ts, good_code, play_code), "ts": ts}
  278. j = kj_request(log, f"good/videoPlay/{good_code}", method="POST", base=SERVER_BASE,
  279. json_body=body, need_auth=False, extra_headers={"Content-Type": "application/json"})
  280. print(j)
  281. return (j or {}).get("media_url") if j and j.get("code") == 0 else None
  282. def build_product(good_item: dict, detail_good: dict, merchant_code: str, merchant_name: str, video_url: str) -> dict:
  283. """把商户列表项 + 商品详情组装成 kj_product_record 一行。
  284. 卡集为拼团模型,zc 表中的直播 / 存储 / 多价字段无对应,统一置 None。
  285. Args:
  286. good_item (dict): merchant?tp=2 列表里的商品项(goodCode/title/pic/price/totalNum/currentNum/status)。
  287. detail_good (dict): detail.good,补 startAt/overAt/state/规格;detail 失败时为空字典。
  288. merchant_code (str): 商户编码。
  289. merchant_name (str): 商户名。
  290. video_url (str): 回放视频 URL。
  291. Returns:
  292. dict: 与 kj_product_record 列对应的数据字典。
  293. """
  294. # 时间戳转文本:ts_to_dt(detail_good.get("startAt"))
  295. imgs = detail_good.get("pic", {}).get("carousel", [])
  296. imgs = ','.join(imgs) if imgs else None
  297. # 拼团系列 / 拼团时间(天):从 detail.good.desc 解码提取
  298. group_series, group_days = parse_desc(detail_good.get("desc", ""))
  299. # print(good_item)
  300. # print('------------------------------')
  301. # print(detail_good)
  302. row = {
  303. "shop_id": merchant_code, # 商户编码
  304. "shop_name": merchant_name, # 商户名
  305. "pid": good_item.get("goodCode"), #
  306. "title": good_item.get("title"), # 标题
  307. "price": good_item.get("price"), # 价格
  308. "total_num": good_item.get("totalNum"), # 总数量
  309. "current_num": good_item.get("currentNum"), # 当前数量
  310. "status": good_item.get("status"), # 状态
  311. "imgs": imgs, # 图片链接, ','分割
  312. "start_at": ts_to_dt(detail_good.get("startAt")), # 开始时间
  313. "over_at": ts_to_dt(detail_good.get("overAt")), # 结束时间
  314. "spec_name": detail_good.get("spec", {}).get("name"), # 规格配置
  315. "spec_content": detail_good.get("spec", {}).get("content"), # 产品规格
  316. "group_series": group_series, # 拼团系列,如 "27-28"
  317. "group_days": group_days, # 拼团时间(天),如 15
  318. # "publisher_sale": detail_good.get("publisher", {}).get("sale"), # 在售
  319. # "publisher_fans": detail_good.get("publisher", {}).get("fans"), # 出售者粉丝
  320. "live_start_at": ts_to_dt((detail_good.get("broadcast") or {}).get("startAt")),
  321. "video_url": video_url,
  322. }
  323. return row
  324. def get_shop_stop_pid(shop_id: str, sql_pool) -> str | None:
  325. """取该商户已入库商品里 start_at 最新那条的 pid,作为增量翻页的停止线。
  326. merchant?tp=2 结果按 start_at 倒序(最新在前),daily 翻页碰到这个 pid 即可早停:
  327. 它及其之后的商品都是上一轮已采过的。依赖 (shop_id, start_at) 复合索引,单值查询很快。
  328. Args:
  329. shop_id (str): 商户编码。
  330. sql_pool (MySQLConnectionPool): MySQL 连接池。
  331. Returns:
  332. str | None: 该商户最新商品的 pid(goodCode);该商户暂无记录时返回 None(触发全量采集)。
  333. """
  334. row = sql_pool.select_one(
  335. "SELECT pid FROM kj_product_record WHERE shop_id = %s ORDER BY start_at DESC LIMIT 1",
  336. (shop_id,),
  337. )
  338. return row[0] if row else None
  339. def get_sold_list(log, shop_id: str, shop_name: str, sql_pool, incremental: bool = True) -> int:
  340. """遍历某商户已售列表翻页,逐商品补详情后写入 kj_product_record。
  341. Args:
  342. log: 日志对象。
  343. shop_id (str): 商户编码(shop_id)。
  344. shop_name (str): 商户名。
  345. sql_pool (MySQLConnectionPool): MySQL 连接池。
  346. incremental (bool, optional): True(daily) 按 start_at 最新 pid 早停,只采新商品;
  347. False(history) 全量深翻所有页。Defaults to True。
  348. Returns:
  349. int: 本商户写入的商品数。
  350. """
  351. page_index = 1
  352. saved = 0
  353. stopped = False
  354. # 增量停止线:该商户上次采到的最新商品 pid(接口按 start_at 新在前,翻页碰到它即早停)
  355. stop_pid = get_shop_stop_pid(shop_id, sql_pool) if incremental else None
  356. while page_index <= MAX_PAGES:
  357. try:
  358. data = get_merchant_sold_page(log, shop_id, page_index)
  359. except Exception as e:
  360. log.error(f"商户 {shop_id} 已售列表第 {page_index} 页请求失败: {e}")
  361. break
  362. if not data:
  363. log.info(f"商户 {shop_id} 已售列表第 {page_index} 页无数据,停止翻页")
  364. break
  365. # 商户无效(已注销 / 不存在):merchant?tp=2 返回 code=1、msg='无效商家'。
  366. # 首页即无效 → 标记 is_deleted=1;之后 kj_main 的「WHERE is_deleted = 0」会自动跳过它,不再浪费请求。
  367. if data.get("code") != 0:
  368. if page_index == 1:
  369. log.info(f"商户 {shop_id} 无效({data.get('msg')}),标记 is_deleted=1")
  370. sql_pool.update_one("UPDATE kj_shop_record SET is_deleted = 1 WHERE shop_id = %s", (shop_id,))
  371. break
  372. good_list = data.get("list", [])
  373. if not good_list:
  374. log.info(f"商户 {shop_id} 已售列表第 {page_index} 页无数据,停止翻页")
  375. break
  376. total_page = data.get("totalPage", 1)
  377. batch = []
  378. for gi in good_list:
  379. code = gi.get("goodCode")
  380. # 增量早停:碰到上次采到的最新商品,它及之后都是已采过的旧数据
  381. if stop_pid and code == stop_pid:
  382. log.info(f"商户 {shop_id} 碰到停止线 pid={stop_pid},增量早停")
  383. stopped = True
  384. break
  385. detail_good = {}
  386. if FETCH_DETAIL:
  387. try:
  388. detail_good = get_good_detail(log, code)
  389. except Exception as e:
  390. log.error(f"商品 {code} detail 请求失败: {e}")
  391. # 如需回放视频地址(media_url):playCode 在 detail_good["broadcast"]["playCode"](空=未就绪)
  392. # 首次拿不到不影响主数据 —— refill_videos 会在拼团完成后兜底补采
  393. try:
  394. video_url = get_video(log, code, (detail_good.get("broadcast") or {}).get("playCode"))
  395. except Exception as e:
  396. log.error(f"商品 {code} 视频获取失败: {e}")
  397. video_url = None
  398. row = build_product(gi, detail_good, shop_id, shop_name, video_url)
  399. # print(row)
  400. if row:
  401. batch.append(row)
  402. if batch:
  403. # 已存在(pid 唯一)则跳过:已完成商品为终态,无需覆盖
  404. sql_pool.insert_many(table="kj_product_record", data_list=batch, ignore=True)
  405. saved += len(batch)
  406. if stopped:
  407. log.info(f"商户 {shop_id} 增量早停,停止翻页")
  408. break
  409. log.info(f"商户 {shop_id} 已售列表第 {page_index}/{total_page} 页完成,本页 {len(good_list)} 商品")
  410. page_index += 1
  411. if page_index > total_page:
  412. log.info(f"商户 {shop_id} 已售列表翻页完成")
  413. break
  414. return saved
  415. # ==================== 三、玩家采集(中奖/参与名单) ====================
  416. def get_player_page(log, good_code: str, fetch_from: int, fetch_size: int = PLAYER_PAGE_SIZE):
  417. """获取某商品玩家名单的一页 good/{code}/result/merge。
  418. Args:
  419. log: 日志对象。
  420. good_code (str): 商品编码。
  421. fetch_from (int): 游标起始位置(首页为 1,逐页按 fetch_size 递增)。
  422. fetch_size (int, optional): 每页条数。Defaults to PLAYER_PAGE_SIZE。
  423. Returns:
  424. dict | None: 响应 JSON(含 code / list / isFetchEnd);失败返回 None。
  425. """
  426. return kj_request(
  427. log, f"good/{good_code}/result/merge", base=PAGE_BASE,
  428. params={"fetchFrom": str(fetch_from), "fetchSize": str(fetch_size), "q": ""},
  429. need_auth=True,
  430. )
  431. def save_players(log, good_code: str, player_list: list, sql_pool) -> int:
  432. """解析玩家名单并写入 kj_player_record。
  433. 匿名玩家(userId=0、userName 为空)用 anonymousCode 兜底作为标识。
  434. Args:
  435. log: 日志对象。
  436. good_code (str): 商品编码(写入 pid 列)。
  437. player_list (list[dict]): result/merge 的 list,每项含 userId/userName/total/anonymousCode。
  438. sql_pool (MySQLConnectionPool): MySQL 连接池。
  439. Returns:
  440. int: 本次写入的玩家记录数。
  441. """
  442. log.info(f"开始保存商品 {good_code} 的玩家名单")
  443. info_list = []
  444. for item in player_list:
  445. anon = item.get("anonymousCode")
  446. user_id = item.get("userId")
  447. user_name = item.get("userName")
  448. if not user_name: # 匿名玩家 userName 为空
  449. user_name = f"匿名_{anon}" if anon else None
  450. if not user_id: # 匿名玩家 userId=0,用匿名码兜底
  451. user_id = anon
  452. info_list.append({
  453. "pid": good_code, # 商品编码,来自函数参数(item 里没有商品标识)
  454. "give_number": item.get("total"), # 该玩家份数:卡集字段是 total
  455. "user_id": str(user_id) if user_id is not None else None,
  456. "user_name": user_name,
  457. })
  458. if info_list:
  459. sql_pool.insert_many(table="kj_player_record", data_list=info_list)
  460. return len(info_list)
  461. def get_player_list(log, good_code: str, sql_pool) -> bool:
  462. """遍历某商品玩家名单翻页并写库。
  463. Args:
  464. log: 日志对象。
  465. good_code (str): 商品编码。
  466. sql_pool (MySQLConnectionPool): MySQL 连接池。
  467. Returns:
  468. bool: True 表示抓到玩家数据,False 表示无数据(如拼团未完成)。
  469. """
  470. fetch_from = 1
  471. has_data = False
  472. while fetch_from <= MAX_PAGES * PLAYER_PAGE_SIZE:
  473. try:
  474. data = get_player_page(log, good_code, fetch_from, PLAYER_PAGE_SIZE)
  475. except Exception as e:
  476. log.error(f"商品 {good_code} 玩家名单 fetch_from={fetch_from} 请求失败: {e}")
  477. break
  478. if not data:
  479. log.info(f"商品 {good_code} 玩家名单 fetch_from={fetch_from} 无数据,停止翻页")
  480. break
  481. if data.get("code") != 0:
  482. # code=1 常见于"拼团未完成的商品",视为暂无玩家
  483. log.info(f"商品 {good_code} 暂无玩家: {data.get('msg')}")
  484. break
  485. plist = data.get("list", [])
  486. if not plist:
  487. log.info(f"商品 {good_code} 玩家名单翻页完成")
  488. break
  489. has_data = True
  490. save_players(log, good_code, plist, sql_pool)
  491. if data.get("isFetchEnd") or len(plist) < PLAYER_PAGE_SIZE:
  492. log.info(f"商品 {good_code} 玩家名单翻页完成")
  493. break
  494. fetch_from += PLAYER_PAGE_SIZE
  495. return has_data
  496. # ==================== 四、视频补采 ====================
  497. def refill_videos(log, sql_pool) -> int:
  498. """补采视频地址:对拼团已完成(player_state=1)但 video_url 仍空的商品,
  499. 重新拉 detail 取 broadcast.playCode → get_video → 更新 video_url。
  500. 覆盖场景:首次入库时该商品还处于「即将拆卡 / 正在拆卡」等中间状态,
  501. broadcast.playCode 为空、视频拿不到;拼团完成、回放就绪后由本函数补回。
  502. playCode 仍空则跳过,留到下一轮再试,直到 video_url 有值。
  503. Args:
  504. log: 日志对象。
  505. sql_pool (MySQLConnectionPool): MySQL 连接池。
  506. Returns:
  507. int: 本轮成功补采的视频数。
  508. """
  509. rows = sql_pool.select_all(
  510. "SELECT pid FROM kj_product_record WHERE video_url IS NULL AND player_state = 1"
  511. )
  512. pids = [r[0] for r in rows] if rows else []
  513. log.info(f"待补视频商品 {len(pids)} 个")
  514. filled = 0
  515. for pid in pids:
  516. try:
  517. detail = get_good_detail(log, pid)
  518. play_code = (detail.get("broadcast") or {}).get("playCode")
  519. if not play_code:
  520. # 回放仍未就绪(正在拆卡 / 即将拆卡等),留到下一轮
  521. continue
  522. url = get_video(log, pid, play_code)
  523. if url:
  524. sql_pool.update_one(
  525. "UPDATE kj_product_record SET video_url = %s WHERE pid = %s",
  526. (url, pid),
  527. )
  528. filled += 1
  529. log.info(f"商品 {pid} 视频已补采 → {url}")
  530. except Exception as e:
  531. log.error(f"商品 {pid} 视频补采失败: {e}")
  532. return filled
  533. # ==================== 主流程 ====================
  534. @retry(stop=stop_after_attempt(100), wait=wait_fixed(3600), after=after_log)
  535. def kj_main(log):
  536. """卡集每日采集主函数:商户发现 → 商品采集 → 玩家采集。
  537. Args:
  538. log: 日志对象。
  539. Raises:
  540. RuntimeError: 数据库连接池异常时抛出以触发重试。
  541. """
  542. log.info(f"开始运行 {sys._getframe().f_code.co_name} 卡集采集任务" + "." * 40)
  543. sql_pool = MySQLConnectionPool(log=log)
  544. if not sql_pool.check_pool_health():
  545. log.error("数据库连接池异常")
  546. raise RuntimeError("数据库连接池异常")
  547. try:
  548. # 1) 商户发现
  549. try:
  550. n = get_shop_list(log, sql_pool)
  551. log.info(f"商户发现完成,去重商户 {n} 个")
  552. except Exception as e:
  553. log.error(f"get_shop_list error: {e}")
  554. time.sleep(5)
  555. # 2) 商品采集:遍历库内所有商户的已售列表
  556. try:
  557. shop_rows = sql_pool.select_all("SELECT shop_id, shop_name FROM kj_shop_record WHERE is_deleted = 0")
  558. log.info(f"待采集商户 {len(shop_rows)} 个")
  559. for shop_id, shop_name in shop_rows:
  560. try:
  561. cnt = get_sold_list(log, shop_id, shop_name, sql_pool)
  562. log.info(f"商户 {shop_id} {shop_name} 商品采集完成,写入 {cnt} 个")
  563. except Exception as e:
  564. log.error(f"get_sold_list error(商户 {shop_id}): {e}")
  565. except Exception as e:
  566. log.error(f"iterate_shop_list error: {e}")
  567. time.sleep(5)
  568. # 3) 玩家采集:遍历尚未成功采集玩家的商品
  569. try:
  570. prod_rows = sql_pool.select_all("SELECT pid FROM kj_product_record WHERE player_state != 1")
  571. pids = [row[0] for row in prod_rows] if prod_rows else []
  572. log.info(f"待采集玩家的商品 {len(pids)} 个")
  573. for pid in pids:
  574. try:
  575. # 先置 1 表示开始采集(对齐 zc 断点标记)
  576. sql_pool.update_one("UPDATE kj_product_record SET player_state = 1 WHERE pid = %s", (pid,))
  577. has_data = get_player_list(log, pid, sql_pool)
  578. if not has_data:
  579. # 无玩家(如拼团未完成)置 2,下轮仍会重试
  580. sql_pool.update_one("UPDATE kj_product_record SET player_state = 2 WHERE pid = %s", (pid,))
  581. except Exception as pid_error:
  582. log.error(f"商品 {pid} 玩家采集失败: {pid_error}")
  583. try:
  584. sql_pool.update_one("UPDATE kj_product_record SET player_state = 3 WHERE pid = %s", (pid,))
  585. except Exception as update_error:
  586. log.error(f"更新商品 {pid} 状态失败: {update_error}")
  587. except Exception as e:
  588. log.error(f"iterate_player_list error: {e}")
  589. # 4) 视频补采:玩家采集之后,对拼团已完成(player_state=1)但 video_url 仍空的商品重取视频
  590. try:
  591. n = refill_videos(log, sql_pool)
  592. log.info(f"视频补采完成,本轮 {n} 条")
  593. except Exception as e:
  594. log.error(f"refill_videos error: {e}")
  595. except Exception as e:
  596. log.error(f"{sys._getframe().f_code.co_name} error: {e}")
  597. finally:
  598. log.info(f"卡集采集 {sys._getframe().f_code.co_name} 运行结束,等待下一轮" + "." * 20)
  599. def schedule_task():
  600. """定时任务入口:每天 00:01 运行一次 kj_main。"""
  601. # 立即运行一次(调试时取消注释)
  602. # kj_main(log=logger)
  603. schedule.every().day.at("00:01").do(kj_main, log=logger)
  604. while True:
  605. schedule.run_pending()
  606. time.sleep(1)
  607. if __name__ == "__main__":
  608. # kj_main(logger)
  609. schedule_task()