Ver Fonte

refactor(mysql_pool): 优化数据库执行函数,支持断连重试和日志改进

- 增加_execute函数支持断开连接时重试机制,提升稳定性
- 对完整性错误做特殊处理,避免打印堆栈污染日志
- insert_one_or_dict函数中针对重复条目调整日志级别为warning,避免堆栈输出
- insert_many异常捕获时改为error级别日志输出

refactor(pokecolor_new_daily_spider): 实现已售列表接口签名和增量翻页逻辑

- 增加接口请求签名计算,防止接口限制
- 构建动态请求头,包含签名相关字段
- 修改已售列表翻页逻辑,支持按日期增量采集
- 增加匿名请求最大翻页限制,规避接口重复数据问题
- 优化日志记录,增加翻页详细信息和采集完成判断
- 改进异常处理,非正常code抛出异常触发重试

refactor(YamlLoader): 增加yaml配置路径解析,支持主脚本目录查找

- 增加_resolve_path函数,在当前工作目录找不到时尝试主脚本目录定位配置文件
- readYaml函数改用_resolve_path定位文件,保证兼容性和灵活性
- 所有文件打开操作统一加入utf-8编码声明

refactor: 移除pokecolor_history_spider和pokecolor_spider/details.py文件

- 删除已弃用pokecolor_history_spider.py及其中相关代码
- 删除pokecolor_spider/details.py文件,移除未使用的订单详情获取函数
charley há 2 semanas atrás
pai
commit
5f2b316fda

+ 46 - 26
pokecolor_spider/YamlLoader.py

@@ -1,6 +1,6 @@
 # -*- coding: utf-8 -*-
 # Author : Charley
-# Python : 3.10.8
+# Python : 3.12.10
 # Date   : 2025/12/22 10:44
 import os, re
 import yaml
@@ -11,10 +11,10 @@ regex = re.compile(r'^\$\{(?P<ENV>[A-Z_\-]+:)?(?P<VAL>[\w.]+)}$')
 class YamlConfig:
     def __init__(self, config):
         self.config = config
-    
+
     def get(self, key: str):
         return YamlConfig(self.config.get(key))
-    
+
     def getValueAsString(self, key: str):
         try:
             match = regex.match(self.config[key])
@@ -25,7 +25,7 @@ class YamlConfig:
             return None
         except:
             return self.config[key]
-    
+
     def getValueAsInt(self, key: str):
         try:
             match = regex.match(self.config[key])
@@ -36,7 +36,7 @@ class YamlConfig:
             return 0
         except:
             return int(self.config[key])
-    
+
     def getValueAsBool(self, key: str):
         try:
             match = regex.match(self.config[key])
@@ -49,30 +49,50 @@ class YamlConfig:
             return bool(self.config[key])
 
 
-def readYaml(path: str = 'application.yml', profile: str = None) -> YamlConfig:
+def _resolve_path(path: str) -> str:
+    """
+    解析 yaml 文件路径,按优先级查找:
+      1) 绝对路径或 cwd 下存在 → 直接用(保留旧行为,向后兼容)
+      2) 调用方主脚本所在目录 → 兜底,方便打包后从任意 cwd 启动
+    :param path: (str) 用户传入的路径,默认 'application.yml'
+    :return: (str) 实际可读取的完整路径;找不到则返回原 path 让 open() 抛错
+    """
+    # 1) 旧行为:cwd 或绝对路径
     if os.path.exists(path):
