Prechádzať zdrojové kódy

feat(buy_record_analysis): 新增得卡DECA购买记录常驻采集脚本及MySQL连接池模块

- 新增 buy_record_spider.py,实现多商品自适应频率动态采集购买记录
- 购买记录基于倒序滑动窗口序列对齐实现高效去重,解决重复采集问题
- 实现全链路免token获取详情和在售列表,支持使用代理避免限流
- 支持快车商品密集采集,动态计算下次采集间隔,防止漏采
- 监控在售商品列表,自动移除已售罄或结束售卖商品,并加入黑名单避免死循环
- 复用 daily 模块的全站在售采集逻辑,定期刷新监控商品队列及商品详情
- 采集进度快照写入进度时间序列表,提供分钟级售卖统计数据
- 引入基于dbutils的MySQL连接池实现,增强数据库连接稳定性和性能
- connection pool 支持断连重试及异常处理,保证查询和插入操作可靠
- 新增 MySQL 相关配置,支持通过环境变量覆盖默认参数
charley 1 mesiac pred
rodič
commit
8b21d9196b

+ 206 - 0
deca_spider/buy_record_analysis/README.md

@@ -0,0 +1,206 @@
+# 得卡 DECA 商品购买记录采集
+
+> 归属:`D:\work\2026-08-02(deca_spider)\buy_record_analysis\`
+> 平台:得卡 DECA App(`api.decalive.com`)
+> 最后更新:2026/08/06
+> 依赖父项目:签名/Token/请求/代理层复用根目录 `deca_sold_core.py`;数据库配置读根目录 `application.yml`;Token 持久化在根目录 `token.json`。
+
+---
+
+## 1. 目标信息
+
+- **App / URL**:得卡 DECA,业务域名 `https://api.decalive.com`
+- **技术栈**:Python 3.12.10 + `requests` + `tenacity` + `loguru` + `charley-utils`(MySQL 连接池)+ 快代理隧道
+- **数据目标**:把商品详情页「购买记录」板块的**买家名单**(`user_id / 脱敏昵称 / 购买份数 / 购买时间`)持续累积到 MySQL 表 `deca_buy_record`,绕开单次响应只有 10 条的限制。
+
+---
+
+## 2. 请求分析
+
+### 2.1 关键接口
+
+`POST /api/v1/app/groupbuy/detail`,请求体:
+
+```json
+{"code": "GB26080417234"}
+```
+
+响应里 `data.purchaseRecords` 是**最近 10 条**购买记录数组,每项:
+
+| 字段 | 说明 | 采集处理 |
+|---|---|---|
+| `userId` | 买家用户 ID | **不脱敏**,作为唯一键的一部分入库 |
+| `nickname` | 买家昵称 | **服务端已脱敏**(`D**l` / `匿名用户` / `*`);原样保留 |
+| `avatarUrl` | 头像 URL | 不入库 |
+| `purchasedAt` | 相对时间文本 | 存原文 + 反推的绝对时间戳 |
+| `cardCount` | 本次购买份数 | 入库;构成唯一键的一部分 |
+
+### 2.2 登录态 / 签名 / 代理
+
+**关键:详情接口 `groupbuy/detail` 无需登录**(实测不带 token 也返回 code=0 + 10 条),所以**白名单模式全程免登录**。只有**商家模式**要拉的在售列表 `on-sale-list` 需要 token(不带会 `code=10002 未登录`),故仅该模式启动时 `ensure_token`。
+
+| 接口 | 用途 | 是否需要 token |
+|---|---|---|
+| `groupbuy/detail` | 拉购买记录(两种模式都用) | ❌ 免登录 |
+| `groupbuy/merchant/on-sale-list` | 商家模式拉在售商品列表 | ✅ 需要 |
+
+签名细节见根目录 `README.md`。本项目通过 `deca_sold_core` 复用(`core.USE_PROXY=True` 开代理;`throttled_do_request` 里 detail 传 `need_auth=False`、on-sale-list 传 `need_auth=True`)。
+
+---
+
+## 3. 关键发现(为什么必须"轮询累积")
+
+反编译产物 + 真机实测双向验证得出:
+
+1. **不存在独立的购买记录接口**:全库枚举 `api/v1/...` 路径 90+ 个,`groupbuy/` 前缀下没有 `purchase-records` / `buy-record` / `purchase/list` 等端点;`GroupBuyPurchaseRecordDto` 只在详情响应里被引用。
+2. **详情接口不接受分页参数**:请求体类 `GroupBuyProductDetailRequest` 只有 `code` 一个字段;实测把 `page` / `pageSize` / `limit` / `count` / `purchaseRecordLimit` 等 15 种试探参数塞进 body,服务端 100% 忽略、恒定返回 10 条。
+3. **响应无分页游标**:`purchaseRecords` 旁边没有 `hasMore` / `nextCursor` / `total`。
+
+**结论**:客户端合法路径无法一次拿到 >10 条。**唯一可行方式 = 周期性拉详情 + 应用层去重累积**。
+
+**顺带发现**:`open-card-report/public/list`(拆卡报告)的 `hit_user_nickname` 字段**不脱敏**(`小柏` / `Windowsky76` / `汉堡`),但那里没有 `user_id`,无法建立 `user_id → 真实昵称` 的反查通路。想拿真实昵称需另找「根据 userId 查用户资料」接口。
+
+---
+
+## 4. 数据模型
+
+**表名**:`deca_buy_record`(DDL 已并入 `schema.sql`,脚本启动时自动 `CREATE TABLE IF NOT EXISTS`)
+
+| 字段 | 类型 | 说明 |
+|---|---|---|
+| `id` | bigint PK | 自增主键 |
+| `product_code` | varchar(32) | 商品编码,如 `GB26080417234` |
+| `user_id` | varchar(32) | 买家用户 ID(真值) |
+| `nickname` | varchar(64) | 服务端脱敏昵称 |
+| `card_count` | int | 本次购买份数 |
+| `purchased_at_text` | varchar(32) | 首次抓到时的相对时间原文,如 `1分钟前` |
+| `purchased_at_ts` | bigint | 反推的绝对购买时间戳(秒),**去重键之一** |
+| `purchased_at` | datetime | 反推的绝对购买时间(可读格式) |
+| `first_seen_at` | datetime | 我们首次抓到该记录的时间 |
+| `gmt_create_time` | datetime | 默认 `CURRENT_TIMESTAMP` |
+| `gmt_modified_time` | datetime | 默认 `CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP` |
+
+**唯一键**:`uk_buy (product_code, user_id, card_count, purchased_at_ts)`(DB 兜底),**主判重逻辑在应用层**。
+
+### 去重逻辑演进(踩了两次坑)
+
+**坑一 · 分钟桶漂移(100% 重复)**:最早用 `(product_code, user_id, card_count, purchased_at_minute)`,测试每人重复入库 2 次。原因:相对时间每 60 秒 +1("21分钟前"→"22分钟前"),反推时间戳漂移 60 秒,跨到下一个分钟桶。
+
+**坑二 · 买家名单粒度(漏单)**:改成 `(product_code, user_id, card_count)` 后零重复,但**会漏掉同一用户对同一商品的多笔相同份数订单**——实测反例:`王**要 x10` 在 10/12 分钟前各一笔、`D**F x6` 在 37/40 分钟前各一笔,都被合并成一条。
+
+**当前解法 · 订单粒度 + 反推时间戳 ±容差判重**:
+
+利用一个数学性质——
+
+- **同一笔订单**在多轮采集里反推的 `purchased_at_ts` 始终落在 `[真实购买时刻, +60秒)`,即漂移永远 **< 60 秒**。
+- **两笔不同订单**(间隔 ≥ 2 分钟 = 120 秒)反推时间戳之差 ∈ `(60, 180)` 秒,即 **≥ 60 秒**。
+
+判重规则(`DEDUP_TOLERANCE_SEC = 70`):新记录的 `purchased_at_ts` 与库里同 `(product_code, user_id, card_count)` 的已有 ts **相差 ≤ 70 秒 → 判为漂移,跳过;> 70 秒 → 判为新订单,入库**。
+
+实现:进程内维护 `_order_idx` 内存索引 `{(code,uid,cc): [ts...]}`;商品首次纳入监控时用 `load_order_index` 从库恢复(避免重启后把老订单当新的)。DB 唯一键含精确 `purchased_at_ts` 仅作并发/异常兜底。
+
+**已知局限**:同一个人在 **~1 分钟内**下两笔**份数完全相同**的单,会被当成一笔(purchasedAt 只精确到分钟,数据上无法与漂移区分)。此情形极罕见。
+
+---
+
+## 5. 文件清单
+
+| 文件 | 职责 |
+|---|---|
+| `buy_record_spider.py` | ⭐ 正式采集脚本:多商品自适应频率、常驻长跑;白名单模式下即单商品测试 |
+| `README.md` | 本文档 |
+
+数据库 DDL 追加在项目根 `schema.sql`(第 204 行起)。
+
+---
+
+## 6. `buy_record_spider.py` 使用说明
+
+### 6.1 两种工作模式(切换只改一个变量)
+
+```python
+# 白名单模式:非空 → 只监控这些 code,跳过 on-sale-list
+WATCH_CODES = ["GB26080417234"]
+
+# 商家模式:WATCH_CODES 为空 → 拉 MERCHANT_ID 的全部在售商品,动态增删
+# WATCH_CODES = []
+MERCHANT_ID = "881226408"
+```
+
+- 单商品测试 → 白名单里放 1 个 code
+- 多商品自选 → 白名单里加多个 code
+- 商家全在售 → 清空 `WATCH_CODES`
+
+### 6.2 自适应频率
+
+```python
+MIN_INTERVAL_SEC = 0      # 单商品最小采集间隔(测试期设 0,未来看限流表现再抬)
+MAX_INTERVAL_SEC = 600    # 单商品最大采集间隔(10 分钟)
+GLOBAL_MIN_GAP_SEC = 0.2  # 全局相邻两次请求最小间隔(≈每秒 5 次上限,防瞬时限流)
+```
+
+**算法**:每次拉完,取 10 条 `purchasedAt` 反推的时间戳跨度 `span_sec`,下次间隔 = `span_sec // 3`,硬夹在 `[MIN, MAX]`。
+
+**典型档位**:
+
+| 10 条时间跨度 | 计算 span/3 | 实际下次间隔 |
+|---|---|---|
+| 3s 内(超级极热) | 1s | 立即再刷(受 GLOBAL_MIN_GAP_SEC 兜底 ≥ 0.2s) |
+| 30s 左右 | 10s | 10s 后 |
+| 5 分钟 | 100s | ~1.5 分钟后 |
+| 30 分钟 | 600s | 触顶 10 分钟 |
+| 空数据(新品无购买) | — | 触顶 10 分钟 |
+
+### 6.3 运行
+
+**前置**:
+- 项目根有 `application.yml`(数据库配置)。
+- MySQL 里 **`deca_buy_record` 表需先建好**:`python init_db.py` 或在客户端跑 `schema.sql` 里那段 DDL(脚本不再自动建表)。
+- 白名单模式**不需要** `token.json`;商家模式(`WATCH_CODES=[]`)才需要根目录有有效 `token.json`。
+
+```bash
+# 前台跑(关终端就停,适合调试)
+cd D:/work/2026-08-02(deca_spider)
+python buy_record_analysis/buy_record_spider.py
+```
+
+```bash
+# 后台跑(关终端不停,Windows Git Bash / WSL)
+cd D:/work/2026-08-02(deca_spider)
+nohup python buy_record_analysis/buy_record_spider.py > logs/buy_record.out 2>&1 &
+tail -f logs/buy_record.out
+```
+
+**推荐**:`tmux` / `screen` 里跑,锁屏 / 断开都不影响,需要时随时 `attach`。
+
+**日志**:`logs/{YYYYMMDD}_buy_record.log`(按天切分,保留 7 天),stderr 同步 INFO 便于观察。
+
+### 6.4 出问题时的重试语义
+
+- **单商品单次采集失败** → 该商品 60s 后再试,不阻塞其他商品
+- **`main_task` 挂掉**(如数据库连接池异常)→ `tenacity @retry(stop=100, wait=3600)`,每小时重试一次,最多 100 次
+- **Token 续期**(仅商家模式):`ensure_token` 自动检查过期 → 走 `refreshToken` 续签(900s / 次);白名单模式不涉及
+
+---
+
+## 7. 注意事项 / 踩坑
+
+1. **昵称脱敏是服务端行为**,客户端拿不到真实名(`D**l` 是接口返回的原样)。要真实昵称得另找「userId → 用户资料」接口,本项目未做。
+2. **去重踩过两次坑**(详见第 4 节演进,🔴 已根治):
+   - 分钟桶漂移 → 100% 重复入库
+   - 名单粒度合并 → 漏掉同用户多笔相同份数订单
+   最终解:`(product_code, user_id, card_count)` 分组 + 反推 `purchased_at_ts` ±70s 应用层判重 + 内存索引重启从库恢复。
+3. **停机 = 断档**:进程挂了、机器重启、网络断了那段时间的买家**永久丢失**(10 条窗口早滚走了)。生产必须用 `nohup` / `tmux` / `supervisor` 保活。
+4. **售罄后接口依然返回**:只是 `purchaseRecords` 不再新增。白名单模式不会自动移除;商家模式靠 on-sale-list 消失自动移除。
+5. **限流**:目前 `MIN_INTERVAL_SEC=0` + `GLOBAL_MIN_GAP_SEC=0.2` 是"测试期激进"配置。若出现频繁 HTTP 非 200 / 服务端拦截,先抬 `MIN_INTERVAL_SEC` 到 3–5,再抬 `GLOBAL_MIN_GAP_SEC` 到 0.3–0.5。
+6. **代理必开**:`core.USE_PROXY = True`(快代理隧道),避免高频请求触发 IP 风控。
+
+---
+
+## 8. 举一反三
+
+- 任何"只显示最近 N 条 + 会滚动"的板块(评论、直播弹幕、观众列表、消息流)都能套本方案:**周期轮询 + 应用层去重累积**。
+- **数据源只精确到分钟、但需要区分"重复推送"和"真实新事件"** 是这类场景的通用难点。核心思路:从相对时间反推绝对时间戳时,量化误差有上界(本项目里 < 60s),只要**真实事件的间隔大于这个上界**,就能靠"±容差判重"完美切开重复与真实。相对时间只精确到分钟 → 上界 60s → 阈值取 60~90s;如果数据源精确到 X 秒,阈值 = X ~ 1.5X。
+- **唯一键设计**:稳定值(`user_id` / `item_id` / `msg_id`)必进;带量化误差的时间戳(如本项目的 `purchased_at_ts`)**可以进 DB 唯一键作兜底**,但**不能作为主判重逻辑**——主判重必须走应用层的 ±容差匹配,否则会出现"分钟桶漂移"型重复。
+- **自适应频率**(`interval = span/N`)在别的滚动场景同样适用;`N=3` 是经验值,可以按业务调 —— N 越小越激进,越容易撞限流。
+- 未来若发现"根据 userId 查真实资料"的接口(比如 `user/detail`),可在本表基础上加个 `real_nickname` 列,做一次批量补拉;或者关联查 `deca_report_record.hit_user_nickname`(那边不脱敏)—— 但目前两表间无 `user_id` 关联字段,只能靠昵称文本模糊 join。

