Quellcode durchsuchen

feat(spider): 新增板块商户发现爬虫补充店铺数据

- 在 mysql_pool.py 中改进 _execute 方法,增加断连重试机制,提升数据库操作稳定性
- 优化对完整性错误的日志处理,避免日志堆栈污染
- 新增 zc_board_shop_spider.py 模块,实现从多个活动板块接口获取店铺信息
- 实现跨页、跨板块去重逻辑,避免重复写库
- 根据板块接口限制,采用 INSERT IGNORE 避免覆盖已有真实销量数据
- 在 zc_new_daily_spider.py 中集成板块商户发现功能,定期补充漏采店铺
- 调整主流程查询条件,取消对 is_deleted 字段的筛选,确保覆盖更多店铺数据
charley vor 1 Monat
Ursprung
Commit
ebdba89dd1
3 geänderte Dateien mit 296 neuen und 22 gelöschten Zeilen
  1. 61 18
      zc_spider/mysql_pool.py
  2. 222 0
      zc_spider/zc_board_shop_spider.py
  3. 13 4
      zc_spider/zc_new_daily_spider.py

+ 61 - 18
zc_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,

+ 222 - 0
zc_spider/zc_board_shop_spider.py

@@ -0,0 +1,222 @@
+# -*- coding: utf-8 -*-
+# Author : Charley
+# Python : 3.12.10
+# Date   : 2026/06/29
+
+"""板块商户发现爬虫(getActList 多板块补充店铺)。
+
+背景:原 get_shop_list 走 getMerMyList(我的商户列表),覆盖不全,会漏掉部分店铺。
+本模块改从 5 个活动板块抓活动列表(接口同为 /zc-api/act/actProduct/getActList,
+仅请求体里的 cateId 不同),从每条活动里解析出店铺信息(merNo / merName),
+补充写入 zc_shop_record,从而把 get_shop_list 漏采的店铺也覆盖到。
+
+注意:板块接口的活动数据里没有销量(spell_number)和粉丝(attentionNumber)字段,
+按主公确认的方案,解析逻辑照搬 zc_new_daily_spider.parse_shop_data —— 这两个字段
+缺失时会写入 None,并沿用 ON DUPLICATE KEY UPDATE 覆盖。
+
+板块接口(getActList)经实测为公开接口,不校验 Authorization,故本模块不需要 token。
+
+对外入口:get_board_shop_list(log, sql_pool),可在 zc_new_daily_spider.zc_main 中直接调用。
+"""
+
+# 基础配置(与 zc_new_daily_spider 保持一致)
+BASE_URL = "https://cashier.yqszpay.com"
+PAGE_SIZE = 10
+
+# 5 个板块的分类ID。来源:抓包数据.txt 的 5 条 getActList 请求体(AES 解密后得到)。
+# 每个板块对应一个 cateId,接口与其余参数完全相同,仅 cateId 不同:
+#   "0"      -> 板块1 宝可梦
+#   "3"      -> 板块2 篮球
+#   "0,6,3"  -> 板块3(多分类聚合) 精选
+#   "1"      -> 板块4 足球
+#   "11"     -> 板块5 综合
+BOARD_CATE_IDS = ["0", "3", "0,6,3", "1", "11"]
+
+# 单个板块翻页保护上限,防止异常情况下无限翻页
+MAX_PAGES = 100
+
+
+def get_board_single_page(log, cate_id, page_num, page_size=PAGE_SIZE):
+    """获取单个板块的活动列表(支持翻页)。
+
+    复用 zc_new_daily_spider.make_encrypted_post_request 完成加密 / 重试 / 代理 / 解密,
+    采用函数内惰性导入,避免与主模块产生循环导入。
+
+    注意:板块接口经实测为公开接口,不校验 Authorization(无 token / 空 token / 假 token
+    返回结果完全一致),故不传 token。
+
+    Args:
+        log: 日志对象,从调用方透传。
+        cate_id (str): 板块分类ID,取值见 BOARD_CATE_IDS。
+        page_num (int): 页码,从 1 开始。
+        page_size (int, optional): 每页条数。Defaults to PAGE_SIZE。
+
+    Returns:
+        dict | None: 解密后的响应字典(含 rows / total);请求失败时返回 None。
+    """
+    # 惰性导入:调用发生在 zc_main 运行期,此时主模块已完全加载,不会触发循环导入
+    from zc_new_daily_spider import make_encrypted_post_request
+
+    log.debug(f"Getting board list, cateId: {cate_id}, page: {page_num}")
+    url = f"{BASE_URL}/zc-api/act/actProduct/getActList"
+    # 请求体字段与抓包解密结果一一对应:
+    #   cateId    -> 板块分类ID(区分 5 个板块的唯一参数)
+    #   actStatus -> 1 表示只取进行中的活动(抓包固定为 1)
+    #   finish / loading -> 前端分页状态位,抓包固定 false / true,照搬即可
+    request_data = {
+        'cateId': cate_id,
+        'pageSize': page_size,
+        'pageNum': page_num,
+        'finish': False,
+        'loading': True,
+        'actStatus': 1
+    }
+    try:
+        # 板块接口不需要 token,仅透传随页变化的 pageNum 头
+        resp = make_encrypted_post_request(
+            log, url, request_data,
+            extra_headers={"pageNum": str(page_num)}
+        )
+    except Exception as e:
+        log.error(f"Error getting board list (cateId={cate_id}, page={page_num}): {e}")
+        resp = None
+    return resp
+
+
+def parse_board_shop_data(log, items, sql_pool, seen=None):
+    """解析板块活动里的店铺数据并写库(INSERT IGNORE,已存在店铺不改动)。
+
+    板块接口(getActList)经实测不返回销量(spell_number)和粉丝(attentionNumber)字段,
+    只能拿到 merNo / merName。因此本函数写 shop_id + shop_name + sold_number=1,用
+    INSERT IGNORE 跳过已存在的店铺——避免把 get_shop_list(getMerMyList,含真实
+    sold_number / fans)采到的数据覆盖掉。
+
+    sold_number 写死哨兵值 1,是为了让板块新发现的店铺能被每日任务
+    「SELECT ... WHERE sold_number != 0」选中进入商品采集管道(板块拿不到真实销量);
+    它只对新插入的店铺生效,已存在店铺被 IGNORE 跳过。fans 走表默认值。
+
+    Args:
+        log: 日志对象,从调用方透传。
+        items (list[dict]): getActList 返回的活动列表(每条含 merNo / merName)。
+        sql_pool (MySQLConnectionPool): MySQL 连接池。
+        seen (set, optional): 已写入过的 shop_id 集合,用于跨页 / 跨板块去重,避免重复写库。
+            Defaults to None(不去重)。
+
+    Returns:
+        int: 本次实际写入的店铺数量。
+    """
+    log.debug(f"Parsing board shop data...........")
+    info_list = []
+    for item in items:
+        shop_id = item.get('merNo')
+        shop_name = item.get('merName')
+
+        # 同一店铺会出现在多条活动 / 多个板块里,去重以减少无谓的写库
+        if seen is not None:
+            if shop_id in seen:
+                continue
+            seen.add(shop_id)
+
+        # 板块接口拿不到真实 sold_number / fans。
+        # sold_number 写死哨兵值 1:让板块新发现的店铺能被每日任务
+        # 「SELECT ... WHERE sold_number != 0」选中,从而进入商品采集管道。
+        # 注意:这是“需采集”标记而非真实销量;配合 INSERT IGNORE,已存在店铺会被整行跳过、
+        # 不受影响,若该店日后被 get_shop_list 覆盖到会刷回真实值。fans 仍走表默认值。
+        data_dict = {
+            'shop_id': shop_id,
+            'shop_name': shop_name,
+            'sold_number': 1
+        }
+        log.debug(f"Parsed board shop data: {data_dict}")
+        info_list.append(data_dict)
+
+    if not info_list:
+        return 0
+
+    # INSERT IGNORE:店铺已存在则跳过,不覆盖 get_shop_list 采到的真实 sold_number / fans
+    sql_pool.insert_many(table='zc_shop_record', data_list=info_list, ignore=True)
+    return len(info_list)
+
+
+def get_single_board_shops(log, cate_id, sql_pool, seen):
+    """抓取单个板块的全部店铺(翻页直至无数据)。
+
+    翻页逻辑参考 get_shop_list:逐页请求,解析写库,遇到空页或不足整页即停止,
+    并受 MAX_PAGES 上限保护。
+
+    Args:
+        log: 日志对象,从调用方透传。
+        cate_id (str): 板块分类ID,取值见 BOARD_CATE_IDS。
+        sql_pool (MySQLConnectionPool): MySQL 连接池。
+        seen (set): 跨板块共享的 shop_id 去重集合。
+
+    Returns:
+        int: 本板块新写入的店铺数量(已去重)。
+    """
+    page_num = 1
+    board_count = 0
+
+    while page_num <= MAX_PAGES:
+        result = get_board_single_page(log, cate_id, page_num, PAGE_SIZE)
+        if result is None:
+            log.error(f"板块 {cate_id} 第 {page_num} 页请求失败,停止翻页")
+            break
+
+        data_list = result.get('rows', [])
+        if len(data_list) == 0:
+            log.info(f"板块 {cate_id} 第 {page_num} 页无数据,停止翻页")
+            break
+
+        board_count += parse_board_shop_data(log, data_list, sql_pool, seen=seen)
+        log.info(f"板块 {cate_id} 第 {page_num} 页完成,本页活动数: {len(data_list)}")
+
+        # 不足整页说明已到末页,停止翻页
+        if len(data_list) < PAGE_SIZE:
+            log.info(f"板块 {cate_id} 第 {page_num} 页不足整页,已到末页,停止翻页")
+            break
+
+        page_num += 1
+
+    return board_count
+
+
+def get_board_shop_list(log, sql_pool):
+    """板块商户发现入口:遍历 5 个板块,解析并补充店铺到 zc_shop_record。
+
+    对外暴露的唯一入口函数,可在 zc_new_daily_spider.zc_main 中直接调用。
+    跨 5 个板块共享同一个去重集合,确保同一店铺只写库一次。
+    板块接口为公开接口,无需 token。
+
+    Args:
+        log: 日志对象,从调用方透传(与 zc_main 同一个 logger)。
+        sql_pool (MySQLConnectionPool): MySQL 连接池,由 zc_main 创建后透传。
+
+    Returns:
+        int: 5 个板块合计去重后新写入的店铺数量。
+    """
+    log.info(f"开始板块商户发现,共 {len(BOARD_CATE_IDS)} 个板块: {BOARD_CATE_IDS}")
+    seen = set()  # 跨板块去重,同一 shop_id 只写库一次
+    total_count = 0
+
+    for cate_id in BOARD_CATE_IDS:
+        try:
+            count = get_single_board_shops(log, cate_id, sql_pool, seen)
+            total_count += count
+            log.info(f"板块 {cate_id} 采集完成,新增店铺(去重后): {count}")
+        except Exception as e:
+            log.error(f"板块 {cate_id} 采集异常: {e}")
+
+    log.info(f"板块商户发现结束,5 个板块合计去重店铺数: {len(seen)},写库次数: {total_count}")
+    return total_count
+
+
+
+if __name__ == '__main__':
+    from loguru import logger
+    from mysql_pool import MySQLConnectionPool
+
+    pool = MySQLConnectionPool(log=logger)
+    try:
+        get_board_shop_list(logger, pool)
+    except Exception as e:
+        logger.error(f'get_board_shop_list error: {e}')

