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