+ 98 - 0
deca_spider/buy_record_analysis/YamlLoader.py

@@ -0,0 +1,98 @@
+# -*- coding: utf-8 -*-
+# Author : Charley
+# Python : 3.12.10
+# Date   : 2025/12/22 10:44
+import os, re
+import yaml
+
+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])
+            group = match.groupdict()
+            if group['ENV'] is not None:
+                env = group['ENV'][:-1]
+                return os.getenv(env, group['VAL'])
+            return None
+        except:
+            return self.config[key]
+
+    def getValueAsInt(self, key: str):
+        try:
+            match = regex.match(self.config[key])
+            group = match.groupdict()
+            if group['ENV'] is not None:
+                env = group['ENV'][:-1]
+                return int(os.getenv(env, group['VAL']))
+            return 0
+        except:
+            return int(self.config[key])
+
+    def getValueAsBool(self, key: str):
+        try:
+            match = regex.match(self.config[key])
+            group = match.groupdict()
+            if group['ENV'] is not None:
+                env = group['ENV'][:-1]
+                return bool(os.getenv(env, group['VAL']))
+            return False
+        except:
+            return bool(self.config[key])
+
+
+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):
+        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 = real_path.rsplit('.', 1)
+        profiledYaml = f'{result[0]}-{profile}.{result[1]}'
+        if os.path.exists(profiledYaml):
+            with open(profiledYaml, encoding='utf-8') as fd:
+                conf.update(yaml.load(fd, Loader=yaml.FullLoader))
+
+    return YamlConfig(conf)

+ 6 - 0
deca_spider/buy_record_analysis/application.yml

@@ -0,0 +1,6 @@
+mysql:
+  host: ${MYSQL_HOST:100.64.0.25}
+  port: ${MYSQL_PROT:3306}
+  username: ${MYSQL_USERNAME:crawler}
+  password: ${MYSQL_PASSWORD:Pass2022}
+  db: ${MYSQL_DATABASE:crawler}

+ 672 - 0
deca_spider/buy_record_analysis/buy_record_spider.py

