# -*- 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()