sold_daily_spider.py 5.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105
  1. # -*- coding: utf-8 -*-
  2. # Author : Charley
  3. # Python : 3.12.10
  4. # Date : 2026/08/04
  5. """得卡 DECA 已售流程——每日完整采集 + 密集时段每小时占坑补采(常驻定时)。
  6. 每天 08:00 完整管道(sold-list 服务端只开放最近 60 条=3 页,全量/增量翻页已一致):
  7. 商家列表 → 每商家已售(翻3页入库) → 详情补抓 → 随机团回补 → 卡密清单 → 拆卡报告 + 视频回放。
  8. 拆卡报告/回放靠 report_state/replay_state 驱动,拿不到(还没生成)下轮重采。
  9. 另在 13:00-06:00 售卖密集时段每小时跑一次 hourly_task:只翻目标商家 3 页已售 INSERT IGNORE
  10. 占坑入库(不跑精加工),防止单商家两次采集间成交 >60 被挤出窗口漏采;精加工仍由 08:00 完整流程扫库补齐。
  11. 从根目录运行:python sold_daily_spider.py
  12. """
  13. import sys
  14. import time
  15. import os
  16. # 挂靠新项目根:sys.path 指向 common、CWD 对齐新根(application.yml / logs / 账号池 DB 生效)
  17. _ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
  18. sys.path.insert(0, os.path.join(_ROOT, "common"))
  19. os.chdir(_ROOT)
  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 deca_sold_core as core # 裸导入(与 common/ 内部各模块的裸导入一致,单实例);跳转靠 IDE 把 common 标为 Sources Root,见 .idea/*.iml
  25. logger.remove()
  26. logger.add("./logs/sold_daily_{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. # 密集时段(13:00-06:00 售卖集中)每小时占坑入库的目标商家配置:
  30. # sold-list 服务端只开放最近 60 条(3页),若相邻两次采集间隔内单商家新成交 >60,最早那批会被挤出
  31. # 窗口永久漏采;故高频补采尽早把 product_code 占坑入库。实测 881226408 单小时成交峰值 41<60,每小时够。
  32. HOURLY_MERCHANT_IDS = ["881226408"] # 目标商家(后续可扩多个)
  33. HOURLY_HOURS = list(range(13, 24)) + list(range(0, 7)) # 13,14,...,23,0,1,...,6 每整点各跑一次
  34. @retry(stop=stop_after_attempt(100), wait=wait_fixed(3600), after=core.after_log)
  35. def main_task(log):
  36. """已售每日完整采集(run_pipeline:商家→已售3页→详情→随机团→卡密→报告)。
  37. Args:
  38. log: 日志对象。
  39. Raises:
  40. RuntimeError: 数据库连接池异常时抛出以触发重试。
  41. """
  42. log.info(f"开始运行 {sys._getframe().f_code.co_name} 已售每日增量采集" + "." * 40)
  43. pool = MySQLConnectionPool(log=log)
  44. core.init_account_pool(pool, task_tag="sold_daily") # 启用账号池:need_auth 请求走 20 号池 + 各号专属 IP
  45. if not pool.check_pool_health():
  46. log.error("数据库连接池异常")
  47. raise RuntimeError("数据库连接池异常")
  48. try:
  49. core.run_pipeline(log, pool, incremental=True)
  50. except Exception as e:
  51. log.error(f"{sys._getframe().f_code.co_name} error: {e}")
  52. finally:
  53. log.info(f"已售每日增量采集 {sys._getframe().f_code.co_name} 运行结束" + "." * 20)
  54. def hourly_task(log):
  55. """密集时段每小时轻量占坑:只翻目标商家 3 页已售 INSERT IGNORE 入库,不跑精加工。
  56. sold-list 服务端只开放最近 60 条(3页)。若相邻两次采集间隔内单商家新成交 >60,最早那批会被
  57. 挤出窗口永久漏采。故在 13:00-06:00 售卖密集时段每小时补采一次,尽早把 product_code 占坑入库;
  58. 详情补抓 / 随机团回补 / 拆卡报告等精加工仍由每天 08:00 的 main_task(run_pipeline) 扫库存量统一
  59. 补齐——两者解耦,product 一旦进库即不丢,晚几小时加工无碍。
  60. Args:
  61. log: 日志对象。
  62. """
  63. log.info(f"开始运行 {sys._getframe().f_code.co_name} 密集时段每小时占坑" + "." * 30)
  64. pool = MySQLConnectionPool(log=log)
  65. core.init_account_pool(pool, task_tag="sold_hourly") # need_auth 请求走账号池 + 各号专属 IP
  66. if not pool.check_pool_health():
  67. log.error("数据库连接池异常,跳过本轮 hourly_task")
  68. return
  69. for mid in HOURLY_MERCHANT_IDS:
  70. try:
  71. core.get_sold_list(log, mid, pool) # 翻 3 页 INSERT IGNORE 占坑,不跑后续精加工
  72. except Exception as e:
  73. log.error(f"hourly_task get_sold_list error(商家 {mid}): {e}")
  74. log.info(f"{sys._getframe().f_code.co_name} 本轮结束" + "." * 20)
  75. def schedule_task():
  76. """定时任务入口:每天 08:00 完整采集(含精加工),13:00-06:00 密集时段每整点占坑补采。"""
  77. main_task(log=logger) # 启动立即完整跑一次
  78. schedule.every().day.at("08:00").do(main_task, log=logger) # 每天 08:00 完整流程(详情/团/报告精加工)
  79. for h in HOURLY_HOURS: # 密集时段每整点轻量占坑入库
  80. schedule.every().day.at(f"{h:02d}:00").do(hourly_task, log=logger)
  81. while True:
  82. schedule.run_pending()
  83. time.sleep(1)
  84. if __name__ == "__main__":
  85. schedule_task()