@@ -0,0 +1,672 @@
+# -*- coding: utf-8 -*-
+# Author : Charley
+# Python : 3.12.10
+# Date   : 2026/08/05
+"""得卡 DECA 购买记录常驻采集脚本(多商品自适应频率)。
+
+功能:
+  一个进程干两件事(方案A,2026/08/08 合并):
+  (1) 每分钟复用 daily 的免 token home/search 拉全站在售,落 deca_onsale_* 三张表(供 deca_on_sale_report 查库出报告);
+  (2) 从库里取本商家(MERCHANT_ID)在售 code,采其购买记录写 deca_buy_record。
+  购买记录按「近 10 条购买的时间跨度」动态调节采集频率——
+  卖得快(跨度短)密采、卖得慢(跨度长)稀采;商品下线、售卖结束或售罄自动停采(含白名单模式)。
+  数据写 deca_buy_record;去重用「倒序滑动窗口序列对齐」:purchaseRecords 是最近 10 条按时间倒序、
+  新单只从顶部推入,故按「本次窗口 (userId,cardCount) 序列相对上次整体下移了几位」求出顶部新增的 k 条,
+  只入库这 k 条——不依赖会漂移的反推时间戳,且能正确区分同用户多笔相同份数的单(靠位移而非时间)。
+  兜底:进程重启首轮 / 窗口整体换新(间隔内卖出≥10 笔)无法对齐时,退回反推时间戳「动态容差」判重
+  (秒 70s / 分钟 90s / 小时 3660s / 天 90000s)+ DB 唯一键,无重叠时告警疑似漏采。
+
+停采(重点):
+  详情返回 availableStock<=0(售罄)或已过 saleEndAt(到结束时间)即判定售卖结束,移出监控并加入
+  _ended_codes 黑名单,避免白名单模式下对已结束商品死循环重采「固定的最后 10 条」造成重复入库。
+
+自适应节奏:
+  下次间隔 ≈ 10 条跨度 / 3,硬夹在 [MIN_INTERVAL_SEC, MAX_INTERVAL_SEC]。
+  典型档位(当前 MIN=0,测试期让服务端限流自己说话):
+    - 极热(10 条 3s 内)→ 立即再刷(由 GLOBAL_MIN_GAP_SEC 兜底≥0.2s)
+    - 热(10 条 1min 内)→ 20s 一次
+    - 中(10 条 5min 内)→ 100s 一次
+    - 慢(10 条 30min 内)→ 600s 一次触顶
+
+登录态:
+  - **全链路免 token**:详情接口 groupbuy/detail 免登录;在售列表改用免 token 的 home/search(复用 daily),
+    落库后查库拿商家在售 code,不再走需 token 的 on-sale-list。彻底摆脱登录/验证码。
+
+保护:
+  - 全局请求节流 GLOBAL_MIN_GAP_SEC,避免瞬时高并发触发限流
+  - main_task @retry(stop=100, wait=3600):挂了每小时重试,跑到手动停
+  - 快代理隧道请求(走 deca_sold_core)
+
+前置:
+  - 表 deca_buy_record 及 deca_onsale_* 需先建好(schema.sql)。
+  - 依赖 on_sale/deca_on_sale_daily_spider.py 的复用函数(get_shop_list/get_onsale_products/fill_product_details);
+    该文件保留、不再单独常驻跑。报告由 deca_on_sale_report.py 独立定时查库生成。
+
+运行:项目根目录 `python buy_record_analysis/buy_record_spider.py`(采集);
+     报告另跑 `python on_sale/deca_on_sale_report.py`(可加 loop 定时,见该文件)。
+"""
+import os
+import re
+import sys
+import time
+from datetime import datetime
+
+# 挂靠项目根:复用核心签名/token/请求/代理层,让 application.yml、token.json 生效
+_ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
+sys.path.insert(0, _ROOT)
+os.chdir(_ROOT)
+
+from loguru import logger
+from tenacity import retry, stop_after_attempt, wait_fixed
+import deca_sold_core as core
+from mysql_pool import MySQLConnectionPool
+# 复用 daily 的免 token 采集/落库(home/search 全站在售 → deca_onsale_* 三张表);daily 文件保留、不再单独常驻
+# 注:本 import 会触发 daily 模块级 logger 配置,但下方 buy_record 的 logger.remove/add 在其后,会覆盖回本脚本日志
+from on_sale.deca_on_sale_daily_spider import get_shop_list, get_onsale_products, fill_product_details
+
+# ==================== 配置 ====================
+core.USE_PROXY = True               # 详情接口(groupbuy/detail)走快代理隧道;全链路免 token
+
+MERCHANT_ID = "881226408"           # 监控商家:采其在售商品的购买记录(WATCH_CODES 为空时生效)
+WATCH_CODES = []                    # 指定商品白名单:非空则只盯这些 code、跳过在售发现;空则走 MERCHANT_ID 全在售(默认)
+ONSALE_INGEST_SEC = 60             # 全站在售落库 + 刷新监控队列间隔(1 分钟,复用 daily 免 token home/search)
+SHOP_DISCOVER_SEC = 600            # 商家发现 + 今日新增补详情间隔(10 分钟,变化慢无需每分钟)
+MIN_INTERVAL_SEC = 0               # 单商品购买记录最小采集间隔(测试期设 0:让实测数据决定是否要抬)
+MAX_INTERVAL_SEC = 600             # 单商品购买记录最大采集间隔(10 分钟)
+# 快车(小批量拼团)专用:份数小的商品售罄极快(实测最快 10 份滑窗约 64s),需比 MAX_INTERVAL_SEC 更小的上限密采,
+# 否则「刚上架 span=0 → 排 600s」会让 3~8 分钟售罄的拼团在两次采集间隙整场漏采(2026/08/10 修复)
+FAST_LANE_MAX_COUNT = 100          # 份数<=此值视为「快车」(本商家快车=30份、大车=313+),涵盖未来 20/40/50 份的小批量快车
+FAST_LANE_MAX_INTERVAL_SEC = 20    # 快车最大采集间隔(实测最快 10 份滑窗~64s,20s 留约 3x 余量)
+GLOBAL_MIN_GAP_SEC = 0.2           # 全局相邻两次详情请求最小间隔(≈每秒 5 次上限)
+IDLE_SLEEP_SEC = 1                 # 主循环空转 sleep(没到点时)
+DEDUP_TOLERANCE_SEC = 70           # 秒级基准判重容差:同(商品,user,份数)下反推时间戳相差<=此值视为同一笔;小时/天级按 _TOL_BY_UNIT 放大
+
+TABLE = "deca_buy_record"
+ONSALE_TABLE = "deca_onsale_product_record"   # 在售商品表:daily 复用函数落库,本脚本查库拿商家在售 code
+PROGRESS_TABLE = "deca_onsale_product_progress_record"   # 2026/08/11 新增:进度时间序列,append-only、变化才写
+DETAIL_PATH = "/api/v1/app/groupbuy/detail"
+
+# 日志:按天切分文件,保留 7 天;stderr 同步 INFO 便于观察
+logger.remove()
+logger.add(os.path.join(_ROOT, "logs", "{time:YYYYMMDD}_buy_record.log"),
+           encoding="utf-8", rotation="00:00",
+           format="[{time:YYYY-MM-DD HH:mm:ss.SSS}] {level} {message}",
+           level="DEBUG", retention="7 day")
+# logger.add(sys.stderr, level="INFO",
+#            format="[{time:HH:mm:ss}] {level} {message}")
+
+# 相对时间文本:捕获数字 + 单位
+_REL_RE = re.compile(r"^(\d+)\s*(秒|分钟|小时|天)前$")
+
+# 动态判重容差(秒):相对时间越粗,反推时间戳随采集时刻漂移越大(幅度≈单位桶宽),容差须≥桶宽才能吸附同一笔的漂移副本。
+# 秒/分钟保持小容差(<同用户多笔的最小间隔~2min)以区分真实多笔;小时/天放大到略大于桶宽(3600/86400)。
+_TOL_BY_UNIT = {"秒": DEDUP_TOLERANCE_SEC, "分钟": 90, "小时": 3660, "天": 90000}
+
+# 全局请求节流:记录上次请求时刻,本轮请求前 sleep 到最小间隔
+_last_req_ts = 0.0
+
+# 订单去重内存索引:{(product_code, user_id, card_count): [已入库订单的反推时间戳...]}
+# 判重靠「反推时间戳 ± 动态容差(_TOL_BY_UNIT)」——同一笔订单变老后的漂移能吸附回锚点,不同笔可区分
+_order_idx: dict = {}
+
+# 已判定售卖结束的商品黑名单:移出监控后加入,_refresh_monitored 跳过,避免白名单模式反复重新纳入
+_ended_codes: set = set()
+
+# 上次窗口快照:{code: [(user_id, card_count), ...]}(index0 最新),供「序列对齐」求本轮顶部新增的 k 条
+_last_window: dict = {}
+
+
+# ==================== 工具函数 ====================
+def after_log(retry_state):
+    """tenacity 重试回调,记录每次尝试的结果。
+
+    Args:
+        retry_state: tenacity RetryCallState。
+    """
+    log = retry_state.args[0] if retry_state.args else 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")
+
+
+def parse_relative_ago(text: str, now_ts: int) -> int:
+    """把「X 秒/分钟/小时/天前」或「刚刚」反推为绝对秒级时间戳。
+
+    Args:
+        text (str): 相对时间原文(来自 purchaseRecords.purchasedAt)。
+        now_ts (int): 采集时刻的秒级 unix 时间戳。
+
+    Returns:
+        int: 反推的绝对购买时间戳(秒);无法识别时兜底返回 now_ts。
+    """
+    if not text:
+        return now_ts
+    t = text.strip()
+    if t in ("刚刚", "刚才", "现在"):
+        return now_ts
+    m = _REL_RE.match(t)
+    if not m:
+        return now_ts
+    n = int(m.group(1))
+    unit = m.group(2)
+    mult = {"秒": 1, "分钟": 60, "小时": 3600, "天": 86400}[unit]
+    return now_ts - n * mult
+
+
+def dedup_tolerance(text: str) -> int:
+    """按相对时间文本的单位返回该记录的判重容差(秒)。
+
+    相对时间越粗(小时/天),反推时间戳随采集时刻的漂移越大(幅度≈单位桶宽),
+    需用大容差才能把同一笔订单变老后的漂移副本吸附回锚点,否则会被误判为新订单重复入库。
+
+    Args:
+        text (str): 相对时间原文(如 "1分钟前" / "3小时前")。
+
+    Returns:
+        int: 该粒度下的判重容差秒数;无法识别时返回秒级基准 DEDUP_TOLERANCE_SEC。
+    """
+    m = _REL_RE.match((text or "").strip())
+    if not m:
+        return DEDUP_TOLERANCE_SEC
+    return _TOL_BY_UNIT.get(m.group(2), DEDUP_TOLERANCE_SEC)
+
+
+def is_sale_ended(data: dict, now_ts: int) -> tuple[bool, str]:
+    """根据详情返回判断商品售卖是否已结束(售罄或到结束时间)。
+
+    Args:
+        data (dict): 详情接口 data 层。
+        now_ts (int): 当前秒级时间戳。
+
+    Returns:
+        tuple[bool, str]: (是否已结束, 原因文本)。未结束时原因为空串。
+    """
+    stock = data.get("availableStock")
+    if stock is not None and stock <= 0:
+        return True, f"售罄(availableStock={stock})"
+    end_text = data.get("saleEndAt")
+    if end_text:
+        try:
+            end_ts = time.mktime(time.strptime(end_text, "%Y-%m-%d %H:%M:%S"))
+            if now_ts >= end_ts:
+                return True, f"已过结束时间(saleEndAt={end_text})"
+        except (ValueError, OverflowError):
+            pass
+    return False, ""
+
+
+def throttled_do_request(log, path: str, body: dict, need_auth: bool = False) -> dict | None:
+    """核心签名 POST 请求 + 全局最小间隔节流。
+
+    做请求前 sleep 到 `_last_req_ts + GLOBAL_MIN_GAP_SEC`,防止多商品并发瞬时打满。
+    详情/在售列表都走这里,共用同一根节流线。
+
+    Args:
+        log: 日志对象。
+        path (str): 接口相对路径。
+        body (dict): 请求体(同时用于签名)。
+        need_auth (bool, optional): 是否带 Bearer token。详情接口不需要(False),
+            在售列表需要(True)。Defaults to False。
+
+    Returns:
+        dict | None: 响应 JSON。
+    """
+    global _last_req_ts
+    now = time.time()
+    wait = _last_req_ts + GLOBAL_MIN_GAP_SEC - now
+    if wait > 0:
+        time.sleep(wait)
+    resp = core.do_request(log, path, body, need_auth=need_auth)
+    _last_req_ts = time.time()
+    return resp
+
+
+def compute_next_interval(span_sec: int, card_count: int | None = None) -> int:
+    """按 10 条时间跨度算下次采集间隔(span/3,硬夹进 [MIN, 上限])。
+
+    上限按是否「快车」区分:小批量商品(card_count<=FAST_LANE_MAX_COUNT,如 30 份拼团)售罄极快
+    (实测最快 10 份滑窗约 64s),上限压到 FAST_LANE_MAX_INTERVAL_SEC 密采;其余商品沿用 MAX_INTERVAL_SEC。
+    span<=0 有两种情形——刚上架还没人买(purchaseRecords 空)、或最近 10 条落在同一相对时间桶(超热)——
+    对快车都返回其小上限快探,避免旧逻辑误判为最慢 600s,导致快售拼团在两次采集间隙整场漏采。
+
+    Args:
+        span_sec (int): 10 条 purchaseRecords 最新与最老之间的秒数差。
+        card_count (int, optional): 商品总份数,用于判定是否快车。None 时按非快车处理。Defaults to None。
+
+    Returns:
+        int: 下次采集间隔(秒)。
+    """
+    is_fast = card_count is not None and card_count <= FAST_LANE_MAX_COUNT
+    cap = FAST_LANE_MAX_INTERVAL_SEC if is_fast else MAX_INTERVAL_SEC
+    if span_sec <= 0:
+        # 快车:刚上架没人买/超热同桶 → 快探到小上限;非快车:维持原「最大间隔兜底」600s
+        return cap if is_fast else MAX_INTERVAL_SEC
+    interval = span_sec // 3
+    return max(MIN_INTERVAL_SEC, min(cap, interval))
+
+
+# ==================== 订单判重 ====================
+def load_order_index(pool, code: str):
+    """从库里恢复某商品已入库订单的反推时间戳到内存索引(进程启动/新商品纳入时调)。
+
+    进程重启会丢失内存索引,靠这步从 DB 恢复,避免重启后把老订单当新订单重复入库。
+
+    Args:
+        pool: MySQL 连接池。
+        code (str): 商品编码。
+    """
+    rows = pool.select_all(
+        f"SELECT user_id, card_count, purchased_at_ts FROM {TABLE} WHERE product_code=%s", (code,)) or []
+    for uid, cc, ts in rows:
+        _order_idx.setdefault((code, str(uid), cc), []).append(int(ts))
+
+
+def is_duplicate_order(code: str, uid: str, cc: int, ts: int, tol: int) -> bool:
+    """判断一条购买记录是否为已入库订单的相对时间漂移(非新订单)。
+
+    在同 (code, uid, cc) 下的已入库时间戳里找是否有一个与 ts 相差 <= tol;
+    有则视为同一笔订单的漂移(重复),无则视为新订单。tol 由该记录粒度动态给出
+    (见 dedup_tolerance),粗粒度用大容差以吸附漂移副本。
+
+    Args:
+        code (str): 商品编码。
+        uid (str): 买家用户 ID。
+        cc (int): 购买份数。
+        ts (int): 本条记录反推的绝对购买时间戳(秒)。
+        tol (int): 本条记录的判重容差(秒),来自 dedup_tolerance。
+
+    Returns:
+        bool: True 表示是已有订单的漂移(应跳过),False 表示新订单(应入库)。
+    """
+    for ts_existing in _order_idx.get((code, str(uid), cc), []):
+        if abs(ts - ts_existing) <= tol:
+            return True
+    return False
+
+
+def align_new_orders(prev_keys: list, curr_keys: list) -> tuple[int, bool]:
+    """用倒序滑动窗口的整体位移,求本轮从顶部新增的订单数。
+
+    purchaseRecords 是倒序滑动窗口:新单只从顶部推入、旧单整体下移。故存在最小 k,使
+    curr_keys[k:] 与 prev_keys[:len(curr_keys)-k] 逐元素 (user_id, card_count) 相等
+    (旧记录整体下移 k 位);则 curr_keys[:k] 即本轮新增。取最小 k(最大重叠)作最保守估计,
+    残余错位由调用方的动态容差 + DB 唯一键二次兜底。
+
+    Args:
+        prev_keys (list): 上次窗口的 (user_id, card_count) 序列(index0 最新);无上次窗口传 None/[]。
+        curr_keys (list): 本次窗口的 (user_id, card_count) 序列(index0 最新)。
+
+    Returns:
+        tuple[int, bool]: (新增条数 k, 是否可靠对齐)。窗口整体换新/无上次窗口无法对齐时,
+            返回 (len(curr_keys), False),表示需退回容差兜底。
+    """
+    n = len(curr_keys)
+    if not prev_keys:
+        return n, False                          # 无上次窗口(重启首轮):交给兜底
+    for k in range(0, n + 1):
+        overlap = n - k
+        if overlap == 0:
+            return n, False                      # 与上次完全不重叠:整窗换新,疑似漏采
+        if overlap > len(prev_keys):
+            continue                             # 上次窗口不够长,尝试更大的 k
+        if curr_keys[k:] == prev_keys[:overlap]:
+            return k, True                       # 旧记录整体下移 k 位对齐成功,顶部 k 条为新增
+    return n, False
+
+
+# ==================== 商品源(复用 daily 免 token 采集 + 查库)====================
+def snapshot_onsale_progress(log, pool) -> int:
+    """采集当前在售商品的进度快照到 deca_onsale_product_progress_record(append-only,变化才写)。
+
+    每次 get_onsale_products 更新完 deca_onsale_product_record 后调一次:先从 progress 表拿每个
+    商品的最新 sold_count 做基线(一次 JOIN 一次性拿全),跟 deca_onsale_product_record 当前值
+    对比,**sold_count 变化 或 从未记录过的商品** 才 INSERT 一行——避免"没卖动"的团重复占位。
+
+    与 deca_onsale_product_daily_record(每日单点)互补:本表是分钟级细粒度时间序列,供后续
+    画进度曲线、算售卖速度、找热销时段等统计。
+
+    Args:
+        log: 日志对象。
+        pool: MySQL 连接池。
+
+    Returns:
+        int: 本轮写入的新快照行数。
+    """
+    # 1) 基线:progress 表每个商品的最新一条 sold_count(表可能为空 → baseline={})
+    baseline_rows = pool.select_all(
+        f"SELECT p.product_code, p.sold_count "
+        f"FROM {PROGRESS_TABLE} p "
+        f"INNER JOIN (SELECT product_code, MAX(captured_at) AS max_ts "
+        f"            FROM {PROGRESS_TABLE} GROUP BY product_code) t "
+        f"  ON p.product_code=t.product_code AND p.captured_at=t.max_ts") or []
+    baseline = {code: sold for code, sold in baseline_rows}
+
+    # 2) 本轮 onsale 表里所有在售商品的当前状态(get_onsale_products 刚 upsert 过)
+    curr = pool.select_all(
+        f"SELECT product_code, merchant_user_id, sold_count, available_stock, card_count, unit_price "
+        f"FROM {ONSALE_TABLE} WHERE is_on_sale=1") or []
+
+    # 3) Python 对比:sold_count 与基线不同 或 从未记录 → 加入待插入
+    now = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
+    rows_to_insert = []
+    for code, mid, sold, avail, card, up in curr:
+        if code is None:
+            continue
+        if code in baseline and baseline[code] == sold:
+            continue                                          # 未变化,跳过
+        pct = None
+        if card and sold is not None:
+            try:
+                pct = round(sold * 100.0 / card, 2)
+            except Exception:
+                pct = None
+        rows_to_insert.append({
+            "product_code": code, "merchant_user_id": mid,
+            "sold_count": sold, "available_stock": avail, "card_count": card,
+            "progress_pct": pct, "unit_price": up, "captured_at": now,
+        })
+    if rows_to_insert:
+        pool.insert_many(table=PROGRESS_TABLE, data_list=rows_to_insert, ignore=False)
+    log.info(f"[进度快照] 变化写入 {len(rows_to_insert)} 条(当前在售 {len(curr)} 个 / 基线 {len(baseline)} 个)")
+    return len(rows_to_insert)
+
+
+def ingest_onsale(log, pool, discover: bool):
+    """复用 daily 逻辑免 token 采集全站在售并落三张表。
+
+    每轮都拉全站在售 get_onsale_products(写 deca_onsale_product_record + 每日快照 + 下架对账);
+    紧跟一步 snapshot_onsale_progress 记录本轮进度变化(分钟级时间序列,供后续统计)。
+    discover=True 时额外做商家发现 get_shop_list 与今日新增商品补详情 fill_product_details——
+    这两步变化慢,按 SHOP_DISCOVER_SEC 降频触发。各步独立 try/except,单步失败不拖垮整轮。
+
+    Args:
+        log: 日志对象。
+        pool: MySQL 连接池。
+        discover (bool): 本轮是否附带商家发现 + 新增补详情(低频触发)。
+    """
+    try:
+        n = get_onsale_products(log, pool)
+        log.info(f"[在售落库] 全站在售写入/更新 {n} 个")
+    except Exception as e:
+        log.error(f"get_onsale_products error: {e}")
+    # 进度时间序列:紧跟 onsale 表更新之后跑(此时表里就是最新一轮的 sold_count/available_stock)
+    try:
+        snapshot_onsale_progress(log, pool)
+    except Exception as e:
+        log.error(f"snapshot_onsale_progress error: {e}")
+    if discover:
+        try:
+            n = get_shop_list(log, pool)
+            log.info(f"[商家发现] 去重商家 {n} 个")
+        except Exception as e:
+            log.error(f"get_shop_list error: {e}")
+        try:
+            n = fill_product_details(log, pool)
+            log.info(f"[补详情] 今日新增补详情 {n} 个")
+        except Exception as e:
+            log.error(f"fill_product_details error: {e}")
+
+
+def fetch_on_sale_products(log, merchant_id: str, pool) -> list:
+    """从库里查某商家当前在售商品(deca_onsale_product_record,由 ingest_onsale 落库)。
+
+    在售数据已由 ingest_onsale 复用 daily 的 home/search(免 token)落库,这里只查库拿该商家
+    is_on_sale=1 的商品,避免再打一次需 token 的 on-sale-list 接口。
+
+    Args:
+        log: 日志对象。
+        merchant_id (str): 商家用户 ID。
+        pool: MySQL 连接池。
+
+    Returns:
+        list[dict]: 每项含 code / title / sold_count / card_count / available_stock。
+    """
+    rows = pool.select_all(
+        f"SELECT product_code, title, sold_count, card_count, available_stock, merchant_user_id, merchant_name "
+        f"FROM {ONSALE_TABLE} WHERE merchant_user_id=%s AND is_on_sale=1", (merchant_id,)) or []
+    out = []
+    for pc, title, sold, cc, stock, mid, mname in rows:
+        out.append({
+            "code": pc,
+            "title": title,
+            "sold_count": sold,
+            "card_count": cc,
+            "available_stock": stock,
+            "merchant_user_id": mid,
+            "merchant_name": mname,
+        })
+    return out
+
+
+# ==================== 单商品采集 ====================
+def build_row(rec: dict, code: str, now_ts: int, now_dt: str, meta: dict) -> dict | None:
+    """把一条 purchaseRecords 项映射为 deca_buy_record 行字典。
+
+    Args:
+        rec (dict): 单条 purchaseRecords 项。
+        code (str): 商品编码。
+        now_ts (int): 本轮采集时刻的秒级时间戳。
+        now_dt (str): now_ts 对应的 datetime 文本。
+        meta (dict): 该商品冗余描述,含 title / merchant_user_id / merchant_name(来自监控队列/查库),随行入库便于查看。
+
+    Returns:
+        dict | None: 行字典;缺 userId 时返回 None。
+    """
+    uid = rec.get("userId")
+    if not uid:
+        return None
+    pat_text = rec.get("purchasedAt") or ""
+    pat_ts = parse_relative_ago(pat_text, now_ts)
+    return {
+        "product_code": code,
+        "merchant_user_id": meta.get("merchant_user_id"),
+        "merchant_name": meta.get("merchant_name"),
+        "title": meta.get("title"),
+        "user_id": str(uid),
+        "nickname": rec.get("nickname"),
+        "card_count": rec.get("cardCount"),
+        "purchased_at_text": pat_text,
+        "purchased_at_ts": pat_ts,
+        "purchased_at": time.strftime("%Y-%m-%d %H:%M:%S", time.localtime(pat_ts)),
+        "first_seen_at": now_dt,
+    }
+
+
+def poll_product(log, pool, code: str, meta: dict) -> tuple:
+    """对单个商品拉一次详情、判售卖是否结束、入库、算 span。
+
+    Args:
+        log: 日志对象。
+        pool: MySQL 连接池。
+        code (str): 商品编码。
+        meta (dict): 该商品冗余描述(title/merchant_user_id/merchant_name),随每条购买记录入库便于查看。
+
+    Returns:
+        tuple[int, int, bool, str]: (span_sec, new_count, ended, end_reason) ——
+            10 条时间跨度秒数、本轮新增入库条数、是否售卖结束、结束原因文本(未结束为空串)。
+    """
+    resp = throttled_do_request(log, DETAIL_PATH, {"code": code})
+    data = (resp or {}).get("data") or {}
+
+    now_ts = int(time.time())
+    ended, end_reason = is_sale_ended(data, now_ts)  # 结束也先把本轮记录收尾入库,再由主循环移出
+
+    recs = data.get("purchaseRecords") or []
+    if not recs:
+        return 0, 0, ended, end_reason
+
+    now_dt = time.strftime("%Y-%m-%d %H:%M:%S", time.localtime(now_ts))
+    rows = [r for r in (build_row(x, code, now_ts, now_dt, meta) for x in recs) if r]
+    if not rows:
+        return 0, 0, ended, end_reason
+
+    # 计算 span:records[0] 最新、records[-1] 最老,均是相对时间反推
+    ts_new = rows[0]["purchased_at_ts"]
+    ts_old = rows[-1]["purchased_at_ts"]
+    span_sec = max(0, ts_new - ts_old)
+
+    # 主判重:倒序滑动窗口序列对齐——只有从顶部新增的 k 条才是本轮新订单
+    curr_keys = [(r["user_id"], r["card_count"]) for r in rows]
+    prev_keys = _last_window.get(code)
+    k, reliable = align_new_orders(prev_keys, curr_keys)
+
+    if reliable:
+        # 对齐可靠:顶部 k 条直接入库,不再走时间容差(否则会误杀同用户短间隔的多笔真实订单)
+        new_rows = rows[:k]
+        for row in new_rows:
+            _order_idx.setdefault((code, row["user_id"], row["card_count"]), []).append(row["purchased_at_ts"])
+    else:
+        # 无法对齐(进程重启首轮 / 窗口整体换新):整窗退回动态容差 + DB 唯一键兜底
+        if prev_keys:
+            log.warning(f"[{code}] 窗口无重叠,疑似漏采(采集间隔内卖出≥{len(rows)}笔),建议提高采集频率")
+        new_rows = []
+        for row in rows:
+            tol = dedup_tolerance(row["purchased_at_text"])
+            if is_duplicate_order(code, row["user_id"], row["card_count"], row["purchased_at_ts"], tol):
+                continue
+            new_rows.append(row)
+            _order_idx.setdefault((code, row["user_id"], row["card_count"]), []).append(row["purchased_at_ts"])
+
+    if new_rows:
+        # DB 唯一键(含 purchased_at_ts)兜底防并发/异常重复;主判重靠窗口对齐已完成
+        pool.insert_many(table=TABLE, data_list=new_rows, ignore=True)
+
+    _last_window[code] = curr_keys              # 更新窗口快照,供下轮对齐
+    return span_sec, len(new_rows), ended, end_reason
+
+
+# ==================== 主循环 ====================
+def _refresh_monitored(log, pool, monitored: dict, merchant_id: str, now: float):
+    """刷新监控队列:新商品立即入队,下线商品移出队列。
+
+    WATCH_CODES 非空 → 只监控白名单里的 code(跳过在售发现,无下线逻辑)。
+    WATCH_CODES 为空 → 查库拿 merchant_id 的在售商品(由 ingest_onsale 落库),动态增删。
+    每个商品首次纳入监控时,从库里恢复其订单去重索引(避免重启后重复入库)。
+
+    Args:
+        log: 日志对象。
+        pool: MySQL 连接池(用于恢复订单去重索引)。
+        monitored (dict): 监控状态字典 {code: {next_run_ts, interval, title, ...}},就地更新。
+        merchant_id (str): 商家 ID。
+        now (float): 当前时间戳。
+    """
+    if WATCH_CODES:
+        # 单/多商品白名单模式:只挂想盯的 code,不动其他
+        for code in WATCH_CODES:
+            if code in _ended_codes:
+                continue                      # 售卖已结束,不再重新纳入监控
+            if code not in monitored:
+                load_order_index(pool, code)  # 从库恢复订单去重索引
+                monitored[code] = {
+                    "next_run_ts": now,
+                    "interval": MIN_INTERVAL_SEC,
+                    "title": f"(白名单 code) {code}",
+                    "sold_count": None,
+                    "card_count": None,
+                    "merchant_user_id": None,
+                    "merchant_name": None,
+                }
+                log.info(f"[+] 纳入监控 {code} | 白名单模式 | 已恢复历史订单索引")
+        log.info(f"当前监控商品数:{len(monitored)}(白名单模式,共 {len(WATCH_CODES)} 个 code)")
+        return
+
+    products = fetch_on_sale_products(log, merchant_id, pool)
+    live_codes = {p["code"] for p in products}
+    # 新增:立即到点采
+    for p in products:
+        if p["code"] in _ended_codes:
+            continue                          # 售卖已结束,跳过(正常也会自然从在售列表消失)
+        if p["code"] not in monitored:
+            load_order_index(pool, p["code"])  # 从库恢复订单去重索引
+            monitored[p["code"]] = {
+                "next_run_ts": now,
+                "interval": MIN_INTERVAL_SEC,
+                "title": p["title"],
+                "sold_count": p["sold_count"],
+                "card_count": p["card_count"],
+                "merchant_user_id": p["merchant_user_id"],
+                "merchant_name": p["merchant_name"],
+            }
+            log.info(f"[+] 纳入监控 {p['code']} | {p['title']} | 售 {p['sold_count']}/{p['card_count']}")
+    # 下线:移出
+    for code in list(monitored.keys()):
+        if code not in live_codes:
+            log.info(f"[-] 移出监控 {code} | {monitored[code].get('title')}")
+            del monitored[code]
+            _last_window.pop(code, None)         # 下线同步清窗口快照,重新上架时按重启首轮兜底
+    log.info(f"当前监控商品数:{len(monitored)}")
+
+
+@retry(stop=stop_after_attempt(100), wait=wait_fixed(3600), after=after_log)
+def main_task(log):
+    """常驻主流程:每分钟落全站在售(免 token) + 采本商家在售商品的购买记录。
+
+    Args:
+        log: 日志对象。
+
+    Raises:
+        RuntimeError: 数据库连接池异常,触发外层每小时重试。
+    """
+    mode = f"白名单模式 codes={WATCH_CODES}" if WATCH_CODES else f"商家模式 merchant={MERCHANT_ID}"
+    log.info(f"购买记录+在售常驻采集启动 | {mode} | 购买记录频率 [{MIN_INTERVAL_SEC}s, {MAX_INTERVAL_SEC}s] "
+             f"| 在售落库每 {ONSALE_INGEST_SEC}s | 全链路免 token | 详情代理={core.USE_PROXY}")
+    pool = MySQLConnectionPool(log=log)
+    if not pool.check_pool_health():
+        log.error("数据库连接池异常")
+        raise RuntimeError("db pool 异常")
+
+    monitored: dict = {}                        # {code: {next_run_ts, interval, title, sold_count, card_count}}
+    last_ingest = 0.0                           # 上次「在售落库 + 刷新监控」时刻
+    last_discover = 0.0                         # 上次「商家发现 + 补详情」时刻(按 SHOP_DISCOVER_SEC 降频)
+
+    while True:
+        now = time.time()
+
+        # 1) 每 ONSALE_INGEST_SEC:商家模式先落全站在售(免token)再查库刷新监控;白名单模式只按固定 code 刷新
+        if now - last_ingest >= ONSALE_INGEST_SEC:
+            try:
+                if not WATCH_CODES:
+                    discover = (now - last_discover >= SHOP_DISCOVER_SEC)
+                    ingest_onsale(log, pool, discover)
+                    if discover:
+                        last_discover = now
+                _refresh_monitored(log, pool, monitored, MERCHANT_ID, now)
+            except Exception as e:
+                log.error(f"在售落库/刷新监控异常: {e}")
+            last_ingest = now
+
+        # 2) 找到点的商品,取 next_run_ts 最早的一个跑
+        due_codes = [c for c, s in monitored.items() if s["next_run_ts"] <= now]
+        if not due_codes:
+            time.sleep(IDLE_SLEEP_SEC)
+            continue
+        code = min(due_codes, key=lambda c: monitored[c]["next_run_ts"])
+
+        try:
+            span_sec, n_rows, ended, end_reason = poll_product(log, pool, code, monitored[code])
+            if ended:
+                # 售卖结束:本轮收尾记录已在 poll_product 入库,这里移出监控 + 拉黑,避免死循环重采固定的最后 10 条
+                _ended_codes.add(code)
+                _last_window.pop(code, None)     # 清窗口快照,避免复活时误对齐
+                title = monitored[code]["title"]
+                monitored.pop(code, None)
+                log.info(f"[{code}] 售卖已结束({end_reason}),停止采集并移出监控 | {title}")
+                continue
+            interval = compute_next_interval(span_sec, monitored[code].get("card_count"))
+            monitored[code]["next_run_ts"] = time.time() + interval
+            monitored[code]["interval"] = interval
+            log.info(f"[{code}] 响应 {n_rows} 条 | 10条跨度 {span_sec}s | 下次 {interval}s 后 | {monitored[code]['title']}")
+        except Exception as e:
+            log.error(f"[{code}] 采集失败: {e}")
+            # 失败退避:60s 后重试;避免异常商品阻塞全局
+            monitored[code]["next_run_ts"] = time.time() + 60
+
+
+def schedule_task():
+    """脚本入口:直接进入 main_task(其 tenacity retry 保证挂了每小时重试)。"""
+    main_task(log=logger)
+
+
+if __name__ == "__main__":
+    schedule_task()

