zc_new_daily_spider.py 26 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690
  1. # -*- coding: utf-8 -*-
  2. # Author : Charley
  3. # Python : 3.10.8
  4. # Date : 2026/2/27 11:22
  5. import random
  6. import time
  7. import inspect
  8. import requests
  9. import schedule
  10. import user_agent
  11. from loguru import logger
  12. from crypto_utils import CryptoHelper
  13. from mysql_pool import MySQLConnectionPool
  14. # 2026/08/06 停用板块发现:get_board_shop_list 的旧 cateId 已随版本升级失效,
  15. # 改用「全部」接口(get_all_shop_list)覆盖全部在售店铺,故移除该导入(文件保留备用)。
  16. from tenacity import retry, stop_after_attempt, wait_fixed
  17. logger.remove()
  18. logger.add("./logs/{time:YYYYMMDD}.log", encoding='utf-8', rotation="00:00",
  19. format="[{time:YYYY-MM-DD HH:mm:ss.SSS}] {level} {message}",
  20. level="DEBUG", retention="7 day")
  21. # 基础配置
  22. BASE_URL = "https://cashier.yqszpay.com"
  23. PAGE_SIZE = 10
  24. headers = {
  25. "User-Agent": user_agent.generate_user_agent(os="android"), # 设置为安卓模拟器
  26. "Connection": "Keep-Alive",
  27. "Accept-Encoding": "gzip",
  28. "Content-Type": "application/json",
  29. "channelNo": "88888888",
  30. "pageSize": str(PAGE_SIZE),
  31. # "pageNum": 1,
  32. # 2026/08/06 版本失效修复:服务端升级后按 version 头做灰度校验,旧值 1.9.9.82537
  33. # 会被静默降级为返回空 rows(HTTP 200 不报错、数据悄悄变空)。改为抓包实测的当前
  34. # 真实版本 1.3.09 后 getActList / getActEndList 均恢复正常返回。
  35. "version": "1.3.14"
  36. }
  37. def after_log(retry_state):
  38. """
  39. retry 回调
  40. :param retry_state: RetryCallState 对象
  41. """
  42. # 检查 args 是否存在且不为空
  43. if retry_state.args and len(retry_state.args) > 0:
  44. log = retry_state.args[0] # 获取传入的 logger
  45. else:
  46. log = logger # 使用全局 logger
  47. if retry_state.outcome.failed:
  48. log.warning(
  49. f"Function '{retry_state.fn.__name__}', Attempt {retry_state.attempt_number} Times")
  50. else:
  51. log.info(f"Function '{retry_state.fn.__name__}', Attempt {retry_state.attempt_number} succeeded")
  52. @retry(stop=stop_after_attempt(5), wait=wait_fixed(1), after=after_log)
  53. def get_proxys(log):
  54. """
  55. 获取代理配置
  56. :param log: 日志对象
  57. :return: 代理字典
  58. """
  59. tunnel = "x371.kdltps.com:15818"
  60. kdl_username = "t13753103189895"
  61. kdl_password = "o0yefv6z"
  62. try:
  63. proxies = {
  64. "http": "http://%(user)s:%(pwd)s@%(proxy)s/" % {"user": kdl_username, "pwd": kdl_password, "proxy": tunnel},
  65. "https": "http://%(user)s:%(pwd)s@%(proxy)s/" % {"user": kdl_username, "pwd": kdl_password, "proxy": tunnel}
  66. }
  67. return proxies
  68. except Exception as e:
  69. log.error(f"Error getting proxy: {e}")
  70. raise e
  71. @retry(stop=stop_after_attempt(5), wait=wait_fixed(1), after=after_log)
  72. def make_encrypted_post_request(log, url: str, request_data: dict, extra_headers: dict = None):
  73. """
  74. 通用加密POST请求函数(带重试机制)
  75. :param log: 日志对象
  76. :param url: 请求URL
  77. :param request_data: 请求数据字典(会被加密)
  78. :param extra_headers: 额外的请求头
  79. :return: 解密后的响应数据,失败返回None
  80. """
  81. request_headers = headers.copy()
  82. if extra_headers:
  83. request_headers.update(extra_headers)
  84. log.debug(f"Request URL: {url}, Data: {request_data}")
  85. encrypted_body = CryptoHelper.encrypt_request_data(request_data)
  86. # print(request_headers)
  87. # response = requests.post(url, headers=request_headers, json=encrypted_body, timeout=22, proxies=get_proxys(log))
  88. response = requests.post(url, headers=request_headers, json=encrypted_body, timeout=(5, 30),
  89. proxies=get_proxys(log))
  90. # response = requests.post(url, headers=request_headers, json=encrypted_body, timeout=(5, 30))
  91. # response.raise_for_status()
  92. if response.status_code == 200:
  93. response_json = response.json()
  94. # log.debug(f"Raw response: {response_json}")
  95. # print(response_json)
  96. if 'data' in response_json:
  97. decrypted = CryptoHelper.decrypt_response_data(response_json)
  98. # log.debug(f"Decrypted response: {decrypted}")
  99. # print(decrypted)
  100. return decrypted
  101. return response_json
  102. else:
  103. log.error(f"请求失败: {response.status_code}, Response: {response.text}")
  104. return None
  105. def get_shop_single_page(log, page_num, page_size=PAGE_SIZE):
  106. """
  107. 获取商户列表(支持翻页)
  108. :param log: 日志对象
  109. :param page_num: 页码
  110. :param page_size: 每页条数
  111. """
  112. log.debug(f"Getting shop list, page: {page_num}")
  113. url = f"{BASE_URL}/zc-api/merchant/getMerMyList"
  114. request_data = {'pageNum': page_num, 'pageSize': page_size}
  115. try:
  116. resp = make_encrypted_post_request(log, url, request_data, extra_headers={"pageNum": str(page_num)})
  117. except Exception as e:
  118. log.error(f"Error getting shop list: {e}")
  119. resp = None
  120. return resp
  121. # 「全部」板块的分类ID。来源:抓包数据.txt 里「全部」两条 getActList 请求体(AES 解密后得到),
  122. # 2026/08/06 版本升级后该值为 "0,6,3"。此接口为公开接口,不校验 Authorization,无需 token。
  123. ALL_SHOP_CATE_ID = "0,6,3"
  124. def get_all_shop_single_page(log, page_num, page_size=PAGE_SIZE):
  125. """获取「全部」板块的活动列表单页(用于发现全部在售店铺)。
  126. 走 getActList「全部」板块(cateId=ALL_SHOP_CATE_ID),从每条活动里解析出
  127. merNo / merName,覆盖当前所有有进行中活动的店铺。为公开接口,无需 token。
  128. Args:
  129. log: 日志对象,从调用方透传。
  130. page_num (int): 页码,从 1 开始。
  131. page_size (int, optional): 每页条数。Defaults to PAGE_SIZE。
  132. Returns:
  133. dict | None: 解密后的响应字典(含 rows);请求失败时返回 None。
  134. """
  135. log.debug(f"Getting all-shop list, page: {page_num}")
  136. url = f"{BASE_URL}/zc-api/act/actProduct/getActList"
  137. # 请求体字段与抓包解密结果一一对应,仅 cateId 固定为「全部」板块:
  138. # finish / loading -> 前端分页状态位,抓包固定 false / true,照搬即可
  139. # actStatus=1 -> 只取进行中的活动
  140. request_data = {
  141. 'cateId': ALL_SHOP_CATE_ID,
  142. 'pageSize': page_size,
  143. 'pageNum': page_num,
  144. 'finish': False,
  145. 'loading': True,
  146. 'actStatus': 1
  147. }
  148. try:
  149. resp = make_encrypted_post_request(log, url, request_data, extra_headers={"pageNum": str(page_num)})
  150. except Exception as e:
  151. log.error(f"Error getting all-shop list, page {page_num}: {e}")
  152. resp = None
  153. return resp
  154. def parse_all_shop_data(log, items, sql_pool, seen):
  155. """解析「全部」接口活动里的店铺并写库(INSERT IGNORE,不覆盖存量)。
  156. 「全部」接口只返回 merNo / merName,拿不到 sold_number(销量) / fans(粉丝)。因此对
  157. 新发现的店铺 sold_number 写死哨兵值 1,让其能被每日任务「WHERE sold_number != 0」
  158. 选中进入商品采集管道;用 INSERT IGNORE 跳过已存在的店铺,避免覆盖 get_shop_list
  159. (getMerMyList)回填的真实 sold_number / fans。
  160. Args:
  161. log: 日志对象,从调用方透传。
  162. items (list[dict]): getActList「全部」板块返回的活动列表(每条含 merNo / merName)。
  163. sql_pool (MySQLConnectionPool): MySQL 连接池。
  164. seen (set): 跨页去重的 shop_id 集合,避免同一店铺重复写库。
  165. Returns:
  166. int: 本次实际写入的店铺数量(已去重)。
  167. """
  168. info_list = []
  169. for item in items:
  170. shop_id = item.get('merNo')
  171. shop_name = item.get('merName')
  172. # 同一店铺会出现在多条活动里,跨页去重以减少无谓写库
  173. if shop_id in seen:
  174. continue
  175. seen.add(shop_id)
  176. # sold_number 写死哨兵值 1:仅作「需采集」标记(板块拿不到真实销量),配合 INSERT IGNORE
  177. # 只对新插入店铺生效,已存在店铺整行被跳过;后续 get_shop_list 会回填真实值
  178. info_list.append({'shop_id': shop_id, 'shop_name': shop_name, 'sold_number': 1})
  179. if not info_list:
  180. return 0
  181. sql_pool.insert_many(table='zc_shop_record', data_list=info_list, ignore=True)
  182. return len(info_list)
  183. def get_all_shop_list(log, sql_pool):
  184. """「全部」接口店铺发现入口:翻页遍历全部在售店铺并写库。
  185. 店铺主来源,替代原来覆盖不全的 getMerMyList「我的商户列表」。翻页采用
  186. 「空页或不足整页即停」的可靠判断,受 MAX 页数保护。
  187. Args:
  188. log: 日志对象,从调用方透传。
  189. sql_pool (MySQLConnectionPool): MySQL 连接池。
  190. Returns:
  191. int: 去重后合计写入的店铺数量。
  192. """
  193. log.info("开始「全部」接口店铺发现......")
  194. page_num = 1
  195. max_pages = 100
  196. seen = set()
  197. while page_num <= max_pages:
  198. result = get_all_shop_single_page(log, page_num, PAGE_SIZE)
  199. if result is None:
  200. log.error(f"「全部」第 {page_num} 页请求失败,停止翻页")
  201. break
  202. data_list = result.get('rows', [])
  203. if len(data_list) == 0:
  204. log.info(f"「全部」第 {page_num} 页无数据,停止翻页")
  205. break
  206. parse_all_shop_data(log, data_list, sql_pool, seen)
  207. log.info(f"「全部」第 {page_num} 页完成,本页活动数: {len(data_list)},累计去重店铺: {len(seen)}")
  208. # 不足整页说明已到末页,停止翻页
  209. if len(data_list) < PAGE_SIZE:
  210. log.info(f"「全部」第 {page_num} 页不足整页,已到末页,停止翻页")
  211. break
  212. page_num += 1
  213. log.info(f"「全部」接口店铺发现结束,去重店铺数: {len(seen)}")
  214. return len(seen)
  215. def get_sold_single_page(log, mer_no, page_num):
  216. """
  217. 获取商品列表(支持翻页)
  218. :param log: 日志对象
  219. :param mer_no: 商户编号
  220. :param page_num: 页码
  221. """
  222. log.info(f"Getting sold items for mer_no: {mer_no}, page: {page_num}")
  223. # url = f"{BASE_URL}/zc-api/act/actProduct/getActList"
  224. # url = "https://cashier.yqszpay.com/zc-api/act/actProduct/getActEndList" 20260407接口变化 修改
  225. url = f"{BASE_URL}/zc-api/act/actProduct/getActEndList"
  226. request_data = {
  227. 'merNo': mer_no,
  228. 'pageNum': page_num,
  229. 'pageSize': PAGE_SIZE,
  230. 'queryType': 1
  231. }
  232. return make_encrypted_post_request(log, url, request_data, extra_headers={"pageNum": str(page_num)})
  233. def get_player_single_page(log, act_id, token, page_num, page_size=PAGE_SIZE):
  234. """
  235. 获取玩家列表(支持翻页)
  236. :param log: 日志对象
  237. :param act_id: 活动ID
  238. :param token: Authorization token
  239. :param page_num: 页码
  240. :param page_size: 每页条数
  241. """
  242. log.debug(f"Getting player list for act_id: {act_id}, page: {page_num}")
  243. url = f"{BASE_URL}/zc-api/act/actOrder/getActOrderPublicDetails"
  244. request_data = {'actId': act_id, 'pageNum': page_num, 'pageSize': page_size}
  245. return make_encrypted_post_request(
  246. log, url, request_data,
  247. extra_headers={"Authorization": token, "pageNum": str(page_num)}
  248. )
  249. def parse_shop_data(log, items, sql_pool):
  250. """
  251. 解析商户数据
  252. :param log: 日志对象
  253. :param items: 商户列表
  254. :param sql_pool: MySQL连接池
  255. :return: 解析后的数据列表
  256. """
  257. log.debug(f"Parsing shop data...........")
  258. info_list = []
  259. for item in items:
  260. # log.debug(f"Processing shop item: {item}")
  261. shop_id = item.get('merNo')
  262. shop_name = item.get('merName')
  263. sold_number = item.get('spell_number')
  264. # link_man = item.get('linkMan')
  265. # user_id = item.get('userId')
  266. fans = item.get('attentionNumber')
  267. data_dict = {
  268. 'shop_id': shop_id,
  269. 'shop_name': shop_name,
  270. 'sold_number': sold_number,
  271. 'fans': fans
  272. }
  273. log.debug(f"Parsed shop data: {data_dict}")
  274. info_list.append(data_dict)
  275. # 保存/更新 根据shop_id判断 是否存在,存在则更新,不存在则插入
  276. sql = "INSERT INTO zc_shop_record (shop_id, shop_name, sold_number, fans) VALUES (%s, %s, %s, %s) ON DUPLICATE KEY UPDATE shop_name=VALUES(shop_name), sold_number=VALUES(sold_number), fans=VALUES(fans)"
  277. # 将字典列表转换为元组列表
  278. args_list = [tuple(d.values()) for d in info_list]
  279. sql_pool.insert_many(query=sql, args_list=args_list)
  280. @retry(stop=stop_after_attempt(3), wait=wait_fixed(1), after=after_log)
  281. def get_video(log, token, pid):
  282. """
  283. 获取活动视频信息
  284. :param log: 日志对象
  285. :param token: Authorization token
  286. :param pid: 活动ID
  287. :return: (live_id, open_time, close_time, video_url)
  288. """
  289. url = "https://cashier.yqszpay.com/zc-api/live/actLive/getMerLiveInfo"
  290. request_data = {'actId': pid}
  291. log.debug(f"获取视频信息,actId: {pid}")
  292. resp_data = make_encrypted_post_request(
  293. log, url, request_data,
  294. extra_headers={"Authorization": token}
  295. )
  296. # log.debug(f"视频响应: {resp_data}")
  297. live_id = resp_data.get('live', {}).get('liveId')
  298. live_open_time = resp_data.get('live', {}).get('openTime')
  299. live_close_time = resp_data.get('live', {}).get('closeTime')
  300. video_url = resp_data.get('live', {}).get('videoUrl')
  301. return live_id, live_open_time, live_close_time, video_url
  302. def parse_sold_data(log, token, items, sql_pool, shop_name):
  303. """
  304. 解析商品数据
  305. :param log: 日志对象
  306. :param token: Authorization token
  307. :param items: 商品列表
  308. :param sql_pool: MySQL连接池
  309. :param shop_name: 商户名称
  310. :return: 解析后的数据列表
  311. """
  312. info_list = []
  313. for item in items:
  314. # log.debug(f"Processing sold item: {item}")
  315. shop_id = item.get('merNo') # 商户编号
  316. pid = item.get('id')
  317. act_day = item.get('actDay') # 活动天数
  318. act_logo = item.get('actLogo')
  319. act_name = item.get('actName') # 活动名称
  320. act_no = item.get('actNo') # 活动编号
  321. act_status = item.get('actStatus') # 活动状态
  322. startDate = item.get('startDate') # 开始时间
  323. endDate = item.get('complete_date') # 结束时间
  324. storageId = item.get('storageId') # 存储ID
  325. storageName = item.get('storageName') # 存储名称
  326. unitPrice = item.get('unitPrice') # 单价
  327. sumPrice = item.get('sumPrice') # 总价
  328. reality_price = item.get('realityPrice') # 实际价格
  329. packageNumber = item.get('packageNumber') # 包配置
  330. schedule_ = item.get('schedule') # 库存
  331. try:
  332. live_id, live_open_time, live_close_time, video_url = get_video(log, token, pid)
  333. except Exception as e:
  334. log.error(f"Error getting video info for pid {pid}: {e}")
  335. live_id, live_open_time, live_close_time, video_url = None, None, None, None
  336. data_dict = {
  337. 'shop_id': shop_id,
  338. 'shop_name': shop_name,
  339. 'pid': pid,
  340. 'act_day': act_day,
  341. 'act_img': act_logo,
  342. 'act_name': act_name,
  343. 'act_no': act_no,
  344. 'act_status': act_status,
  345. 'start_date': startDate,
  346. 'end_date': endDate,
  347. 'storage_id': storageId,
  348. 'storage_name': storageName,
  349. 'unit_price': unitPrice,
  350. 'sum_price': sumPrice,
  351. 'reality_price': reality_price,
  352. 'package_number': packageNumber,
  353. 'schedule': schedule_,
  354. 'live_id': live_id,
  355. 'live_open_time': live_open_time,
  356. 'live_close_time': live_close_time,
  357. 'video_url': video_url
  358. }
  359. # log.debug(f"Parsed sold data: {data_dict}")
  360. # { 'live_close_time': None, 'video_url': None}
  361. info_list.append(data_dict)
  362. # 保存数据
  363. sql_pool.insert_many(table='zc_product_record', data_list=info_list, ignore=True)
  364. def parse_player_data(log, items, sql_pool):
  365. """
  366. 解析玩家数据
  367. :param log: 日志对象
  368. :param items: 玩家列表
  369. :param sql_pool: MySQL连接池
  370. :return: 解析后的数据列表
  371. """
  372. log.debug(f"Parsing player data...........")
  373. info_list = []
  374. for item in items:
  375. # log.debug(f"Processing player item: {item}")
  376. pid = item.get('actId')
  377. give_number = item.get('giveNumber') # 份数
  378. user_id = item.get('userId')
  379. user_name = item.get('userName')
  380. data_dict = {
  381. 'pid': pid,
  382. 'give_number': give_number,
  383. 'user_id': user_id,
  384. 'user_name': user_name
  385. }
  386. # log.debug(f"Parsed player data: {data_dict}")
  387. info_list.append(data_dict)
  388. # 保存数据
  389. sql_pool.insert_many(table='zc_player_record', data_list=info_list, ignore=True)
  390. def get_shop_list(log, sql_pool):
  391. """
  392. 商户列表翻页生成器
  393. :param log: 日志对象
  394. :param sql_pool: MySQL连接池
  395. """
  396. page_num = 1
  397. max_pages = 100
  398. # 2026/08/06 翻页 bug 修复:原逻辑 total 初值为 0,却用「total is None」判断是否取总数,
  399. # 导致 total 永远停在 0,停止条件「(page_num-1)*PAGE_SIZE >= total」在第 1 页跑完后即
  400. # 0 >= 0 成立、立刻 break,只能采到第 1 页 10 个店铺。改为「空页 / 不足整页即停」的可靠判断。
  401. while page_num <= max_pages:
  402. result = get_shop_single_page(log, page_num, PAGE_SIZE)
  403. if result is None:
  404. log.error(f"第 {page_num} 页请求失败,停止翻页")
  405. break
  406. data_list = result.get('rows', [])
  407. # 空页直接停止,避免用空列表调用 insert_many 触发 ValueError
  408. if len(data_list) == 0:
  409. log.info(f"第 {page_num} 页无数据,停止翻页")
  410. break
  411. # getMerMyList 含真实 sold_number / fans,走 ON DUPLICATE KEY UPDATE 回填覆盖哨兵值
  412. parse_shop_data(log, data_list, sql_pool)
  413. log.info(f"第 {page_num} 页查询完成,本页条数: {len(data_list)}")
  414. # 不足整页说明已到末页,停止翻页
  415. if len(data_list) < PAGE_SIZE:
  416. log.info(f"第 {page_num} 页不足整页,已到末页,停止翻页")
  417. break
  418. page_num += 1
  419. def get_sold_list(log, shop_id, token, sql_pool, shop_name):
  420. """
  421. 商品列表翻页生成器
  422. :param log: 日志对象
  423. :param shop_id: shop_id
  424. :param token: Authorization token
  425. :param sql_pool: MySQL连接池
  426. :param shop_name: 商户名称
  427. """
  428. page_num = 1
  429. max_pages = 10
  430. while page_num <= max_pages:
  431. result = get_sold_single_page(log, shop_id, page_num)
  432. time.sleep(random.uniform(0.5, 1)) # 添加随机延迟,防止对目标服务器造成过大负载
  433. # print(result)
  434. if result is None:
  435. log.error(f"第 {page_num} 页请求失败,停止翻页")
  436. break
  437. data_list = result.get('rows', [])
  438. if not data_list:
  439. log.info(f"第 {page_num} 页无数据,停止翻页")
  440. # log.info(f'该店铺{shop_name}无数据, 修改为店铺注销状态..........')
  441. # # 更新店铺状态为注销状态
  442. # sql_pool.update_one("UPDATE zc_shop_record SET is_deleted = 1 WHERE shop_id = %s", (shop_id,))
  443. break
  444. parse_sold_data(log, token, data_list, sql_pool, shop_name)
  445. # 检查是否有数据
  446. if len(data_list) < 10:
  447. log.info(f"第 {page_num} 页无数据,停止翻页")
  448. break
  449. log.info(f"第 {page_num} 页查询完成,本页条数: {len(data_list)}")
  450. page_num += 1
  451. def get_player_list(log, act_id, token, sql_pool):
  452. """
  453. 玩家列表翻页生成器
  454. :param log: 日志对象
  455. :param act_id: 活动ID
  456. :param token: Authorization token
  457. :param sql_pool: MySQL连接池
  458. :return: has_data (True: 有数据, False: 无数据)
  459. """
  460. page_num = 1
  461. max_pages = 100
  462. has_data = False
  463. while page_num <= max_pages:
  464. result = get_player_single_page(log, act_id, token, page_num)
  465. if result is None:
  466. log.error(f"第 {page_num} 页请求失败,停止翻页")
  467. break
  468. data_list = result.get('rows', [])
  469. # 如果有数据才解析
  470. if len(data_list) > 0:
  471. has_data = True
  472. parse_player_data(log, data_list, sql_pool)
  473. # 检查是否有数据
  474. if len(data_list) < 10:
  475. log.info(f"第 {page_num} 页无数据,停止翻页")
  476. break
  477. log.info(f"第 {page_num} 页查询完成,本页条数: {len(data_list)}")
  478. page_num += 1
  479. return has_data
  480. @retry(stop=stop_after_attempt(100), wait=wait_fixed(3600), after=after_log)
  481. def zc_main(log):
  482. """
  483. 主函数
  484. :param log: logger对象
  485. """
  486. log.info(
  487. f'开始运行 {inspect.currentframe().f_code.co_name} 爬虫任务....................................................')
  488. # 配置 MySQL 连接池
  489. sql_pool = MySQLConnectionPool(log=log)
  490. if not sql_pool.check_pool_health():
  491. log.error("数据库连接池异常")
  492. raise RuntimeError("数据库连接池异常")
  493. try:
  494. # 获取 token
  495. token_row = sql_pool.select_one("SELECT token FROM zc_token WHERE id = 1")
  496. if not token_row:
  497. log.error("未查询到 token")
  498. return
  499. token = token_row[0]
  500. # player test
  501. # has_data = get_player_list(log, 1800, token, sql_pool)
  502. # 店铺主来源:「全部」接口发现全部在售店铺(替代覆盖不全的 getMerMyList)
  503. # 2026/08/06 停用板块发现 get_board_shop_list:其 5 个旧 cateId 已随版本升级失效,
  504. # 且「全部」接口(cateId=0,6,3)已覆盖所有在售店铺,板块发现冗余。
  505. try:
  506. get_all_shop_list(log, sql_pool)
  507. except Exception as e:
  508. log.error(f'get_all_shop_list error: {e}')
  509. # 补充:getMerMyList 回填「我的店铺」真实 sold_number / fans(「全部」接口拿不到这两个字段)
  510. try:
  511. get_shop_list(log, sql_pool)
  512. except Exception as e:
  513. log.error(f'iterate_shop_list error: {e}')
  514. time.sleep(5)
  515. # 获取sold data - 遍历所有商户
  516. try:
  517. # 从 shop 表查询所有 merNo 2026/5/9增加is_deleted判断是否注销店铺
  518. mer_no_rows = sql_pool.select_all(
  519. # "SELECT shop_id, shop_name FROM zc_shop_record WHERE is_deleted = 0 AND sold_number != 0"
  520. "SELECT shop_id, shop_name FROM zc_shop_record WHERE sold_number != 0"
  521. )
  522. log.info(f"查询到 {len(mer_no_rows)} 个商户编号: {mer_no_rows}")
  523. for shop_id, shop_name in mer_no_rows:
  524. log.info(f"开始爬取商户 {shop_id}, {shop_name} 的商品数据")
  525. # get_sold_list(log, shop_id, shop_name, token, sql_pool)
  526. get_sold_list(log, shop_id, token, sql_pool, shop_name)
  527. except Exception as e:
  528. log.error(f'get_sold_list error: {e}')
  529. time.sleep(5)
  530. # 获取player data - 遍历所有活动
  531. try:
  532. # 从 sold 表查询所有 actId
  533. act_id_rows = sql_pool.select_all("SELECT pid FROM zc_product_record WHERE player_state != 1")
  534. act_id_list = [row[0] for row in act_id_rows] if act_id_rows else []
  535. log.info(f"查询到 {len(act_id_list)} 个活动ID")
  536. for act_id in act_id_list:
  537. try:
  538. # 先将当前 pid 的状态改为 1,表示开始查询
  539. sql_pool.update_one("UPDATE zc_product_record SET player_state = 1 WHERE pid = %s", (act_id,))
  540. log.info(f"将 pid: {act_id} 的状态更新为 1(开始查询)")
  541. log.info(f"开始爬取pid: {act_id} 的玩家数据")
  542. has_data = get_player_list(log, act_id, token, sql_pool)
  543. # 根据是否有数据更新状态
  544. if has_data:
  545. log.info(f"pid: {act_id} 查询到数据,状态保持为 1")
  546. else:
  547. log.info(f"pid: {act_id} 没有数据,状态更新为 2")
  548. sql_pool.update_one("UPDATE zc_product_record SET player_state = 2 WHERE pid = %s", (act_id,))
  549. except Exception as pid_error:
  550. # 如果查询失败,将状态改为 3
  551. log.error(f"pid: {act_id} 查询失败,错误: {pid_error}")
  552. try:
  553. sql_pool.update_one("UPDATE zc_product_record SET player_state = 3 WHERE pid = %s", (act_id,))
  554. log.info(f"已将 pid: {act_id} 的状态更新为 3(查询异常)")
  555. except Exception as update_error:
  556. log.error(f"更新 pid: {act_id} 状态失败: {update_error}")
  557. except Exception as e:
  558. log.error(f'iterate_player_list error: {e}')
  559. except Exception as e:
  560. log.error(f'{inspect.currentframe().f_code.co_name} error: {e}')
  561. finally:
  562. log.info(f'爬虫程序 {inspect.currentframe().f_code.co_name} 运行结束,等待下一轮的采集任务............')
  563. def schedule_task():
  564. """
  565. 爬虫模块 定时任务 的启动文件
  566. """
  567. # 立即运行一次任务
  568. zc_main(log=logger)
  569. # 设置定时任务
  570. schedule.every().day.at("06:01").do(zc_main, log=logger)
  571. while True:
  572. schedule.run_pending()
  573. time.sleep(1)
  574. if __name__ == '__main__':
  575. # zc_main(logger)
  576. schedule_task()