-        with open(path) as fd:
-            conf = yaml.load(fd, Loader=yaml.FullLoader)
-    
+        return path
+
+    # 2) 主脚本目录(__main__.__file__)
+    try:
+        import __main__
+        main_file = getattr(__main__, '__file__', None)
+        if main_file:
+            candidate = os.path.join(os.path.dirname(os.path.abspath(main_file)), path)
+            if os.path.exists(candidate):
+                return candidate
+    except Exception:
+        pass
+
+    return path
+
+
+def readYaml(path: str = 'application.yml', profile: str = None) -> YamlConfig:
+    """
+    读取 yaml 配置。
+    :param path: (str) yaml 文件路径,默认 'application.yml'。
+                       优先 cwd / 绝对路径(保留旧行为),找不到再 fallback 到主脚本所在目录。
+    :param profile: (str) 可选环境后缀,如 'dev' 会额外加载 'application-dev.yml' 并 update
+    :return: (YamlConfig) 配置访问对象
+    :raises FileNotFoundError: cwd 和主脚本目录都找不到时抛出
+    """
+    real_path = _resolve_path(path)
+    with open(real_path, encoding='utf-8') as fd:
+        conf = yaml.load(fd, Loader=yaml.FullLoader)
+
     if profile is not None:
-        result = path.split('.')
+        result = real_path.rsplit('.', 1)
         profiledYaml = f'{result[0]}-{profile}.{result[1]}'
         if os.path.exists(profiledYaml):
-            with open(profiledYaml) as fd:
+            with open(profiledYaml, encoding='utf-8') as fd:
                 conf.update(yaml.load(fd, Loader=yaml.FullLoader))
-    
-    return YamlConfig(conf)
-
-# res = readYaml()
-# mysqlConf = res.get('mysql')
-# print(mysqlConf)
 
-# print(res.getValueAsString("host"))
-# mysqlYaml = mysqlConf.getValueAsString("host")
-# print(mysqlYaml)
-# host = mysqlYaml.get("host").split(':')[-1][:-1]
-# port = mysqlYaml.get("port").split(':')[-1][:-1]
-# username = mysqlYaml.get("username").split(':')[-1][:-1]
-# password = mysqlYaml.get("password").split(':')[-1][:-1]
-# mysql_db = mysqlYaml.get("db").split(':')[-1][:-1]
-# print(host,port,username,password)
+    return YamlConfig(conf)

+ 0 - 14
pokecolor_spider/details.py

@@ -1,14 +0,0 @@
-# -*- coding: utf-8 -*-
-# Author : Charley
-# Python : 3.10.8
-# Date   : 2026/3/17 16:50
-
-# -----------------------------------------------------------------------------------------------
-# def get_details(order_id):
-#     """
-#     获取订单详情
-#     """
-#     # url = "https://api.pokecolor.cn/api/h5/trade/sellorder/XS2601281100b85329/"
-#     url = f"https://api.pokecolor.cn/api/h5/trade/sellorder/{order_id}/"
-#     response = requests.get(url, headers=headers)
-#     print(response.text)

+ 61 - 18
pokecolor_spider/mysql_pool.py

@@ -1,6 +1,6 @@
 # -*- coding: utf-8 -*-
 # Author : Charley
-# Python : 3.10.8
+# Python : 3.12.10
 # Date   : 2025/3/25 14:14
 import re
 import pymysql
@@ -50,27 +50,69 @@ class MySQLConnectionPool:
             write_timeout=30  # 写入超时时间(秒)
         )
 
+    # def _execute(self, query, args=None, commit=False):
+    #     """
+    #     执行SQL
+    #     :param query: SQL语句
+    #     :param args: SQL参数
+    #     :param commit: 是否提交事务
+    #     :return: 查询结果
+    #     """
+    #     try:
+    #         with self.pool.connection() as conn:
+    #             with conn.cursor() as cursor:
+    #                 cursor.execute(query, args)
+    #                 if commit:
+    #                     conn.commit()
+    #                 self.log.debug(f"sql _execute, Query: {query}, Rows: {cursor.rowcount}")
+    #                 return cursor
+    #     except Exception as e:
+    #         if commit and conn:
+    #             conn.rollback()
+    #         self.log.exception(f"Error executing query: {e}, Query: {query}, Args: {args}")
+    #         raise e
+
     def _execute(self, query, args=None, commit=False):
         """
-        执行SQL
+        执行SQL(带断连重试)
         :param query: SQL语句
         :param args: SQL参数
         :param commit: 是否提交事务
         :return: 查询结果
         """