+ 13 - 4
zc_spider/zc_new_daily_spider.py

@@ -10,6 +10,7 @@ import user_agent
 from loguru import logger
 from crypto_utils import CryptoHelper
 from mysql_pool import MySQLConnectionPool
+from zc_board_shop_spider import get_board_shop_list
 from tenacity import retry, stop_after_attempt, wait_fixed
 
 logger.remove()
@@ -393,9 +394,9 @@ def get_sold_list(log, shop_id, token, sql_pool, shop_name):
         data_list = result.get('rows', [])
         if not data_list:
             log.info(f"第 {page_num} 页无数据,停止翻页")
-            log.info(f'该店铺{shop_name}无数据, 修改为店铺注销状态..........')
-            # 更新店铺状态为注销状态
-            sql_pool.update_one("UPDATE zc_shop_record SET is_deleted = 1 WHERE shop_id = %s", (shop_id,))
+            # log.info(f'该店铺{shop_name}无数据, 修改为店铺注销状态..........')
+            # # 更新店铺状态为注销状态
+            # sql_pool.update_one("UPDATE zc_shop_record SET is_deleted = 1 WHERE shop_id = %s", (shop_id,))
             break
         parse_sold_data(log, token, data_list, sql_pool, shop_name)
 
@@ -475,6 +476,12 @@ def zc_main(log):
         # player test
         # has_data = get_player_list(log, 1800, token, sql_pool)
 
+        # 获取首页板块商户列表
+        try:
+            get_board_shop_list(log, sql_pool)
+        except Exception as e:
+            log.error(f'get_board_shop_list error: {e}')
+
         # 获取shop data
         try:
             get_shop_list(log, sql_pool)
@@ -487,7 +494,9 @@ def zc_main(log):
         try:
             # 从 shop 表查询所有 merNo  2026/5/9增加is_deleted判断是否注销店铺
             mer_no_rows = sql_pool.select_all(
-                "SELECT shop_id, shop_name FROM zc_shop_record WHERE is_deleted = 0 AND sold_number != 0")
+                # "SELECT shop_id, shop_name FROM zc_shop_record WHERE is_deleted = 0 AND sold_number != 0"
+                "SELECT shop_id, shop_name FROM zc_shop_record WHERE sold_number != 0"
+            )
             log.info(f"查询到 {len(mer_no_rows)} 个商户编号: {mer_no_rows}")
             for shop_id, shop_name in mer_no_rows:
                 log.info(f"开始爬取商户 {shop_id}, {shop_name} 的商品数据")