+ 671 - 0
deca_spider/buy_record_analysis/mysql_pool.py

@@ -0,0 +1,671 @@
+# -*- coding: utf-8 -*-
+# Author : Charley
+# Python : 3.12.10
+# Date   : 2025/3/25 14:14
+import re
+import pymysql
+import YamlLoader
+from loguru import logger
+from dbutils.pooled_db import PooledDB
+
+# 获取yaml配置
+yaml = YamlLoader.readYaml()
+mysqlYaml = yaml.get("mysql")
+sql_host = mysqlYaml.getValueAsString("host")
+sql_port = mysqlYaml.getValueAsInt("port")
+sql_user = mysqlYaml.getValueAsString("username")
+sql_password = mysqlYaml.getValueAsString("password")
+sql_db = mysqlYaml.getValueAsString("db")
+
+
+class MySQLConnectionPool:
+    """
+    MySQL连接池
+    """
+
+    def __init__(self, mincached=1, maxcached=2, maxconnections=3, log=None):
+        """
+        初始化连接池
+        :param mincached: 初始化时,链接池中至少创建的链接,0表示不创建
+        :param maxcached: 池中空闲连接的最大数目(0 或 None 表示池大小不受限制)
+        :param maxconnections: 允许的最大连接数(0 或 None 表示任意数量的连接)
+        :param log: 自定义日志记录器
+        """
+        # 使用 loguru 的 logger,如果传入了其他 logger,则使用传入的 logger
+        self.log = log or logger
+        self.pool = PooledDB(
+            creator=pymysql,
+            mincached=mincached,
+            maxcached=maxcached,
+            maxconnections=maxconnections,
+            blocking=True,  # 连接池中如果没有可用连接后,是否阻塞等待。True,等待;False,不等待然后报错
+            host=sql_host,
+            port=sql_port,
+            user=sql_user,
+            password=sql_password,
+            database=sql_db,
+            ping=2,  # 每次执行前检查连接有效性,防止使用已断开的连接
+            connect_timeout=5,  # 连接超时时间(秒)
+            # read_timeout=30,  # 读取超时时间(秒)
+            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(带断连重试)
+        :param query: SQL语句
+        :param args: SQL参数
+        :param commit: 是否提交事务
+        :return: 查询结果
+        """
+        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):
+        """
+        执行查询,返回单个结果
+        :param query: 查询语句
+        :param args: 查询参数
+        :return: 查询结果
+        """
+        cursor = self._execute(query, args)
+        return cursor.fetchone()
+
+    def select_all(self, query, args=None):
+        """
+        执行查询,返回所有结果
+        :param query: 查询语句
+        :param args: 查询参数
+        :return: 查询结果
+        """
+        cursor = self._execute(query, args)
+        return cursor.fetchall()
+
+    def insert_one(self, query, args):
+        """
+        执行单条插入语句
+        :param query: 插入语句
+        :param args: 插入参数
+        """
+        self.log.info('>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>data insert_one 入库中>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>')
+        cursor = self._execute(query, args, commit=True)
+        return cursor.lastrowid  # 返回插入的ID
+
+    def insert_all(self, query, args_list):
+        """
+        执行批量插入语句,如果失败则逐条插入
+        :param query: 插入语句
+        :param args_list: 插入参数列表
+        """
+        conn = None
+        cursor = None
+        try:
+            conn = self.pool.connection()
+            cursor = conn.cursor()
+            cursor.executemany(query, args_list)
+            conn.commit()
+            self.log.debug(f"sql insert_all, SQL: {query[:100]}..., Rows: {cursor.rowcount}")
+            self.log.info('>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>data insert_all 入库中>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>')
+        except pymysql.err.IntegrityError as e:
+            if "Duplicate entry" in str(e):
+                conn.rollback()
+                self.log.warning(f"批量插入遇到重复,开始逐条插入。错误: {e}")
+                rowcount = 0
+                for args in args_list:
+                    try:
+                        self.insert_one(query, args)
+                        rowcount += 1
+                    except pymysql.err.IntegrityError as e2:
+                        if "Duplicate entry" in str(e2):
+                            self.log.debug(f"跳过重复条目: {e2}")
+                        else:
+                            self.log.error(f"插入失败: {e2}")
+                    except Exception as e2:
+                        self.log.error(f"插入失败: {e2}")
+                self.log.info(f"逐条插入完成: {rowcount}/{len(args_list)}条")
+            else:
+                conn.rollback()
+                self.log.exception(f"数据库完整性错误: {e}")
+                raise e
+        except Exception as e:
+            conn.rollback()
+            self.log.exception(f"批量插入失败: {e}")
+            raise e
+        finally:
+            if cursor:
+                cursor.close()
+            if conn:
+                conn.close()
+
+    def insert_one_or_dict(self, table=None, data=None, query=None, args=None, commit=True, ignore=False):
+        """
+        单条插入(支持字典或原始SQL)
+        :param table: 表名(字典插入时必需)
+        :param data: 字典数据 {列名: 值}
+        :param query: 直接SQL语句(与data二选一)
+        :param args: SQL参数(query使用时必需)
+        :param commit: 是否自动提交
+        :param ignore: 是否使用ignore
+        :return: 最后插入ID
+        """
+        if data is not None:
+            if not isinstance(data, dict):
+                raise ValueError("Data must be a dictionary")
+
+            keys = ', '.join([self._safe_identifier(k) for k in data.keys()])
+            values = ', '.join(['%s'] * len(data))
+
+            # 构建 INSERT IGNORE 语句
+            ignore_clause = "IGNORE" if ignore else ""
+            query = f"INSERT {ignore_clause} INTO {self._safe_identifier(table)} ({keys}) VALUES ({values})"
+            args = tuple(data.values())
+        elif query is None:
+            raise ValueError("Either data or query must be provided")
+
+        try:
+            cursor = self._execute(query, args, commit)
+            self.log.info(f"sql insert_one_or_dict, Table: {table}, Rows: {cursor.rowcount}")
+            self.log.info('>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>data insert_one_or_dict 入库中>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>')
+            return cursor.lastrowid
+        except pymysql.err.IntegrityError as e:
+            if "Duplicate entry" in str(e):
+                # 重复条目用 warning 简短输出,不打印堆栈
+                self.log.warning(f"插入跳过-重复条目 Table: {table}, {e.args[1] if len(e.args) > 1 else e}")
+                return -1  # 返回 -1 表示重复条目被跳过
+            else:
+                self.log.error(f"数据库完整性错误 Table: {table}, Error: {e}")
+                raise
+        except Exception as 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,
+                    ignore=False):
+        """
+        批量插入(支持字典列表或原始SQL)
+        :param table: 表名(字典插入时必需)
+        :param data_list: 字典列表 [{列名: 值}]
+        :param query: 直接SQL语句(与data_list二选一)
+        :param args_list: SQL参数列表(query使用时必需)
+        :param batch_size: 分批大小
+        :param commit: 是否自动提交
+        :param ignore: 是否使用ignore
+        :return: 影响行数
+        """
+        if data_list is not None:
+            if not data_list or not isinstance(data_list[0], dict):
+                raise ValueError("Data_list must be a non-empty list of dictionaries")
+
+            keys = ', '.join([self._safe_identifier(k) for k in data_list[0].keys()])
+            values = ', '.join(['%s'] * len(data_list[0]))
+
+            # 构建 INSERT IGNORE 语句
+            ignore_clause = "IGNORE" if ignore else ""
+            query = f"INSERT {ignore_clause} INTO {self._safe_identifier(table)} ({keys}) VALUES ({values})"
+            args_list = [tuple(d.values()) for d in data_list]
+        elif query is None:
+            raise ValueError("Either data_list or query must be provided")
+
+        total = 0
+        for i in range(0, len(args_list), batch_size):
+            batch = args_list[i:i + batch_size]
+            try:
+                with self.pool.connection() as conn:
+                    with conn.cursor() as cursor:
+                        cursor.executemany(query, batch)
+                        if commit:
+                            conn.commit()
+                        total += cursor.rowcount
+            except pymysql.err.IntegrityError as e:
+                # 处理唯一索引冲突
+                if "Duplicate entry" in str(e):
+                    if ignore:
+                        # 如果使用了 INSERT IGNORE,理论上不会进这里,但以防万一
+                        self.log.warning(f"批量插入遇到重复条目(ignore模式): {e}")
+                    else:
+                        # 没有使用 IGNORE,降级为逐条插入
+                        self.log.warning(f"批量插入遇到重复条目,开始逐条插入。错误: {e}")
+                        if commit:
+                            conn.rollback()
+                        
+                        rowcount = 0
+                        for j, args in enumerate(batch):
+                            try:
+                                if data_list:
+                                    # 字典模式
+                                    self.insert_one_or_dict(
+                                        table=table,
+                                        data=dict(zip(data_list[0].keys(), args)),
+                                        commit=commit,
+                                        ignore=False  # 单条插入时手动捕获重复
+                                    )
+                                else:
+                                    # 原始SQL模式
+                                    self.insert_one(query, args)
+                                rowcount += 1
+                            except pymysql.err.IntegrityError as e2:
+                                if "Duplicate entry" in str(e2):
+                                    self.log.debug(f"跳过重复条目[{i+j+1}]: {e2}")
+                                else:
+                                    self.log.error(f"插入失败[{i+j+1}]: {e2}")
+                            except Exception as e2:
+                                self.log.error(f"插入失败[{i+j+1}]: {e2}")
+                        total += rowcount
+                        self.log.info(f"批次逐条插入完成: 成功{rowcount}/{len(batch)}条")
+                else:
+                    # 其他完整性错误
+                    self.log.exception(f"数据库完整性错误: {e}")
+                    if commit:
+                        conn.rollback()
+                    raise e
+            except Exception as e:
+                # 其他数据库错误
+                self.log.exception(f"批量插入失败: {e}")
+                if commit:
+                    conn.rollback()
+                raise e
+        if table:
+            self.log.info(f"sql insert_many, Table: {table}, Total Rows: {total}")
+        else:
+            self.log.info(f"sql insert_many, Query: {query}, Total Rows: {total}")
+        return total
+
+    def insert_many_two(self, table=None, data_list=None, query=None, args_list=None, batch_size=1000, commit=True,
+                        ignore=False):
+        """
+        批量插入(支持字典列表或原始SQL) - 备用方法
+        :param table: 表名(字典插入时必需)
+        :param data_list: 字典列表 [{列名: 值}]
+        :param query: 直接SQL语句(与data_list二选一)
+        :param args_list: SQL参数列表(query使用时必需)
+        :param batch_size: 分批大小
+        :param commit: 是否自动提交
+        :param ignore: 是否使用INSERT IGNORE
+        :return: 影响行数
+        """
+        if data_list is not None:
+            if not data_list or not isinstance(data_list[0], dict):
+                raise ValueError("Data_list must be a non-empty list of dictionaries")
+            keys = ', '.join([self._safe_identifier(k) for k in data_list[0].keys()])
+            values = ', '.join(['%s'] * len(data_list[0]))
+            ignore_clause = "IGNORE" if ignore else ""
+            query = f"INSERT {ignore_clause} INTO {self._safe_identifier(table)} ({keys}) VALUES ({values})"
+            args_list = [tuple(d.values()) for d in data_list]
+        elif query is None:
+            raise ValueError("Either data_list or query must be provided")
+    
+        total = 0
+        for i in range(0, len(args_list), batch_size):
+            batch = args_list[i:i + batch_size]
+            try:
+                with self.pool.connection() as conn:
+                    with conn.cursor() as cursor:
+                        cursor.executemany(query, batch)
+                        if commit:
+                            conn.commit()
+                        total += cursor.rowcount
+            except pymysql.err.IntegrityError as e:
+                if "Duplicate entry" in str(e) and not ignore:
+                    self.log.warning(f"批量插入遇到重复,降级为逐条插入: {e}")
+                    if commit:
+                        conn.rollback()
+                    rowcount = 0
+                    for args in batch:
+                        try:
+                            self.insert_one(query, args)
+                            rowcount += 1
+                        except pymysql.err.IntegrityError as e2:
+                            if "Duplicate entry" in str(e2):
+                                self.log.debug(f"跳过重复条目: {e2}")
+                            else:
+                                self.log.error(f"插入失败: {e2}")
+                        except Exception as e2:
+                            self.log.error(f"插入失败: {e2}")
+                    total += rowcount
+                else:
+                    self.log.exception(f"数据库完整性错误: {e}")
+                    if commit:
+                        conn.rollback()
+                    raise e
+            except Exception as e:
+                self.log.exception(f"批量插入失败: {e}")
+                if commit:
+                    conn.rollback()
+                raise e
+        self.log.info(f"sql insert_many_two, Table: {table}, Total Rows: {total}")
+        return total
+
+    def insert_too_many(self, query, args_list, batch_size=1000):
+        """
+        执行批量插入语句,分片提交, 单次插入大于十万+时可用, 如果失败则降级为逐条插入
+        :param query: 插入语句
+        :param args_list: 插入参数列表
+        :param batch_size: 每次插入的条数
+        """
+        self.log.info(f"sql insert_too_many, Query: {query}, Total Rows: {len(args_list)}")
+        for i in range(0, len(args_list), batch_size):
+            batch = args_list[i:i + batch_size]
+            try:
+                with self.pool.connection() as conn:
+                    with conn.cursor() as cursor:
+                        cursor.executemany(query, batch)
+                        conn.commit()
+                        self.log.debug(f"insert_too_many -> Total Rows: {len(batch)}")
+            except Exception as e:
+                self.log.error(f"insert_too_many error. Trying single insert. Error: {e}")
+                # 当前批次降级为单条插入
+                for args in batch:
+                    self.insert_one(query, args)
+
+    def update_one(self, query, args):
+        """
+        执行单条更新语句
+        :param query: 更新语句
+        :param args: 更新参数
+        """
+        self.log.info('>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>data update_one 更新中>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>')
+        return self._execute(query, args, commit=True)
+
+    def update_all(self, query, args_list):
+        """
+        执行批量更新语句,如果失败则逐条更新
+        :param query: 更新语句
+        :param args_list: 更新参数列表
+        """
+        conn = None
+        cursor = None
+        try:
+            conn = self.pool.connection()
+            cursor = conn.cursor()
+            cursor.executemany(query, args_list)
+            conn.commit()
+            self.log.debug(f"sql update_all, SQL: {query}, Rows: {len(args_list)}")
+            self.log.info('>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>data update_all 更新中>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>')
+        except Exception as e:
+            conn.rollback()
+            self.log.error(f"Error executing query: {e}")
+            # 如果批量更新失败,则逐条更新
+            rowcount = 0
+            for args in args_list:
+                self.update_one(query, args)
+                rowcount += 1
+            self.log.debug(f'Batch update failed. Updated {rowcount} rows individually.')
+        finally:
+            if cursor:
+                cursor.close()
+            if conn:
+                conn.close()
+
+    def update_one_or_dict(self, table=None, data=None, condition=None, query=None, args=None, commit=True):
+        """
+        单条更新(支持字典或原始SQL)
+        :param table: 表名(字典模式必需)
+        :param data: 字典数据 {列名: 值}(与 query 二选一)
+        :param condition: 更新条件,支持以下格式:
+            - 字典: {"id": 1} → "WHERE id = %s"
+            - 字符串: "id = 1" → "WHERE id = 1"(需自行确保安全)
+            - 元组: ("id = %s", [1]) → "WHERE id = %s"(参数化查询)
+        :param query: 直接SQL语句(与 data 二选一)
+        :param args: SQL参数(query 模式下必需)
+        :param commit: 是否自动提交
+        :return: 影响行数
+        :raises: ValueError 参数校验失败时抛出
+        """
+        # 参数校验
+        if data is not None:
+            if not isinstance(data, dict):
+                raise ValueError("Data must be a dictionary")
+            if table is None:
+                raise ValueError("Table name is required for dictionary update")
+            if condition is None:
+                raise ValueError("Condition is required for dictionary update")
+
+            # 构建 SET 子句
+            set_clause = ", ".join([f"{self._safe_identifier(k)} = %s" for k in data.keys()])
+            set_values = list(data.values())
+
+            # 解析条件
+            condition_clause, condition_args = self._parse_condition(condition)
+            query = f"UPDATE {self._safe_identifier(table)} SET {set_clause} WHERE {condition_clause}"
+            args = set_values + condition_args
+
+        elif query is None:
+            raise ValueError("Either data or query must be provided")
+
+        # 执行更新
+        cursor = self._execute(query, args, commit)
+        # self.log.debug(
+        #     f"Updated table={table}, rows={cursor.rowcount}, query={query[:100]}...",
+        #     extra={"table": table, "rows": cursor.rowcount}
+        # )
+        return cursor.rowcount
+
+    def _parse_condition(self, condition):
+        """
+        解析条件为 (clause, args) 格式
+        :param condition: 字典/字符串/元组
+        :return: (str, list) SQL 子句和参数列表
+        """
+        if isinstance(condition, dict):
+            clause = " AND ".join([f"{self._safe_identifier(k)} = %s" for k in condition.keys()])
+            args = list(condition.values())
+        elif isinstance(condition, str):
+            clause = condition  # 注意:需调用方确保安全
+            args = []
+        elif isinstance(condition, (tuple, list)) and len(condition) == 2:
+            clause, args = condition[0], condition[1]
+            if not isinstance(args, (list, tuple)):
+                args = [args]
+        else:
+            raise ValueError("Condition must be dict/str/(clause, args)")
+        return clause, args
+
+    def update_many(self, table=None, data_list=None, condition_list=None, query=None, args_list=None, batch_size=500,
+                    commit=True):
+        """
+        批量更新(支持字典列表或原始SQL)
+        :param table: 表名(字典插入时必需)
+        :param data_list: 字典列表 [{列名: 值}]
+        :param condition_list: 条件列表(必须为字典,与data_list等长)
+        :param query: 直接SQL语句(与data_list二选一)
+        :param args_list: SQL参数列表(query使用时必需)
+        :param batch_size: 分批大小
+        :param commit: 是否自动提交
+        :return: 影响行数
+        """
+        if data_list is not None:
+            if not data_list or not isinstance(data_list[0], dict):
+                raise ValueError("Data_list must be a non-empty list of dictionaries")
+            if condition_list is None or len(data_list) != len(condition_list):
+                raise ValueError("Condition_list must be provided and match the length of data_list")
+            if not all(isinstance(cond, dict) for cond in condition_list):
+                raise ValueError("All elements in condition_list must be dictionaries")
+
+            # 获取第一个数据项和条件项的键
+            first_data_keys = set(data_list[0].keys())
+            first_cond_keys = set(condition_list[0].keys())
+
+            # 构造基础SQL
+            set_clause = ', '.join([self._safe_identifier(k) + ' = %s' for k in data_list[0].keys()])
+            condition_clause = ' AND '.join([self._safe_identifier(k) + ' = %s' for k in condition_list[0].keys()])
+            base_query = f"UPDATE {self._safe_identifier(table)} SET {set_clause} WHERE {condition_clause}"
+            total = 0
+
+            # 分批次处理
+            for i in range(0, len(data_list), batch_size):
+                batch_data = data_list[i:i + batch_size]
+                batch_conds = condition_list[i:i + batch_size]
+                batch_args = []
+
+                # 检查当前批次的结构是否一致
+                can_batch = True
+                for data, cond in zip(batch_data, batch_conds):
+                    data_keys = set(data.keys())
+                    cond_keys = set(cond.keys())
+                    if data_keys != first_data_keys or cond_keys != first_cond_keys:
+                        can_batch = False
+                        break
+                    batch_args.append(tuple(data.values()) + tuple(cond.values()))
+
+                if not can_batch:
+                    # 结构不一致,转为单条更新
+                    for data, cond in zip(batch_data, batch_conds):
+                        self.update_one_or_dict(table=table, data=data, condition=cond, commit=commit)
+                        total += 1
+                    continue
+
+                # 执行批量更新
+                try:
+                    with self.pool.connection() as conn:
+                        with conn.cursor() as cursor:
+                            cursor.executemany(base_query, batch_args)
+                            if commit:
+                                conn.commit()
+                            total += cursor.rowcount
+                            self.log.debug(f"Batch update succeeded. Rows: {cursor.rowcount}")
+                except Exception as e:
+                    if commit:
+                        conn.rollback()
+                    self.log.error(f"Batch update failed: {e}")
+                    # 降级为单条更新
+                    for args, data, cond in zip(batch_args, batch_data, batch_conds):
+                        try:
+                            self._execute(base_query, args, commit=commit)
+                            total += 1
+                        except Exception as e2:
+                            self.log.error(f"Single update failed: {e2}, Data: {data}, Condition: {cond}")
+            self.log.info(f"Total updated rows: {total}")
+            return total
+        elif query is not None:
+            # 处理原始SQL和参数列表
+            if args_list is None:
+                raise ValueError("args_list must be provided when using query")
+
+            total = 0
+            for i in range(0, len(args_list), batch_size):
+                batch_args = args_list[i:i + batch_size]
+                try:
+                    with self.pool.connection() as conn:
+                        with conn.cursor() as cursor:
+                            cursor.executemany(query, batch_args)
+                            if commit:
+                                conn.commit()
+                            total += cursor.rowcount
+                            self.log.debug(f"Batch update succeeded. Rows: {cursor.rowcount}")
+                except Exception as e:
+                    if commit:
+                        conn.rollback()
+                    self.log.error(f"Batch update failed: {e}")
+                    # 降级为单条更新
+                    for args in batch_args:
+                        try:
+                            self._execute(query, args, commit=commit)
+                            total += 1
+                        except Exception as e2:
+                            self.log.error(f"Single update failed: {e2}, Args: {args}")
+            self.log.info(f"Total updated rows: {total}")
+            return total
+        else:
+            raise ValueError("Either data_list or query must be provided")
+
+    def check_pool_health(self):
+        """
+        检查连接池中有效连接数
+
+        # 使用示例
+        # 配置 MySQL 连接池
+        sql_pool = MySQLConnectionPool(log=log)
+        if not sql_pool.check_pool_health():
+            log.error("数据库连接池异常")
+            raise RuntimeError("数据库连接池异常")
+        """
+        try:
+            with self.pool.connection() as conn:
+                conn.ping(reconnect=True)
+                return True
+        except Exception as e:
+            self.log.error(f"Connection pool health check failed: {e}")
+            return False
+
+    def close(self):
+        """
+        关闭连接池,释放所有连接
+        """
+        try:
+            if hasattr(self, 'pool') and self.pool:
+                self.pool.close()
+                self.log.info("数据库连接池已关闭")
+        except Exception as e:
+            self.log.error(f"关闭连接池失败: {e}")
+
+    @staticmethod
+    def _safe_identifier(name):
+        """SQL标识符安全校验"""
+        if not re.match(r'^[a-zA-Z_][a-zA-Z0-9_]*$', name):
+            raise ValueError(f"Invalid SQL identifier: {name}")
+        return name
+
+
+if __name__ == '__main__':
+    sql_pool = MySQLConnectionPool()
+    data_dic = {'card_type_id': 111, 'card_type_name': '补充包 继承的意志【OPC-13】', 'card_type_position': 964,
+                'card_id': 5284, 'card_name': '蒙奇·D·路飞', 'card_number': 'OP13-001', 'card_rarity': 'L',
+                'card_img': 'https://source.windoent.com/OnePiecePc/Picture/1757929283612OP13-001.png',
+                'card_life': '4', 'card_attribute': '打', 'card_power': '5000', 'card_attack': '-',
+                'card_color': '红/绿', 'subscript': 4, 'card_features': '超新星/草帽一伙',
+                'card_text_desc': '【咚!!×1】【对方的攻击时】我方处于活跃状态的咚!!不多于5张的场合,可以将我方任意张数的咚!!转为休息状态。每有1张转为休息状态的咚!!,本次战斗中,此领袖或我方最多1张拥有《草帽一伙》特征的角色力量+2000。',
+                'card_offer_type': '补充包 继承的意志【OPC-13】', 'crawler_language': '简中'}
+    sql_pool.insert_one_or_dict(table="one_piece_record", data=data_dic)