-        try:
-            with self.pool.connection() as conn:
-                with conn.cursor() as cursor:
-                    cursor.execute(query, args)
-                    if commit:
-                        conn.commit()
-                    self.log.debug(f"sql _execute, Query: {query}, Rows: {cursor.rowcount}")
-                    return cursor
-        except Exception as e:
-            if commit and conn:
-                conn.rollback()
-            self.log.exception(f"Error executing query: {e}, Query: {query}, Args: {args}")
-            raise e
+        conn = None
+        for attempt in range(2):  # 最多重试1次
+            try:
+                with self.pool.connection() as conn:
+                    with conn.cursor() as cursor:
+                        cursor.execute(query, args)
+                        if commit:
+                            conn.commit()
+                        self.log.debug(f"sql _execute, Query: {query}, Rows: {cursor.rowcount}")
+                        return cursor
+            except pymysql.err.InterfaceError as e:
+                # 连接已断开,重试一次
+                if attempt == 0:
+                    self.log.warning(f"数据库连接断开,正在重试... Error: {e}")
+                    continue
+                self.log.error(f"重试后仍失败: {e}, Query: {query}")
+                raise e
+            except pymysql.err.IntegrityError:
+                # 完整性错误(如重复条目)交由上层处理,避免在此打印完整堆栈污染日志
+                if commit and conn:
+                    try:
+                        conn.rollback()
+                    except Exception:
+                        pass
+                raise
+            except Exception as e:
+                if commit and conn:
+                    try:
+                        conn.rollback()
+                    except Exception:
+                        pass
+                self.log.exception(f"Error executing query: {e}, Query: {query}, Args: {args}")
+                raise e
 
     def select_one(self, query, args=None):
         """
@@ -180,13 +222,14 @@ class MySQLConnectionPool:
             return cursor.lastrowid
         except pymysql.err.IntegrityError as e:
             if "Duplicate entry" in str(e):
-                self.log.warning(f"插入失败:重复条目,已跳过。错误详情: {e}")
+                # 重复条目用 warning 简短输出,不打印堆栈
+                self.log.warning(f"插入跳过-重复条目 Table: {table}, {e.args[1] if len(e.args) > 1 else e}")
                 return -1  # 返回 -1 表示重复条目被跳过
             else:
-                self.log.exception(f"数据库完整性错误: {e}")
+                self.log.error(f"数据库完整性错误 Table: {table}, Error: {e}")
                 raise
         except Exception as e:
-            self.log.exception(f"未知错误: {e}")
+            self.log.error(f"insert_one_or_dict 失败 Table: {table}, Error: {e}")
             raise
 
     def insert_many(self, table=None, data_list=None, query=None, args_list=None, batch_size=1000, commit=True,

+ 0 - 223
pokecolor_spider/pokecolor_history_spider.py

@@ -1,223 +0,0 @@
-# -*- coding: utf-8 -*-
-# Author : Charley
-# Python : 3.10.8
-# Date   : 2026/3/17 14:57
-import inspect
-import requests
-import user_agent
-from loguru import logger
-from mysql_pool import MySQLConnectionPool
-from tenacity import retry, stop_after_attempt, wait_fixed
-
-logger.remove()
-logger.add("./logs/{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")
-"""
-https://pokecolor.cn/h5/pages-card/card/history
-"""
-
-headers = {
-    "user-agent": user_agent.generate_user_agent()
-}
-
-
-def after_log(retry_state):
-    """
-    retry 回调
-    :param retry_state: RetryCallState 对象
-    """
-    # 检查 args 是否存在且不为空
-    if retry_state.args and len(retry_state.args) > 0:
-        log = retry_state.args[0]  # 获取传入的 logger
-    else:
-        log = logger  # 使用全局 logger
-
-    if retry_state.outcome.failed:
-        log.warning(
-            f"Function '{retry_state.fn.__name__}', Attempt {retry_state.attempt_number} Times")
-    else:
-        log.info(f"Function '{retry_state.fn.__name__}', Attempt {retry_state.attempt_number} succeeded")
-
-
-@retry(stop=stop_after_attempt(5), wait=wait_fixed(1), after=after_log)
-def get_proxys(log):
-    """
-    获取代理配置
-
-    :param log: 日志对象
-    :return: 代理字典
-    """
-    tunnel = "x371.kdltps.com:15818"
-    kdl_username = "t13753103189895"
-    kdl_password = "o0yefv6z"
-    try:
-        proxies = {
-            "http": "http://%(user)s:%(pwd)s@%(proxy)s/" % {"user": kdl_username, "pwd": kdl_password, "proxy": tunnel},
-            "https": "http://%(user)s:%(pwd)s@%(proxy)s/" % {"user": kdl_username, "pwd": kdl_password, "proxy": tunnel}
-        }
-        return proxies
-    except Exception as e:
-        log.error(f"Error getting proxy: {e}")
-        raise e
-
-
-@retry(stop=stop_after_attempt(5), wait=wait_fixed(1), after=after_log)
-def get_sold_single_page(log, page_num=1):
-    """
-    获取已售列表单页数据
-
-    :param log: logger对象
-    :param page_num: 页码
-    :return: 订单列表和总条数
-    """
-    log.debug(f"Fetching page {page_num}...........")
-    url = "https://api.pokecolor.cn/api/h5/trade/sellorder/"
-    params = {
-        "page": str(page_num),
-        "page_size": "20",
-        "order_by": "-order_deal_on",
-        "status_filter": "sold_finish"
-    }
-    response = requests.get(url, headers=headers, params=params, proxies=get_proxys(log), timeout=22)
-    data = response.json()
-
-    if data.get("code") == 100 and data.get("msg") == "ok":
-        results = data.get("data", {}).get("results", [])
-        total_count = data.get("data", {}).get("count", 0)
-        log.info(f"Successfully fetched page {page_num} with {len(results)} items")
-        return results, total_count
-    else:
-        log.error(f"Error fetching page {page_num}: {data}")
-        return None, None
-
-
-def parse_list_results(log, results, sql_pool):
-    """
-    解析API返回的订单数据
-
-    :param log: logger对象
-    :param results: 订单数据列表
-    :param sql_pool: 数据库连接池对象
-    """
-    log.debug("Parsing order results...........")
-    parsed_orders = []
-
-    def convert_iso_time(time_str):
-        """
-        从 ISO 8601 时间格式中提取日期部分(年月日)
-        :param time_str: ISO 格式时间字符串(如:2026-03-16T22:34:17.338857+08:00)
-        :return: 日期字符串(YYYY-MM-DD)或 None
-        """
-        if not time_str:
-            return None
-        try:
-            # 使用 split('T') 分割,取日期部分
-            return time_str.split('T')[0]
-        except (IndexError, AttributeError):
-            return None
-
-    for order in results:
-        # 解析订单信息
-        parsed_order = {
-            'uid': order.get('uid'),
-            'image': order.get('image'),
-            'title': order.get('display_title'),
-            # 'status': order.get('status'),
-            'display_status': order.get('display_status'),
-            # 'sell_link': order.get('sell_link'),
-            'price': order.get('price'),
-            # 'start_bid_price': order.get('start_bid_price'),
-            'order_deal_on': convert_iso_time(order.get('order_deal_on')),
-            # 'bid_total': order.get('bid_total'),
-            'order_finish_on': convert_iso_time(order.get('order_finish_on')),
-            'card_category': order.get('card_category'),
-            'rate_type': order.get('rate_type'),
-            'rate_score_display': order.get('rate_score_display'),
-        }
-        special_performance = order.get('special_performance')
-        if special_performance:
-            parsed_order['sell_id'] = special_performance.get('id')
-            parsed_order['sell_name'] = special_performance.get('name')
-        else:
-            parsed_order['sell_id'] = None
-            parsed_order['sell_name'] = None
-
-        # print(parsed_order)
-        parsed_orders.append(parsed_order)
-
-    if parsed_orders:
-        sql_pool.insert_many(table="pokecolor_sold_record", data_list=parsed_orders, ignore=True)
-
-    return parsed_orders
-
-
-def get_sold_list(log, sql_pool):
-    """
-    获取已售列表数据
-
-    :param log: logger对象
-    :param sql_pool: 数据库连接池对象
-    """
-    page_size = 20
-    total_collected = 0
-    page = 1
-
-    while True:
-        result = get_sold_single_page(log, page)
-
-        if result[0] is None:
-            log.error(f"Error fetching page {page}, stopping...")
-            break
-
-        orders, total_count = result
-        total_collected += len(parse_list_results(log, orders, sql_pool))
-
-        # 第一页时打印总数
-        if page == 1:
-            total_pages = (total_count + page_size - 1) // page_size
-            log.info(f"Total items: {total_count}, Total pages: {total_pages}")
-
-        if len(orders) < page_size:
-            log.info(f"Less than {page_size} items on page {page}, stopping...")
-            break
-
-        # 检查是否还有下一页
-        if page * page_size >= total_count:
-            log.info(f"No more pages, stopping...")
-            break
-
-        page += 1
-
-    log.info(f"Total orders collected: {total_collected}")
-    return total_collected
-
-
-@retry(stop=stop_after_attempt(100), wait=wait_fixed(3600), after=after_log)
-def kl_main(log):
-    """
-    主函数
-
-    :param log: logger对象
-    """
-    log.info(
-        f'开始运行 {inspect.currentframe().f_code.co_name} 爬虫任务....................................................')
-
-    # 配置 MySQL 连接池
-    sql_pool = MySQLConnectionPool(log=log)
-    if not sql_pool.check_pool_health():
-        log.error("数据库连接池异常")
-        raise RuntimeError("数据库连接池异常")
-
-    try:
-        get_sold_list(log, sql_pool)
-    except Exception as e:
-        log.error(f'{inspect.currentframe().f_code.co_name} error: {e}')
-    finally:
-        log.info(f'爬虫程序 {inspect.currentframe().f_code.co_name} 运行结束,等待下一轮的采集任务............')
-
-
-if __name__ == '__main__':
-    kl_main(log=logger)
-    # get_details('XS2601281100b85329')
-

+ 127 - 32
pokecolor_spider/pokecolor_new_daily_spider.py

@@ -4,6 +4,10 @@
 # Date   : 2026/3/17 14:57
 import time
 import inspect
+import hashlib
+import random
+import string
+from datetime import datetime, timedelta
 import requests
 import schedule
 import user_agent
@@ -19,9 +23,49 @@ logger.add("./logs/{time:YYYYMMDD}.log", encoding='utf-8', rotation="00:00",
 https://pokecolor.cn/h5/pages-card/card/history
 """
 
