courtyard_spider.py 26 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648
  1. # -*- coding: utf-8 -*-
  2. # Author : Charley
  3. # Python : 3.10.8
  4. # Date : 2025/12/2 14:08
  5. import shutil
  6. import threading
  7. import time
  8. import inspect
  9. import requests
  10. import schedule
  11. import user_agent
  12. from loguru import logger
  13. from parsel import Selector
  14. from datetime import datetime
  15. from mysql_pool import MySQLConnectionPool
  16. from DrissionPage import ChromiumPage, ChromiumOptions
  17. from tenacity import retry, stop_after_attempt, wait_fixed
  18. """
  19. 扣驾的
  20. """
  21. logger.remove()
  22. logger.add("./logs/{time:YYYYMMDD}.log", encoding='utf-8', rotation="00:00",
  23. format="[{time:YYYY-MM-DD HH:mm:ss.SSS}] {level} {message}",
  24. level="DEBUG", retention="7 day")
  25. headers = {
  26. "accept": "application/json",
  27. "referer": "https://courtyard.io/",
  28. "user-agent": user_agent.generate_user_agent()
  29. }
  30. # 全局变量标识首次运行是否完成
  31. detail_first_run_completed = False
  32. # 详情页浏览器需拦截的资源模式。目标: 该页只需 DOM 结构做 xpath 解析(Activity history),
  33. # 故把「不影响 DOM 数据」的重资源全部拦掉——省请求、提速、降内存, 且不会 OOM。
  34. # 注意: 拦 CSS/字体不影响 xpath(xpath 只认 DOM 不认样式), 图片已由 no_imgs 关闭。
  35. # 视频/大媒体: 最占带宽与内存
  36. VIDEO_PATTERNS = [
  37. '*.mp4',
  38. '*.webm',
  39. '*.avi',
  40. '*.mov',
  41. '*.flv',
  42. '*.m3u8',
  43. '*.ts',
  44. '*googlevideo.com*',
  45. '*videoplayback*',
  46. '*video.twimg.com*',
  47. '*videocdn*',
  48. '*.mpd', # DASH视频流
  49. ]
  50. # 字体: 纯展示资源, 拦掉可省请求且绝不影响 DOM/JS。
  51. # 注意: 不要拦 *.css! courtyard.io 是 Next.js 应用, CSS 与 JS 同属路由 chunk 依赖图,
  52. # 拦掉 CSS 会导致 chunk-loading Promise 被 reject → 触发 React 错误边界(client-side exception)白屏。
  53. STYLE_FONT_PATTERNS = [
  54. '*.woff',
  55. '*.woff2',
  56. '*.ttf',
  57. '*.otf',
  58. '*.eot',
  59. ]
  60. # 第三方埋点/统计/广告/监控: 与业务数据无关, 纯拖慢加载
  61. TRACKER_PATTERNS = [
  62. '*google-analytics.com*',
  63. '*googletagmanager.com*',
  64. '*doubleclick.net*',
  65. '*facebook.net*',
  66. '*connect.facebook*',
  67. '*segment.com*',
  68. '*segment.io*',
  69. '*sentry.io*',
  70. '*mixpanel.com*',
  71. '*hotjar.com*',
  72. '*intercom.io*',
  73. '*fullstory.com*',
  74. '*amplitude.com*',
  75. ]
  76. # 汇总: 实际下发给浏览器的拦截清单
  77. BLOCKED_URL_PATTERNS = VIDEO_PATTERNS + STYLE_FONT_PATTERNS + TRACKER_PATTERNS
  78. # 连续失败达到该阈值即判定浏览器卡死, 触发重启
  79. MAX_CONSECUTIVE_FAILURES = 3
  80. # 每连续处理该条数就主动重启一次浏览器, 释放 Chromium 累积内存(V8 堆/渲染进程/缓存随页面数线性增长)
  81. BROWSER_RESTART_INTERVAL = 100
  82. # 每处理该条数就清理一次缓存/cookies, 减缓单实例内存增长
  83. BROWSER_CLEAR_CACHE_INTERVAL = 10
  84. # auto_port 临时用户目录的基路径(放 D 盘专属子目录)。auto_port 每次用随机端口, 会在此
  85. # 目录下不断新建 userData/{port} 且退出不自动删除, 故每次新建浏览器前会清空整个此目录
  86. BROWSER_TMP_PATH = r'D:\Drissionpage_temp\courtyard'
  87. def after_log(retry_state):
  88. """
  89. retry 回调
  90. :param retry_state: RetryCallState 对象
  91. """
  92. # 检查 args 是否存在且不为空
  93. if retry_state.args and len(retry_state.args) > 0:
  94. log = retry_state.args[0] # 获取传入的 logger
  95. else:
  96. log = logger # 使用全局 logger
  97. if retry_state.outcome.failed:
  98. log.warning(
  99. f"Function '{retry_state.fn.__name__}', Attempt {retry_state.attempt_number} Times")
  100. else:
  101. log.info(f"Function '{retry_state.fn.__name__}', Attempt {retry_state.attempt_number} succeeded")
  102. @retry(stop=stop_after_attempt(5), wait=wait_fixed(1), after=after_log)
  103. def get_proxys(log):
  104. """
  105. 获取代理
  106. :return: 代理
  107. """
  108. tunnel = "x371.kdltps.com:15818"
  109. kdl_username = "t13753103189895"
  110. kdl_password = "o0yefv6z"
  111. try:
  112. proxies = {
  113. "http": "http://%(user)s:%(pwd)s@%(proxy)s/" % {"user": kdl_username, "pwd": kdl_password, "proxy": tunnel},
  114. "https": "http://%(user)s:%(pwd)s@%(proxy)s/" % {"user": kdl_username, "pwd": kdl_password, "proxy": tunnel}
  115. }
  116. return proxies
  117. except Exception as e:
  118. log.error(f"Error getting proxy: {e}")
  119. raise e
  120. @retry(stop=stop_after_attempt(5), wait=wait_fixed(1), after=after_log)
  121. def get_goods_list(log, sql_pool):
  122. """
  123. 获取商品列表
  124. :param log: logger对象
  125. :param sql_pool: MySQL连接池对象
  126. :return:
  127. """
  128. log.info(f"========================== 开始获取商品列表 ==========================")
  129. url = "https://api.courtyard.io/vending-machines"
  130. response = requests.get(url, headers=headers, timeout=22)
  131. # print(response.text)
  132. response.raise_for_status()
  133. # 同样用 `or []` 兜底: 防止 vendingMachines 为 null 时下方 for 迭代抛 TypeError
  134. vendingMachines = response.json().get("vendingMachines") or []
  135. for item in vendingMachines:
  136. bag_id = item.get("id")
  137. bag_title = item.get("title")
  138. # sealed_pack_animation = item.get("sealedPackAnimation")
  139. # sealed_pack_image = item.get("sealedPackImage")
  140. category_title = item.get("category", {}).get("title")
  141. price = item.get("saleDetails", {}).get("salePriceUsd")
  142. data_dict = {
  143. "bag_id": bag_id,
  144. "bag_title": bag_title,
  145. "category": category_title,
  146. "price":price
  147. }
  148. # log.info(f'get_goods_list: {data_dict}')
  149. try:
  150. get_goods_detail(log, data_dict, sql_pool)
  151. except Exception as e:
  152. log.error(f"Error processing item: {e}")
  153. # 保存数据
  154. # if info_list:
  155. # log.info(f"获取商品列表成功, 共 {len(info_list)} 条数据")
  156. # sql_pool.insert_many(table="courtyard_vending_machines_record", data_list=info_list, ignore= True)
  157. @retry(stop=stop_after_attempt(5), wait=wait_fixed(1), after=after_log)
  158. def get_goods_detail(log, query_dict: dict, sql_pool=None):
  159. """
  160. 获取商品详情
  161. :param log: logger对象
  162. :param query_dict: query_dict
  163. :param sql_pool: MySQL连接池对象
  164. :return:
  165. """
  166. log.info(f"========================== 获取商品详情 ==========================")
  167. url = "https://api.courtyard.io/index/query/recent-pulls"
  168. params = {
  169. "limit": "250",
  170. # "vendingMachineIds": "pkmn-basic-pack"
  171. "vendingMachineIds": query_dict["bag_id"]
  172. }
  173. response = requests.get(url, headers=headers, params=params, timeout=22)
  174. # print(response.text)
  175. response.raise_for_status()
  176. resp_json = response.json()
  177. # 接口偶发返回 {"error": ...}(无 assets 键,多为限流/暂无数据),显式跳过并告警,避免静默吞掉
  178. if isinstance(resp_json, dict) and "assets" not in resp_json:
  179. log.warning(f"recent-pulls 返回无 assets, bag_id: {query_dict['bag_id']}, resp: {str(resp_json)[:150]}")
  180. return
  181. # 注意: dict.get(key, default) 的 default 仅在 key 缺失时生效;
  182. # 若 key 存在但值为 null(None) 仍返回 None, 故统一用 `or []` 兜底, 防止后续 len()/迭代抛 TypeError
  183. pulls = resp_json.get("assets") or []
  184. info_list = []
  185. for item in pulls:
  186. detail_title = item.get("title")
  187. detail_id = item.get("proof_of_integrity")
  188. if not detail_id:
  189. log.error(f"信息异常, detail_id: {detail_id}")
  190. continue
  191. # asset_pictures 实测存在为 null 的情况, 原写法 len(None) 是本次 TypeError 崩溃主因, 用 `or []` 兜底
  192. asset_pictures = item.get("asset_pictures") or []
  193. img_front = asset_pictures[0] if len(asset_pictures) > 0 else None
  194. img_back = asset_pictures[1] if len(asset_pictures) > 1 else None
  195. # crawl_date = time.strftime("%Y-%m-%d", time.localtime())
  196. data_dict = {
  197. "bag_id": query_dict["bag_id"],
  198. "bag_title": query_dict["bag_title"],
  199. "category": query_dict["category"],
  200. "price": query_dict["price"],
  201. "detail_id": detail_id,
  202. "detail_title": detail_title,
  203. "img_front": img_front,
  204. "img_back": img_back,
  205. # "crawler_date": crawl_date
  206. }
  207. # log.info(f'data_dict:{data_dict}')
  208. info_list.append(data_dict)
  209. # 保存数据
  210. if info_list:
  211. log.info(f"获取商品详情成功, 共 {len(info_list)} 条数据")
  212. sql_pool.insert_many(table="courtyard_list_record", data_list=info_list, ignore=True)
  213. def convert_time_format(time_str):
  214. """
  215. 将时间字符串转换为标准格式
  216. :param time_str: 原始时间字符串,如 "December 3, 2025 at 4:29 PM"
  217. :return: 标准时间格式字符串,如 "2025-12-03 16:29:00"
  218. """
  219. if not time_str:
  220. return None
  221. try:
  222. dt_obj = datetime.strptime(time_str, "%B %d, %Y at %I:%M %p")
  223. return dt_obj.strftime("%Y-%m-%d %H:%M:%S")
  224. except ValueError as e:
  225. logger.warning(f"时间转换失败: {time_str}, 错误: {e}")
  226. return None
  227. @retry(stop=stop_after_attempt(5), wait=wait_fixed(1), after=after_log)
  228. def get_sale_detail_single_page(log, page, sql_id, detail_id, sql_pool=None):
  229. """
  230. 获取商品详情
  231. :param log: logger对象
  232. :param page: page对象
  233. :param sql_id: 数据库id
  234. :param detail_id: 商品详情id
  235. :param sql_pool: MySQL连接池对象
  236. :return:
  237. """
  238. log.info(f"========================== 获取商品 <sale> 详情, sql_id: {sql_id} ==========================")
  239. # page_url = "https://courtyard.io/asset/a4f0bbebd858370567f1779fddf0f55630810116d80965e33940fc8ff5ac94b4"
  240. page_url = f"https://courtyard.io/asset/{detail_id}"
  241. page.get(page_url)
  242. # 加载策略为 none, get() 立即返回; 这里只等目标节点 "Activity history" 渲染出现即可开始解析,
  243. # 不等整页 load 完成——目标节点在, 就代表要抓的数据已就绪, 大幅缩短单条耗时
  244. target = page.ele('xpath://h6[text()="Activity history"]', timeout=25)
  245. if not target:
  246. log.error(f'{inspect.currentframe().f_code.co_name} -> 目标节点未出现(可能加载失败/被限流), 重试..........')
  247. raise Exception('目标节点未出现, 重新加载........') # 抛出异常以便重试
  248. log.debug(f'{inspect.currentframe().f_code.co_name} -> 目标节点已就绪, url: {page_url}')
  249. html = page.html
  250. if not html:
  251. log.error(f'{inspect.currentframe().f_code.co_name} -> 页面加载失败...........')
  252. raise Exception('页面加载失败, 重新加载........') # 抛出异常以便重试
  253. selector = Selector(text=html)
  254. # 方法一:通过文本内容匹配(优先)
  255. correlation_spans = selector.xpath('//span[contains(text(), ":")]/text()')
  256. correlation_id = None
  257. for text_selector in correlation_spans:
  258. correlation_id = text_selector.get() # ✅ 获取字符串
  259. # match = re.search(r'[\w\s]+:\s*(\d+)', text)
  260. # if match:
  261. # correlation_id = match.group(1)
  262. # break # 获取第一个有效 ID
  263. # correlation_spans = selector.xpath('//span[contains(text(), ":")]')
  264. # correlation_text = None
  265. #
  266. # for span in correlation_spans:
  267. # text = span.get()
  268. # if text and ":" in text:
  269. # correlation_text = text
  270. # break
  271. # 如果方法一失败,使用方法二:通过结构定位(备用)
  272. if not correlation_id:
  273. correlation_span = selector.xpath('//a[contains(@href, "cgccards.com")]/preceding-sibling::span[1]/text()')
  274. correlation_id = correlation_span.get()
  275. # 初始化所有可能的字段为None
  276. data_dict = {"detail_id": detail_id, "correlation_id": correlation_id, "burn_from": None,
  277. "burn_to": None, "burn_time": None, "sale_price": None, "sale_from": None, "sale_to": None,
  278. "sale_time": None, "mint_price": None, "mint_from": None, "mint_to": None, "mint_time": None}
  279. # 获取 "Activity history" 后面的 div
  280. activity_div = selector.xpath('//h6[text()="Activity history"]/following-sibling::div[1]/div')
  281. for tag_div in activity_div:
  282. tag_name = tag_div.xpath('./div[1]/div/span/text()').get()
  283. if not tag_name:
  284. continue
  285. if tag_name == "Burn":
  286. data_dict["burn_from"] = tag_div.xpath('./div[2]/div[1]//h6/text()').get()
  287. data_dict["burn_to"] = tag_div.xpath('./div[2]/div[2]//h6/text()').get()
  288. data_dict["burn_time"] = tag_div.xpath('./div[2]/div[3]/span/@aria-label').get()
  289. # December 3, 2025 at 4:29 PM 转换时间格式
  290. data_dict["burn_time"] = convert_time_format(data_dict["burn_time"])
  291. elif tag_name == "Sale":
  292. sale_price = tag_div.xpath('./div[2]/span/text()').get()
  293. if sale_price:
  294. sale_price = sale_price.replace("$", "").replace(",", "")
  295. data_dict["sale_price"] = sale_price
  296. data_dict["sale_from"] = tag_div.xpath('./div[3]/div[1]//h6/text()').get()
  297. data_dict["sale_to"] = tag_div.xpath('./div[3]/div[2]//h6/text()').get()
  298. data_dict["sale_time"] = tag_div.xpath('./div[3]/div[3]/span/@aria-label').get()
  299. # December 3, 2025 at 4:29 PM 转换时间格式
  300. data_dict["sale_time"] = convert_time_format(data_dict["sale_time"])
  301. elif tag_name == "Mint":
  302. mint_price = tag_div.xpath('./div[2]/span/text()').get()
  303. if mint_price:
  304. mint_price = mint_price.replace("$", "").replace(",", "")
  305. data_dict["mint_price"] = mint_price
  306. data_dict["mint_from"] = tag_div.xpath('./div[3]/div[1]//h6/text()').get()
  307. data_dict["mint_to"] = tag_div.xpath('./div[3]/div[2]//h6/text()').get()
  308. data_dict["mint_time"] = tag_div.xpath('./div[3]/div[3]/span/@aria-label').get()
  309. # December 3, 2025 at 4:29 PM 转换时间格式
  310. data_dict["mint_time"] = convert_time_format(data_dict["mint_time"])
  311. # log.info(f'Sale detail data: {data_dict}')
  312. # 保存数据
  313. sql_pool.insert_one_or_dict(table="courtyard_detail_record", data=data_dict, ignore=True)
  314. sql_pool.update_one("UPDATE courtyard_list_record SET state = 1 WHERE id = %s", (sql_id,))
  315. def _create_detail_browser(log):
  316. """创建并配置用于详情采集的浏览器实例。
  317. 集中管理浏览器启动参数、视频拦截规则与超时, 供首次启动和卡死重启复用。
  318. Args:
  319. log: loguru logger 对象, 从调用方透传。
  320. Returns:
  321. ChromiumPage: 已配置拦截规则与超时的浏览器页面对象。
  322. """
  323. # auto_port 的临时用户目录退出不会自动清理, 每次随机端口都会在 BROWSER_TMP_PATH 下留一份;
  324. # 故新建实例前先清空整个目录, 保证不残留缓存、不累积占磁盘(重启时旧实例已 quit, 无占用)
  325. shutil.rmtree(BROWSER_TMP_PATH, ignore_errors=True)
  326. options = ChromiumOptions()
  327. # 无需登录态, 不做任何持久化: auto_port 自动分配空闲端口(每次全新进程, 确保内存被彻底释放),
  328. # 相比固定端口: 分批重启不会因端口被残留进程占用而失败/误连到旧浏览器
  329. options.auto_port(True)
  330. # 临时用户目录基路径放到 D 盘专属子目录(默认在 C 盘 %TEMP%/DrissionPage)
  331. options.set_tmp_path(BROWSER_TMP_PATH)
  332. # options.set_proxy("http://" + tunnel)
  333. options.no_imgs(True)
  334. # 加载策略设为 none: page.get() 不等整页 load 完成立即返回, 之后由业务代码只等
  335. # 目标元素(Activity history)出现即开始解析——这是提速最关键的一招(SPA 整页 load
  336. # 往往还挂着一堆无关请求/长连接, 傻等会白白拖慢每一条)
  337. options.set_load_mode('none')
  338. # 禁止媒体相关设置
  339. options.set_argument('--autoplay-policy=user-gesture-required')
  340. options.set_argument('--disable-features=PreloadMediaEngagementData')
  341. # ---- 省资源(真正有效且无副作用的项): 减少下载与缓存占用, 不掐内存上限 ----
  342. # 说明: 不再使用 --max-old-space-size / --renderer-process-limit=1 / --process-per-site。
  343. # 前者会在重 SPA 需要更多 V8 堆时直接触发渲染进程 OOM 崩溃(即"喔唷崩溃啦 Out of Memory");
  344. # 后两者强制单渲染进程复用, 反而关闭了浏览器靠"进程销毁"回收内存的机制, 让内存只增不减。
  345. # 真正的内存回收依赖每 BROWSER_RESTART_INTERVAL 条整浏览器重启(见 get_sale_detail_list)。
  346. options.set_argument('--disk-cache-size=1') # 几乎禁用磁盘缓存
  347. options.set_argument('--media-cache-size=1') # 几乎禁用媒体缓存
  348. options.set_argument('--disable-application-cache') # 禁用应用缓存
  349. options.set_argument('--disable-dev-shm-usage') # 不占用共享内存, 避免共享内存耗尽
  350. options.set_argument('--disable-extensions') # 禁用扩展
  351. options.set_argument('--disable-background-networking') # 关闭后台网络活动
  352. # 最大化
  353. options.set_argument("--start-maximized")
  354. options.set_argument("--disable-gpu")
  355. options.set_argument("-accept-lang=en-US")
  356. page = ChromiumPage(options)
  357. # 使用 set.blocked_urls() 拦截视频/样式/字体/埋点等无关资源(不影响 DOM 取数)
  358. page.set.blocked_urls(BLOCKED_URL_PATTERNS)
  359. # 设置超时: 页面加载最多等 30s, 基础操作 20s, 避免浏览器卡死时无限等待
  360. page.set.timeouts(base=20, page_load=30)
  361. return page
  362. def _restart_detail_browser(log, page, reason='卡死'):
  363. """关闭旧浏览器并重建一个新实例。
  364. 先尽力 quit 旧实例(失败忽略), 等待端口释放后重新创建。既用于连续采集失败时的
  365. 自愈重启, 也用于每处理 BROWSER_RESTART_INTERVAL 条后的主动内存释放重启。
  366. Args:
  367. log: loguru logger 对象, 从调用方透传。
  368. page (ChromiumPage): 待关闭的旧浏览器实例。
  369. reason (str, optional): 重启原因, 仅用于日志区分。Defaults to '卡死'。
  370. Returns:
  371. ChromiumPage: 全新的浏览器页面对象。
  372. """
  373. try:
  374. page.quit()
  375. except Exception as e:
  376. # 卡死的浏览器 quit 本身也可能超时/报错, 忽略即可
  377. log.warning(f'关闭旧浏览器失败(忽略): {e}')
  378. # 等待端口(9138)与用户数据目录释放, 再重建, 避免端口占用导致启动失败
  379. time.sleep(3)
  380. log.warning(f'浏览器{reason}, 正在重启新实例..........')
  381. return _create_detail_browser(log)
  382. @retry(stop=stop_after_attempt(5), wait=wait_fixed(1), after=after_log)
  383. def get_sale_detail_list(log, sql_pool=None):
  384. """
  385. 获取商品详情
  386. :param log: logger对象
  387. # :param detail_id_list: 详情id列表
  388. :param sql_pool: MySQL连接池对象
  389. :return:
  390. """
  391. log.info(f"========================== 获取商品 <sale> 详情 LIST ==========================")
  392. page = _create_detail_browser(log)
  393. # 连续失败计数: 用于区分"单条数据问题"与"浏览器整体卡死", 后者需重启浏览器
  394. consecutive_failures = 0
  395. try:
  396. sql_detail_id_list = sql_pool.select_all("SELECT id, detail_id FROM courtyard_list_record WHERE state = 0")
  397. for idx, sql_detail_id in enumerate(sql_detail_id_list):
  398. sql_id = sql_detail_id[0]
  399. detail_id = sql_detail_id[1]
  400. # 每处理 BROWSER_RESTART_INTERVAL 条主动重启浏览器, 彻底释放 Chromium 累积内存;
  401. # 单实例连续跑几百页时 V8 堆/渲染进程/缓存会线性增长到数 GB, 这是内存/CPU 高的主因
  402. if idx > 0 and idx % BROWSER_RESTART_INTERVAL == 0:
  403. log.info(f'已连续处理 {idx} 条, 达到批次阈值, 主动重启浏览器释放内存..........')
  404. page = _restart_detail_browser(log, page, reason='达到批次重启阈值')
  405. consecutive_failures = 0
  406. try:
  407. get_sale_detail_single_page(log, page, sql_id, detail_id, sql_pool)
  408. consecutive_failures = 0 # 成功即清零, 只累计"连续"失败
  409. # 定期清理缓存/cookies, 减缓单实例在两次批次重启之间的内存增长
  410. if idx > 0 and idx % BROWSER_CLEAR_CACHE_INTERVAL == 0:
  411. try:
  412. page.clear_cache()
  413. except Exception as ce:
  414. log.warning(f'清理浏览器缓存失败(忽略): {ce}')
  415. except Exception as e:
  416. consecutive_failures += 1
  417. log.error(f'get_sale_detail_single_page error '
  418. f'({consecutive_failures}/{MAX_CONSECUTIVE_FAILURES}), sql_id: {sql_id}: {e}')
  419. if consecutive_failures >= MAX_CONSECUTIVE_FAILURES:
  420. # 连续多条失败, 判定浏览器已卡死: 不标记 state=2(避免误伤),
  421. # 保留 state=0 待下轮重试, 并重启浏览器后继续处理后续条目
  422. page = _restart_detail_browser(log, page)
  423. consecutive_failures = 0
  424. else:
  425. # 未达阈值, 视为该条数据自身问题, 标记 state=2 跳过
  426. sql_pool.update_one("UPDATE courtyard_list_record SET state = 2 WHERE id = %s", (sql_id,))
  427. except Exception as e:
  428. log.error(f'get_response error: {e}')
  429. raise # 直接透传原始异常(原写法 raise 字符串会抛 TypeError 掩盖真因)
  430. finally:
  431. try:
  432. page.quit()
  433. except Exception as e:
  434. log.warning(f'最终关闭浏览器失败(忽略): {e}')
  435. @retry(stop=stop_after_attempt(100), wait=wait_fixed(3600), after=after_log)
  436. def list_main(log):
  437. """
  438. 主函数 自动售货机
  439. :param log: logger对象
  440. """
  441. log.info(
  442. f'开始运行 {inspect.currentframe().f_code.co_name} 爬虫任务....................................................')
  443. start = time.time()
  444. # 配置 MySQL 连接池
  445. sql_pool = MySQLConnectionPool(log=log)
  446. if not sql_pool.check_pool_health():
  447. log.error("数据库连接池异常")
  448. raise RuntimeError("数据库连接池异常")
  449. try:
  450. try:
  451. log.debug('------------------- 开始获取商品列表 -------------------')
  452. get_goods_list(log, sql_pool)
  453. except Exception as e:
  454. log.error(f'get_goods_list error: {e}')
  455. except Exception as e:
  456. log.error(f'{inspect.currentframe().f_code.co_name} error: {e}')
  457. finally:
  458. log.info(f'爬虫程序 {inspect.currentframe().f_code.co_name} 运行结束,等待下一轮的采集任务............')
  459. end = time.time()
  460. elapsed_time = end - start
  461. log.info(f'============================== 本次爬虫运行时间:{elapsed_time:.2f} 秒 ===============================')
  462. return elapsed_time
  463. @retry(stop=stop_after_attempt(100), wait=wait_fixed(3600), after=after_log)
  464. def detail_main(log):
  465. """
  466. 主函数 自动售货机
  467. :param log: logger对象
  468. """
  469. log.info(
  470. f'开始运行 {inspect.currentframe().f_code.co_name} 爬虫任务....................................................')
  471. # 配置 MySQL 连接池
  472. sql_pool = MySQLConnectionPool(log=log)
  473. if not sql_pool.check_pool_health():
  474. log.error("数据库连接池异常")
  475. raise RuntimeError("数据库连接池异常")
  476. global detail_first_run_completed
  477. try:
  478. # 获取详情页信息
  479. try:
  480. log.debug('------------------- 获取商品 detail 数据 -------------------')
  481. get_sale_detail_list(log, sql_pool)
  482. except Exception as e:
  483. log.error(f'get_sale_detail_list error: {e}')
  484. except Exception as e:
  485. log.error(f'{inspect.currentframe().f_code.co_name} error: {e}')
  486. finally:
  487. detail_first_run_completed = True
  488. log.info(f'爬虫程序 {inspect.currentframe().f_code.co_name} 运行结束,等待下一轮的采集任务............')
  489. def control_list_mask(log):
  490. """
  491. 控制列表爬虫任务 每10分钟运行
  492. :param log: logger对象
  493. """
  494. while True:
  495. log.info(
  496. f'--------------------- 开始运行 {inspect.currentframe().f_code.co_name} 新一轮的爬虫任务 ---------------------')
  497. elapsed_time = list_main(log)
  498. # 计算剩余时间
  499. wait_time = max(0, 300 - int(elapsed_time))
  500. if wait_time > 0:
  501. log.info(f"程序运行时间{elapsed_time:.2f}秒, 小于 5 分钟,等待 {wait_time:.2f} 秒后再开始下一轮任务")
  502. time.sleep(wait_time)
  503. else:
  504. log.info("程序运行时间大于等于5分钟,直接开始下一轮任务")
  505. def scheduled_detail_main(log):
  506. """定时任务调用的包装函数"""
  507. global detail_first_run_completed
  508. if detail_first_run_completed:
  509. detail_main(log)
  510. else:
  511. log.info("Skipping scheduled task as first run is not completed yet")
  512. def run_threaded(job_func, *args, **kwargs):
  513. """
  514. 在新线程中运行给定的函数,并传递参数。
  515. :param job_func: 要运行的目标函数
  516. :param args: 位置参数
  517. :param kwargs: 关键字参数
  518. """
  519. job_thread = threading.Thread(target=job_func, args=args, kwargs=kwargs)
  520. job_thread.start()
  521. def schedule_task():
  522. """
  523. 设置定时任务
  524. """
  525. # 启动 control_list_mask 任务线程
  526. list_thread = threading.Thread(target=control_list_mask, args=(logger,))
  527. list_thread.daemon = True # 设置为守护线程,主程序退出时自动结束
  528. list_thread.start()
  529. # 启动 detail_main 任务线程(首次运行)
  530. detail_thread = threading.Thread(target=detail_main, args=(logger,))
  531. detail_thread.daemon = True
  532. detail_thread.start()
  533. # 设置定时任务 每天
  534. # schedule.every().day.at("00:01").do(run_threaded, detail_main, logger)
  535. schedule.every().day.at("00:01").do(run_threaded, scheduled_detail_main, logger)
  536. while True:
  537. schedule.run_pending()
  538. time.sleep(1)
  539. if __name__ == '__main__':
  540. schedule_task()
  541. # detail_main(log=logger)
  542. # get_sale_detail_list(log, ((1, 'a4f0bbebd858370567f1779fddf0f55630810116d80965e33940fc8ff5ac94b4'),
  543. # (2, 'a4f0bbebd858370567f1779fddf0f55630810116d80965e33940fc8ff5ac94b4')), sql_pool)