-headers = {
-    "user-agent": user_agent.generate_user_agent()
-}
+SIGNATURE_SECRET = "pk_code_a_b_4321_sc"
+
+# 单页条数(与接口 page_size 参数保持一致)
+PAGE_SIZE = 20
+# 未登录(匿名)时服务端把已售列表锁死在前 600 条:第 1~30 页正常,
+# 第 31 页起返回的是第 30 页的重复数据(next=False)。故匿名场景最多翻到第 30 页。
+MAX_ANON_PAGE = 30
+
+
+def build_x_r_content():
+    alphabet = string.ascii_letters + string.digits
+    return "".join(random.choice(alphabet) for _ in range(15)) + random.choice("ABCDEF")
+
+
+def build_verify_token(method, timestamp, x_r_content, params):
+    rounds = int(x_r_content[-1], 16)
+    param_string = "&".join(
+        f"{key}={value}"
+        for key, value in sorted(params.items())
+        if value is not None and value != ""
+    )
+    token_source = f"{method.upper()}{timestamp}{x_r_content}{param_string}{SIGNATURE_SECRET}"
+    for _ in range(rounds):
+        token_source = hashlib.md5(token_source.encode("utf-8")).hexdigest()
+    return token_source
+
+
+def build_headers(method, params):
+    timestamp = str(int(time.time() * 1000))
+    x_r_content = build_x_r_content()
+    return {
+        "accept": "application/json, text/plain, */*",
+        "accept-language": "zh-CN,zh;q=0.9,en;q=0.8",
+        "origin": "https://pokecolor.cn",
+        "Paytype": "kl",
+        "Platform": "web",
+        "referer": "https://pokecolor.cn/",
+        "user-agent": user_agent.generate_user_agent(),
+        "v": "1.11.0.481",
+        "X-R-Content": x_r_content,
+        "X-Timestamp": timestamp,
+        "X-Verify-Token": build_verify_token(method, timestamp, x_r_content, params),
+    }
 
 
 def after_log(retry_state):
@@ -81,17 +125,22 @@ def get_sold_single_page(log, page_num=1):
         "order_by": "-order_deal_on",
         "status_filter": "sold_finish"
     }
-    response = requests.get(url, headers=headers, params=params, proxies=get_proxys(log), timeout=22)
+    response = requests.get(url, headers=build_headers("GET", params), params=params, proxies=get_proxys(log), timeout=22)
     data = response.json()
 
     if data.get("code") == 100 and data.get("msg") == "ok":
-        results = data.get("data", {}).get("results", [])
-        total_count = data.get("data", {}).get("count", 0)
-        log.info(f"Successfully fetched page {page_num} with {len(results)} items")
-        return results, total_count
-    else:
-        log.error(f"Error fetching page {page_num}: {data}")
-        return None, None
+        payload = data.get("data", {}) or {}
+        results = payload.get("results", []) or []
+        total_count = payload.get("count", 0)
+        has_next = payload.get("next", False)
+        log.info(f"Successfully fetched page {page_num} with {len(results)} items, next={has_next}")
+        return results, total_count, has_next
+
+    # 非 100 一律抛出,交给 tenacity 重试。
+    # 注意:接口限频时会返回 code=197、data="参数错误"(用「参数错误」伪装限流),
+    # 并非真的参数错,直接重试/退避即可,不能当作正常结束而 break。
+    raise RuntimeError(
+        f"page {page_num} 返回异常 code={data.get('code')} msg={data.get('msg')} data={data.get('data')}")
 
 
 def parse_list_results(log, results, sql_pool):
@@ -155,45 +204,91 @@ def parse_list_results(log, results, sql_pool):
     return parsed_orders
 
 
-def get_sold_list(log, sql_pool):
+def parse_deal_date(order):
+    """从订单里取成交日期(仅年月日)。
+
+    结果为 page 内最旧一条与 cutoff 比较用,接口按 -order_deal_on 倒序返回,
+    故一页里最后一条即最旧一条。
+
+    Args:
+        order (dict): 单条订单数据。
+
+    Returns:
+        datetime.date: 成交日期;取不到或解析失败时返回 None。
     """
-    获取已售列表数据
+    time_str = order.get("order_deal_on")
+    if not time_str:
+        return None
+    try:
+        # order_deal_on 形如 2026-07-19T07:31:05.142972+08:00,取 T 前的日期部分即可
+        return datetime.strptime(time_str.split("T")[0], "%Y-%m-%d").date()
+    except (ValueError, AttributeError, IndexError):
+        return None
 
-    :param log: logger对象
-    :param sql_pool: 数据库连接池对象
+
+def get_sold_list(log, sql_pool, days_back=1):
+    """增量采集已售列表。
+
+    按成交时间倒序从第 1 页往前翻,翻到「本页最旧成交日 < cutoff」即停,只补最近几天的新增。
+    cutoff = 今天 - days_back,且判定为严格小于,故 cutoff 当天的数据仍会被收下。
+    例:00:01 定时跑、days_back=1 时 cutoff=昨天,翻到「前天」的数据即停 —— 正好收完整的昨天。
+
+    匿名请求下服务端最多放行前 600 条(30 页),故硬上限为 MAX_ANON_PAGE;若翻到上限仍未
+    追到 cutoff,说明近 days_back 天新增已售超过 600 条(本站单日成交量约数百条,days_back
+    调大很容易超上限),匿名无法继续深翻,记 warning 提示可能遗漏。
+
+    Args:
+        log: logger 对象。
+        sql_pool: 数据库连接池对象。
+        days_back (int, optional): 增量回溯天数,cutoff = 今天 - days_back。Defaults to 1。
+
+    Returns:
+        int: 本轮解析并尝试入库的订单条数(含被 ignore 的重复项)。
     """
-    page_size = 20
+    cutoff = (datetime.now() - timedelta(days=days_back)).date()
+    log.info(f"增量采集已售列表,cutoff={cutoff}(成交日早于此日期即停),匿名翻页上限={MAX_ANON_PAGE}")
+
     total_collected = 0
     page = 1
-    max_pages = 100
+    reached_cutoff = False
 
-    while page <= max_pages:
-        result = get_sold_single_page(log, page)
+    while page <= MAX_ANON_PAGE:
+        orders, total_count, _ = get_sold_single_page(log, page)
 
-        if result[0] is None:
-            log.error(f"Error fetching page {page}, stopping...")
+        if not orders:
+            log.info(f"Page {page} 无数据,停止")
             break
 
-        orders, total_count = result
-        total_collected += len(parse_list_results(log, orders, sql_pool))
+        inserted_count = len(parse_list_results(log, orders, sql_pool))
+        total_collected += inserted_count
 
-        # 第一页时打印总数
         if page == 1:
-            total_pages = (total_count + page_size - 1) // page_size
-            log.info(f"Total items: {total_count}, Total pages: {total_pages}")
+            log.info(f"Total items(接口累计已售): {total_count}")
+
+        oldest = parse_deal_date(orders[-1])
+        log.info(f"Page {page} processed, fetched={len(orders)}, parsed={inserted_count}, 本页最旧成交日={oldest}")
 
-        if len(orders) < page_size:
-            log.info(f"Less than {page_size} items on page {page}, stopping...")
+        # 已翻到早于 cutoff 的数据,增量部分采集完毕
+        if oldest and oldest < cutoff:
+            log.info(f"本页最旧成交日 {oldest} < cutoff {cutoff},增量采集完成,停止")
+            reached_cutoff = True
             break
 
-        # 检查是否还有下一页
-        if page * page_size >= total_count:
-            log.info(f"No more pages, stopping...")
+        # 末页(接口返回不足一页)
+        if len(orders) < PAGE_SIZE:
+            log.info(f"Page {page} 不足 {PAGE_SIZE} 条,已到列表末尾,停止")
+            reached_cutoff = True
             break
 
         page += 1
 
-    log.info(f"Total orders collected: {total_collected}")
+    if not reached_cutoff:
+        log.warning(
+            f"已翻到匿名上限第 {MAX_ANON_PAGE} 页仍未追到 cutoff {cutoff},"
+            f"说明近 {days_back} 天新增已售 > {MAX_ANON_PAGE * PAGE_SIZE} 条,"
+            f"未登录无法继续深翻,本轮可能有遗漏——建议提高采集频率或改带登录 token 深翻")
+
+    log.info(f"Total orders collected(含重复忽略前): {total_collected}")
     return total_collected