ソースを参照

feat(kaji_spider): 添加卡集每日采集爬虫及签名模块

- 新增 kj_auth.py,实现卡集接口 Open-Auth-Sig 请求签名功能
- 新增 kj_daily_spider.py,完成商户发现、商品采集、玩家采集及视频补采逻辑
- 实现自动生成请求头 Authorization,支持无登录 token 长期无人值守
- 支持抓取在售商户列表,更新商户存活状态
- 实现商户已售商品列表增量采集,并补充商品详情及拆卡回放视频地址
- 支持玩家名单分页采集和匿名用户处理
- 新增视频补采功能,针对拼团完成但无视频地址的商品重新拉取视频
- 配置 mysql 连接信息,支持环境变量覆写
- 日志记录规范化,支持重试机制及代理获取功能
charley 1 ヶ月 前
コミット
554a2a7f2d

+ 151 - 0
kaji_spider/README.md

@@ -0,0 +1,151 @@
+# 卡集 App 逆向分析文档(抓包 + Open-Auth-Sig 请求签名 + 爬虫)
+
+## 1. 目标信息
+
+- **App / 包名**:卡集,`card.kaji.com`
+- **版本**:versionName `2.5.39`,versionCode `10003`(huawei 渠道)
+- **APK**:`卡集_2.5.39.apk`(约 86MB)
+- **业务域名**:`server.ssl1.kaji6.com`(主列表 / videoPlay)、`page.ssl1.kaji6.com`(商户列表 / 详情 / 玩家)
+- **技术栈**:**uni-app**(DCloud HBuilderX 打包,Vue2 + uni-v3 / Weex 渲染引擎)
+  - 业务逻辑是 JavaScript:`assets/apps/__UNI__EFAB5FD/www/app-service.js`(3.6MB)、`app-view.js`(4.9MB)
+  - Application 类:`io.dcloud.application.DCloudApplication`;入口 `io.dcloud.PandoraEntry`
+  - 关键 so:`libweexcore.so` / `libweexjss.so` / `libuts-runtime.so`
+  - **排除项**:无 `libflutter.so` / `flutter_assets`(非 Flutter);无 `index.android.bundle` / `libhermes`(非 React Native)
+- **是否加固**:**未整包加固**。四个 dex 可正常反编译,类名语义完整(`io.dcloud.*`、`com.taobao.weex.*`),无壳 Stub。
+  - `assets/39285EFA.dex` + `lib39285EFA.so` 是 DCloud uni-app 自带的分包/加密机制,不是第三方壳。
+  - 存在**局部防护 SDK**:`libzxprotect.so`(局部代码保护)、`libaliyunaf.so`(阿里聚安全)、`libsecuritydevice.so`(阿里 dtf 设备风控,约 3.9MB,疑似带反调试/反 Frida,静态未坐实)。
+- **反爬 / 校验措施**:
+  - **没有应用级硬 SSL Pinning**(JS 层、DCloud 框架层、native 层三层排查,`pin-sha256` / `CertificatePinner` 配置值均为 0)。握手拦截来自**标准证书链校验**(`targetSdk 30` + 无 `networkSecurityConfig`)+ **DCloud 的 HostnameVerifier 主机名校验**。
+  - 接口鉴权:请求头 `Authorization: Open-Auth-Sig`(**每请求签名**,非登录 token),详见第 4 节。
+  - 设备风控 SDK(`libsecuritydevice.so`)负责环境检测,与 HTTP 校验无关。
+
+---
+
+## 2. 抓包突破(SSL 握手失败)
+
+### 抓包工具与环境
+- 抓包工具:**Reqable**(PC + 手机代理)
+- 手机环境:已 **root**,**Magisk** 管理,**LSPosed** 框架
+- 证书处理:Magisk 模块 `Reqable Certificate Installer`(装 CA 到系统库)+ `Move Certificates`(用户库证书迁移到系统库)
+- unpinning 历程:LSPosed 模块 **JustTrustMe(失败)→ JustTrustMe++/算法助手(API 能抓、WebView 图片不显示)→ JustTrustMePro(完全成功,含图片)**
+
+### 现象与根因
+- 一开抓包,卡集立即报 **「客户端 SSL 握手失败」**,App 网络异常;系统级证书已装、JustTrustMe 已勾选重启仍失败。
+- 反编译确认:业务请求走 **DCloud shade(重命名)版 OkHttp `dc.squareup.okhttp3.*`**(第三方 SDK 才走标准 `okhttp3.*`),证书装配点是 `io.dcloud.common.adapter.util.DCloudTrustManager`。
+- **JustTrustMe 失效根因**:原版 hook 名单写死标准类名(`okhttp3.CertificatePinner.check`、`javax.net.ssl.*`),而卡集类前缀是 `dc.squareup.okhttp3.*`,名字对不上 → hook 没挂上。
+- **JustTrustMe++ 能抓 API 但图片不显示**:它从 `SSLContext.init` 底层兜底放行了走网络栈的 API,但**没 hook WebView 的证书校验**(`onReceivedSslError`),而 uni-app 的 view 层是 WebView、图片由 WebView 拉取。
+- **JustTrustMePro 完全成功**:它额外 hook 了 WebView 等安全连接方法,WebView 拉图片的握手也放行,图片正常显示。
+
+> 经验:unpinning 能抓到什么流量,取决于它 hook 的覆盖面;「API 能抓、图片/WebView 内容不显示」通常就是 WebView 那条 SSL 校验路径没被 hook。
+
+---
+
+## 3. Open-Auth-Sig 请求签名逆向
+
+### 目标
+抓包发现部分接口带 `Authorization: Open-Auth-Sig Swap="..",Timestamp=..,Nonce=..,Signature=".."`,且每请求都变。要让爬虫长期运行,必须还原其生成算法(而非抠固定值——`Timestamp` 有时效)。
+
+### 定位与还原
+- 技术栈是 uni-app,业务逻辑在明文 `app-service.js`,直接搜 `Open-Auth-Sig` / `Signature` / `Swap` 定位。
+- 签名入口:`app-service.js` webpack 模块 `"665c"` 的 **`GetCrypto(path)`**(偏移约 1511900),入参为接口相对路径。
+- 核心常量:256 项置换表(偏移约 1513290 的数组 `i`,0..255 全排列);MD5 库 `n("e133")`;`requestVersion = "/api/v4/"`。
+- videoPlay 的 body 签名在偏移 1203166:`sign = md5(ts + "_" + goodCode + "_" + playCode + "_videoPlayKsj")`。
+- 注入方式:**不是全局拦截器**,而是 http 封装里的 `PostWithCrypto` / `GetWithCrypto` 显式调用——所以只有部分接口带签名(解释了抓包「有的带有的不带」)。
+
+### 验证
+用三条真实抓包样本正推 `Signature`,**逐字节匹配**;videoPlay 的 `sign` 用 `goodCode=MC5388325` 正推命中 `aba9d5be1f290cec5029e12c3d747d9f`。算法确认无误,已落地为 `kj_auth.py`。
+
+---
+
+## 4. 算法要点
+
+### 4.1 Open-Auth-Sig 请求头(`kj_auth.gen_authorization(path)`)
+
+输入:接口相对路径 `path`(不含域名 / `/api/v4/` 前缀 / query,如 `goodlist/forsale/main`)。
+
+1. `noce`(Nonce)= 随机整数 `[1, 500]`
+2. `swap` = 从 0..255 取 **8 个不重复字节** → 转 16 位小写 hex(**每请求随机**)
+3. `ts`(Timestamp)= 当前 Unix 秒
+4. `full_path = ("/api/v4/" + path)`,去掉 `/dataApi` 子串,去掉 `?query`
+5. 拼接串 `l = f"{swap}_{ts}_{noce}_{full_path}"`(顺序固定:swap _ ts _ noce _ path)
+6. `f = MD5(l)` → 32 位小写 hex
+7. **Signature 编码**:
+   - `r = f.encode('ascii').hex()`(把 32 字符 md5-hex 当 ASCII 再转 hex → 64 字符)
+   - 复制 256 项置换表 `a = list(I_TABLE)`
+   - 按 swap 的 8 个字节做 4 对两两交换:`for s in range(0,8,2): a[c[s]], a[c[s+1]] = a[c[s+1]], a[c[s]]`
+   - 对 `r` 每个字符 `ch`:取 `a[ord(ch)]` 转 2 位大写 hex,拼接 → 128 位 Signature
+8. 组装:`Open-Auth-Sig Swap="{swap}",Timestamp={ts},Nonce={noce},Signature="{sig}"`
+
+> **置换表 `I_TABLE`**:256 项(0..255 全排列),来自 `app-service.js` 偏移约 1513290,完整值见 `kj_auth.py`。
+
+### 4.2 videoPlay body 签名(`kj_auth.video_sign(ts, good_code, play_code)`)
+
+```
+sign = MD5(f"{ts}_{good_code}_{play_code}_videoPlayKsj")   # 固定盐 videoPlayKsj
+```
+POST `good/videoPlay/{goodCode}`,body = `{"playCode":.., "sign":.., "ts":..}`,**不需要 Authorization**,响应的 `media_url` 即回放视频地址。
+
+### 4.3 playCode 来源
+
+`playCode` 在 **`detail.good.broadcast.playCode`**:拼团完成 / 回放就绪时非空,否则为空字符串。
+完整取视频链路:`detail → broadcast.playCode → videoPlay → media_url`。
+
+### 4.4 为什么 Signature 每次都变、字节值集中
+
+- 编码的下标是 md5-hex 字符的 ASCII 码(仅 `0-9`、`a-f` 共十来种取值),查表后输出字节自然集中在约 10 种值,看起来像「伪 hex」。
+- `swap` 每请求随机 → 置换表被洗牌 → **同一 MD5 每次输出的 Signature 都不同**,这是它的抗重放 / 防特征提取设计。
+
+### 4.5 关键结论
+
+签名 = **标准 MD5 + 256 项置换表 + swap 洗牌**,拼接串只有 `swap/ts/noce/path`,**无任何密钥 / appSecret / 登录 token**。因此爬虫**无需登录、无需维护 token、永不过期**,天然适合长期无人值守运行(比需维护 token 的项目更省心)。
+
+---
+
+## 5. 接口清单与爬虫实现
+
+### 接口清单
+
+| 接口(`/api/v4/` 之后) | 方法 | 域名 | 签名 | 作用 |
+|---|---|---|---|---|
+| `goodlist/forsale/main` | GET | server | Open-Auth-Sig | 全局在售商品主列表(发现商户) |
+| `merchant/1/goodlist/{商户码}?tp=2` | GET | page | 免签名 | 某商户已售/已完成商品列表 |
+| `good/{code}/1/detail` | GET | page | Open-Auth-Sig | 商品详情(时间/规格/商户 publisher/broadcast) |
+| `good/{code}/result/merge` | GET | page | Open-Auth-Sig | 玩家/中奖名单(仅拼团完成的商品有) |
+| `good/videoPlay/{code}` | POST | server | body sign | 拆卡回放视频(返回 media_url) |
+
+### 代码文件(两个脚本 + 一个签名模块)
+
+- **`kj_auth.py`**:签名模块。`gen_authorization(path)` 生成请求头、`video_sign(...)` 生成 videoPlay 签名、`gen_signature/_md5_hex` 为底层,含 256 项置换表。三处均以真实样本验证通过。
+- **`kj_daily_spider.py`**:**日常增量脚本**(`schedule` 每天 00:01 定时)+ 全部公共函数(请求/签名/翻页/解析/`get_video`/`refill_videos`)。三级管道 = **在售主列表发现商户 → 遍历各商户已售列表(+detail 补充) → 玩家名单**,`player_state` 断点标记(0 未抓 / 1 已抓 / 2 无数据 / 3 异常),loguru 日志 + tenacity 重试。
+  - **增量早停**:`get_sold_list(incremental=True)` 采某商户前,先用 `get_shop_stop_pid` 取该商户 `start_at` 最新那条的 `pid`(`SELECT pid ... WHERE shop_id=%s ORDER BY start_at DESC LIMIT 1`,依赖 `(shop_id, start_at)` 复合索引,单值查询);`merchant?tp=2` 按 start_at 倒序(新在前),翻页用 `goodCode` 匹配到这个 `pid` 即 `break`——它之后都是已采过的,不必再翻、也不发 detail。首次采(无记录)则全量。
+  - **商户注销跳过**:`merchant?tp=2` 返回 `code=1`/`msg='无效商家'` → 标记 `is_deleted=1`;`kj_main` 查商户带 `WHERE is_deleted=0` 自动跳过。商户重新在售时 `save_shops` 的 `ON DUPLICATE KEY UPDATE ... is_deleted=0` 会把它「复活」。
+  - **视频补采**:`refill_videos` 针对拼团完成(`player_state=1`)但 `video_url` 仍空的商品,重取 `detail.broadcast.playCode` → videoPlay,覆盖首次采集时回放未就绪(正在/即将拆卡,playCode 为空)的情况;下一轮自愈直到拿到地址。
+  - 开关:`FETCH_DETAIL`(是否每商品补 detail)、`USE_PROXY`(遇 IP 风控再开)。
+- **`kj_history_spider.py`**:**一次性历史全量脚本**(跑一次即止,无 schedule)。`from kj_daily_spider import ...` 复用全部公共函数,只写 `history_main` 调度:商户发现 → `get_sold_list(incremental=False)` **全量深翻所有页(不早停)** → 玩家全采 → 视频补采。日常增量与历史全量共用同一套函数,靠 `incremental` 开关区分,不重复维护。
+
+### 常用字段(实测)
+- forsale goodList 项:`merchantAlias`(商户码) / `merchantName` / `goodCode` / `title` / `pic` / `price` / `totalNum` / `currentNum` / `startAt` / `overAt` / `state`
+- detail.good:`startAt` / `overAt` / `state` / `spec{name,content}` / `publisher{alias,name,deal(成交),fans,level}` / `broadcast{playCode,name,state,roomId}`
+- result/merge list 项:`userId`(匿名为 0) / `userName`(匿名为空) / `total`(份数) / `anonymous` / `anonymousCode`
+
+---
+
+## 6. 踩坑记录
+
+- **shade / 重打包包名是老 unpinning 模块的盲区**:uni-app 把 `okhttp3` 改成 `dc.squareup.okhttp3`,原版 JustTrustMe、`frida-multiple-unpinning`、objection 默认脚本按标准类名匹配一律打不中;换 JustTrustMePro / 手动补 hook。
+- **「API 能抓,图片不显示」= WebView 未被 hook**:uni-app 图片走 WebView,需 unpinning 模块覆盖 `onReceivedSslError`(JustTrustMePro 覆盖了)。
+- **Authorization 是每请求签名,不是 token**:抓包抠的固定值只能用几分钟(Timestamp 时效);长期运行必须还原算法实时生成。
+- **签名拼接串里的 path 不含 query、需带 `/api/v4/` 前缀**:`gen_authorization` 内部已处理(补前缀 + 裁 query)。
+- **商品标识是字符串 goodCode(前缀 MC/ZC 等多样),不是数字**:`pid` 列已由 int 改为 `varchar`(product、player 两张表都改)。
+- **playCode 为空 = 回放未就绪**:`broadcast.state=1`「即将拆卡」时 `playCode` 为空,`get_video` 已对空值返回 None。
+- **反调试提醒**:`libsecuritydevice.so`(阿里 dtf)疑似带反 Frida 检测(静态未坐实)。若后续上 Frida,attach 闪退用 `magisk-frida` / 改名 gadget 隐身。
+
+---
+
+## 7. 举一反三
+
+- **uni-app / DCloud 应用逆向套路**:技术栈判定看 `DCloudApplication` + `assets/apps/*/www/app-service.js`;抓包 unpinning 用 JustTrustMePro;业务签名/加密/token 基本都在明文 `app-service.js`,搜关键字符串(签名头名、`sign`、`Authorization`)定位 webpack 模块即可。
+- **「自定义字母表编码的 hash」识别法**:签名是定长 hex、但字节值只在少数几种里循环 → 多半是「标准 hash(MD5/SHA) + 查表/替换编码」。先反解字节还原出合法 hash hex,再回源码找表与拼接串。
+- **判断签名能否纯语言复现**:拼接串里若无隐藏密钥 / appSecret(本例只有 swap/ts/noce/path),即可纯 Python 复现,无需 execjs / JS 引擎;反之考虑抠 JS 用 PyMiniRacer/execjs。
+- **请求签名类接口做爬虫**:优先判断有无时效字段(Timestamp/nonce)与密钥;无密钥的实时签名最省心(免登录、免 token、永不过期)。
+- **不同技术栈 unpinning 对应**:原生 OkHttp → objection / frida-multiple-unpinning(注意 shade);Flutter → reFlutter / hook `ssl_verify_peer_cert`;Native pinning → so 里 hook 校验函数;加固 → 先脱壳。

+ 98 - 0
kaji_spider/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
kaji_spider/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}

+ 117 - 0
kaji_spider/kj_auth.py

@@ -0,0 +1,117 @@
+# -*- coding: utf-8 -*-
+# Author : Charley
+# Python : 3.12.10
+# Date   : 2026/07/02
+"""卡集 App Open-Auth-Sig 请求签名模块。
+
+卡集(card.kaji.com,uni-app / DCloud)的鉴权头 Authorization 采用 Open-Auth-Sig
+方案:每个请求用 [随机 swap + 当前时间戳 + 随机 nonce + 接口路径] 拼串做标准 MD5,
+再经一张 256 项固定置换表编码成 Signature。全程不含任何密钥 / appSecret / 登录 token,
+纯客户端即可复现,因此爬虫可长期无人值守运行,不存在凭证过期问题。
+
+算法来源:app-service.js(APK 反编译产物)模块 "665c" 的 GetCrypto 函数,
+已用真实抓包样本逐字节验证通过。videoPlay 的 body 签名来自偏移 1203166 处的 video_sign。
+"""
+import time
+import random
+import hashlib
+
+# 256 项置换表:来自 app-service.js 偏移约 1513290 处的常量数组 i,是 0..255 的一个全排列。
+# gen_signature 会先按 swap 的 4 对字节对它做两两交换,再逐字符查表。
+I_TABLE = [
+    194, 115, 46, 172, 195, 57, 45, 12, 78, 181, 246, 168, 16, 218, 68, 248,
+    242, 184, 7, 190, 216, 59, 160, 189, 14, 117, 113, 72, 210, 0, 120, 69,
+    15, 48, 141, 128, 34, 155, 19, 38, 110, 63, 76, 240, 35, 49, 56, 47,
+    229, 196, 118, 244, 9, 174, 238, 252, 132, 143, 219, 21, 165, 22, 158, 119,
+    147, 85, 29, 233, 236, 146, 88, 27, 30, 5, 123, 221, 255, 83, 24, 247,
+    70, 111, 11, 4, 31, 159, 103, 6, 185, 52, 98, 75, 95, 142, 50, 3,
+    2, 203, 178, 109, 193, 64, 108, 179, 62, 94, 133, 180, 192, 17, 105, 214,
+    245, 28, 226, 87, 227, 215, 102, 201, 107, 202, 149, 124, 13, 208, 138, 188,
+    1, 54, 127, 139, 239, 209, 162, 243, 126, 144, 204, 100, 130, 187, 131, 207,
+    211, 217, 43, 80, 122, 112, 151, 154, 156, 161, 36, 205, 39, 140, 74, 25,
+    37, 164, 177, 33, 41, 61, 66, 234, 92, 44, 96, 79, 90, 99, 166, 153,
+    182, 150, 20, 106, 77, 93, 157, 135, 152, 104, 223, 173, 81, 134, 175, 230,
+    199, 250, 82, 10, 235, 222, 200, 114, 89, 148, 237, 65, 167, 125, 42, 198,
+    60, 129, 101, 228, 170, 86, 249, 55, 8, 97, 176, 73, 251, 213, 26, 136,
+    254, 169, 145, 225, 186, 32, 212, 197, 163, 241, 84, 67, 91, 253, 51, 224,
+    40, 191, 231, 53, 183, 58, 220, 116, 206, 232, 18, 71, 171, 23, 121, 137
+]
+
+# videoPlay body 签名的固定盐,来自 app-service.js 偏移 1203166:
+#   sign = md5(f"{ts}_{goodCode}_{playCode}_videoPlayKsj")
+VIDEO_SIGN_SALT = "videoPlayKsj"
+
+# 接口统一前缀,对应 app-service.js 里 o.app.requestVersion
+API_PREFIX = "/api/v4/"
+
+
+def _md5_hex(text: str) -> str:
+    """计算字符串的标准 MD5 十六进制摘要。
+
+    Args:
+        text (str): 待哈希的明文串。
+
+    Returns:
+        str: 32 位小写十六进制 MD5 值。
+    """
+    return hashlib.md5(text.encode("utf-8")).hexdigest()
+
+
+def gen_signature(f_hex: str, c: list[int]) -> str:
+    """将 MD5 十六进制串经置换表编码为 128 位 Signature。
+
+    对应 app-service.js 模块 665c 的编码逻辑:先把 32 字符的 md5-hex 当作 ASCII
+    字符串再转 hex(得到 64 字符),复制置换表并按 swap 的 4 对字节做两两交换,
+    最后对每个字符以其 ASCII 码作下标查表,输出 2 位大写 hex 拼接。
+
+    Args:
+        f_hex (str): 拼接串的 MD5 值(32 位小写 hex)。
+        c (list[int]): swap 对应的 8 个字节(0-255),决定置换表的 4 对交换。
+
+    Returns:
+        str: 128 位大写十六进制的 Signature。
+    """
+    r = f_hex.encode("ascii").hex()          # md5-hex 当 ASCII 再转 hex -> 64 字符
+    a = list(I_TABLE)                        # 复制置换表,避免污染全局常量
+    for s in range(0, len(c), 2):            # 按 swap 的 4 对字节两两交换
+        a[c[s]], a[c[s + 1]] = a[c[s + 1]], a[c[s]]
+    out = [format(a[ord(ch)], "02X") for ch in r]   # 逐字符以 ASCII 码查表
+    return "".join(out)
+
+
+def gen_authorization(path: str) -> str:
+    """为指定接口路径生成完整的 Open-Auth-Sig 请求头值。
+
+    每次调用都用当前时间戳与随机 swap / nonce 现算,天然免过期、免登录 token。
+
+    Args:
+        path (str): 接口相对路径,不含域名与 query,如 "goodlist/forsale/main"。
+            函数内部自动补 API_PREFIX 前缀、去掉 "/dataApi" 子串与 query。
+
+    Returns:
+        str: 形如 'Open-Auth-Sig Swap="..",Timestamp=..,Nonce=..,Signature=".."'。
+    """
+    ts = int(time.time())                                       # 当前 Unix 秒
+    noce = random.randint(1, 500)                               # 随机 nonce,闭区间 [1,500]
+    full_path = (API_PREFIX + path).replace("/dataApi", "").split("?")[0]
+    c = random.sample(range(256), 8)                            # 8 个不重复字节
+    swap = "".join(f"{b:02x}" for b in c)                       # 小写 hex,16 字符
+    body = f"{swap}_{ts}_{noce}_{full_path}"                    # 拼接顺序:swap_ts_noce_path
+    sig = gen_signature(_md5_hex(body), c)
+    return f'Open-Auth-Sig Swap="{swap}",Timestamp={ts},Nonce={noce},Signature="{sig}"'
+
+
+def video_sign(ts: int, good_code: str, play_code: str) -> str:
+    """计算 good/videoPlay 接口 body 里的 sign。
+
+    对应 app-service.js 偏移 1203166:sign = md5(f"{ts}_{goodCode}_{playCode}_videoPlayKsj")。
+
+    Args:
+        ts (int): Unix 秒时间戳,与 body 里的 ts 字段保持一致。
+        good_code (str): 商品编码,取自请求 URL good/videoPlay/{goodCode}。
+        play_code (str): 播放码,来自商品详情接口。
+
+    Returns:
+        str: 32 位小写十六进制的 sign。
+    """
+    return _md5_hex(f"{ts}_{good_code}_{play_code}_{VIDEO_SIGN_SALT}")

+ 742 - 0
kaji_spider/kj_daily_spider.py

@@ -0,0 +1,742 @@
+# -*- coding: utf-8 -*-
+# Author : Charley
+# Python : 3.12.10
+# Date   : 2026/07/02
+"""卡集(card.kaji.com)每日采集爬虫。
+
+三级采集管道(结构对齐 zc_new_daily_spider):
+    1. 商户发现:遍历「在售主列表」forsale/main,从每个在售商品里提取商户,写入 kj_shop_record。
+    2. 商品采集:遍历库内每个商户的「已售/已完成列表」merchant?tp=2,逐商品补 detail 后写入 kj_product_record。
+    3. 玩家采集:遍历 kj_product_record 中未采集过玩家的商品,抓 result/merge 中奖/参与名单写入 kj_player_record。
+
+鉴权:卡集接口的 Authorization 采用 Open-Auth-Sig 请求签名(见 kj_auth.py),
+每请求用当前时间实时生成,无登录 token、无过期问题,适合长期无人值守。
+其中商户已售列表 merchant?tp=2 为公开接口,不需要签名。
+"""
+import re
+import sys
+import time
+from datetime import datetime
+from urllib.parse import unquote
+import requests
+import schedule
+from loguru import logger
+from tenacity import retry, stop_after_attempt, wait_fixed
+from mysql_pool import MySQLConnectionPool
+
+import kj_auth
+
+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")
+
+# ==================== 基础配置 ====================
+# 两个业务域名(来源:抓包):主列表 / videoPlay 走 server 域,商户列表 / 详情 / 玩家走 page 域
+SERVER_BASE = "https://server.ssl1.kaji6.com"
+PAGE_BASE = "https://page.ssl1.kaji6.com"
+
+PAGE_SIZE = 10  # 列表接口每页条数
+PLAYER_PAGE_SIZE = 30  # 玩家名单每页条数(抓包实测为 30)
+MAX_PAGES = 100  # 单个列表翻页保护上限,防异常时无限翻页
+
+# 是否为每个商品补拉 detail(拿开始/结束时间、商户销量/粉丝、规格)。
+# 关闭可大幅减少请求量,但 start_date/end_date/规格/商户 sold_number/fans 会缺失。
+FETCH_DETAIL = True
+# 是否使用代理。卡集实测可直连;若遇 IP 风控再置 True 并配置 get_proxys。
+USE_PROXY = False
+
+# 固定 UA(来源:抓包,卡集为 uni-app WebView)
+UA = ("Mozilla/5.0 (Linux; Android 11; Pixel 5 Build/RQ3A.211001.001; wv) "
+      "AppleWebKit/537.36 (KHTML, like Gecko) Version/4.0 Chrome/148.0.7778.120 "
+      "Mobile Safari/537.36 uni-app Html5Plus/1.0 (Immersed/52.727272)")
+
+# 基础请求头(Authorization 每请求单独生成,不放这里)
+BASE_HEADERS = {
+    "deviceType": "phone",
+    "Accept": "application/json, text/plain, */*",
+    "plat": "android",
+    "version": "2.5.39",  # 来源:抓包 App 版本
+    "appVersionCode": "10003",  # 来源:抓包 App versionCode
+    "user-agent": UA,
+}
+
+
+def after_log(retry_state):
+    """tenacity 重试回调,记录每次尝试的结果。
+
+    Args:
+        retry_state: tenacity 传入的 RetryCallState 对象,含调用参数与结果。
+    """
+    # 约定业务函数首个位置参数为 log;取不到时回退全局 logger
+    if retry_state.args and len(retry_state.args) > 0:
+        log = retry_state.args[0]
+    else:
+        log = 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):
+    """获取隧道代理配置(默认不启用,见 USE_PROXY)。
+
+    Args:
+        log: 日志对象。
+
+    Returns:
+        dict: requests 可用的 proxies 字典。
+
+    Raises:
+        Exception: 组装代理配置异常时向上抛出以触发重试。
+    """
+    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
+
+
+def ts_to_dt(ts) -> str | None:
+    """将秒级 Unix 时间戳转成 'YYYY-MM-DD HH:MM:SS' 字符串。
+
+    Args:
+        ts (int | None): 秒级时间戳;为 None / 0 时返回 None。
+
+    Returns:
+        str | None: 格式化时间字符串;无有效时间戳时返回 None。
+    """
+    if not ts:
+        return None
+    return datetime.fromtimestamp(ts).strftime("%Y-%m-%d %H:%M:%S")
+
+
+def parse_desc(desc: str) -> tuple[str | None, int | None]:
+    """从 detail.good.desc(URL 编码)中提取拼团系列与拼团天数。
+
+    desc 解码后各字段以回车(\\r)分隔,形如:
+        拼团系列:27-28
+        拼团规格:1张/包,1包/盒,1盒/箱,共3包 3张
+        拼团份数:1205份
+        拼团时间:15天
+
+    Args:
+        desc (str): detail.good.desc 原始 URL 编码串;为空时返回 (None, None)。
+
+    Returns:
+        tuple[str | None, int | None]: (拼团系列, 拼团天数)。系列为文本(如 "27-28"),
+            天数为整数(如 15);对应项提取失败为 None。
+    """
+    if not desc:
+        return None, None
+    text = unquote(desc)  # URL 解码:%EF%BC%9A→:  %0D→\r
+    m_series = re.search(r"拼团系列[::]\s*(.*?)\s*(?:[\r\n]|$)", text)
+    m_days = re.search(r"拼团时间[::]\s*(\d+)", text)  # 只取数字,"15天"→15
+    series = m_series.group(1).strip() if m_series else None
+    days = int(m_days.group(1)) if m_days else None
+    return series, days
+
+
+@retry(stop=stop_after_attempt(5), wait=wait_fixed(1), after=after_log)
+def kj_request(log, path: str, method: str = "GET", base: str = PAGE_BASE,
+               params: dict = None, json_body: dict = None,
+               need_auth: bool = True, extra_headers: dict = None):
+    """卡集通用请求函数(带重试)。
+
+    根据 need_auth 决定是否为该接口路径实时生成 Open-Auth-Sig 签名并注入 Authorization。
+
+    Args:
+        log: 日志对象。
+        path (str): 接口相对路径,不含域名 / '/api/v4/' 前缀 / query,如 "goodlist/forsale/main"。
+        method (str, optional): 请求方法,"GET" 或 "POST"。Defaults to "GET"。
+        base (str, optional): 域名前缀,SERVER_BASE 或 PAGE_BASE。Defaults to PAGE_BASE。
+        params (dict, optional): URL query 参数。Defaults to None。
+        json_body (dict, optional): POST 的 JSON body。Defaults to None。
+        need_auth (bool, optional): 是否注入 Authorization 签名。Defaults to True。
+        extra_headers (dict, optional): 追加 / 覆盖的请求头。Defaults to None。
+
+    Returns:
+        dict | None: 响应 JSON;非 200 时抛异常触发重试。
+
+    Raises:
+        RuntimeError: HTTP 状态码非 200 时抛出。
+    """
+    url = f"{base}/api/v4/{path}"
+    req_headers = BASE_HEADERS.copy()
+    if need_auth:
+        # 签名基于不含 query 的 path,kj_auth 内部会补 '/api/v4/' 前缀并裁掉 query
+        req_headers["Authorization"] = kj_auth.gen_authorization(path)
+    if extra_headers:
+        req_headers.update(extra_headers)
+
+    proxies = get_proxys(log) if USE_PROXY else None
+    if method.upper() == "POST":
+        resp = requests.post(url, headers=req_headers, json=json_body, params=params, timeout=(5, 30), proxies=proxies)
+    else:
+        resp = requests.get(url, headers=req_headers, params=params, timeout=(5, 30), proxies=proxies)
+
+    if resp.status_code != 200:
+        log.error(f"请求失败 {resp.status_code}: {url}")
+        raise RuntimeError(f"HTTP {resp.status_code}")
+    return resp.json()
+
+
+# ==================== 一、商户发现(在售主列表) ====================
+def get_forsale_page(log, fetch_from: int, fetch_size: int = PAGE_SIZE):
+    """获取一页在售主列表 goodlist/forsale/main。
+
+    Args:
+        log: 日志对象。
+        fetch_from (int): 游标起始位置(首页为 1,逐页按 fetch_size 递增)。
+        fetch_size (int, optional): 每页条数。Defaults to PAGE_SIZE。
+
+    Returns:
+        dict | None: 响应 JSON(含 goodList / isFetchEnd);失败返回 None。
+    """
+    return kj_request(
+        log, "goodlist/forsale/main", base=SERVER_BASE,
+        params={"fetchFrom": str(fetch_from), "fetchSize": str(fetch_size)},
+        need_auth=True,
+    )
+
+
+def save_shops(log, good_list: list, sql_pool, seen: set) -> int:
+    """从在售商品列表提取商户并写入 kj_shop_record(存在则更新店名)。
+
+    在售商品项只带 merchantAlias / merchantName,故本步只写 shop_id + shop_name;
+    sold_number / fans 由商品采集阶段的 detail.publisher 回填。
+
+    Args:
+        log: 日志对象。
+        good_list (list[dict]): forsale 返回的 goodList,每项含 merchantAlias / merchantName。
+        sql_pool (MySQLConnectionPool): MySQL 连接池。
+        seen (set): 跨页去重的 shop_id 集合,避免重复写库。
+
+    Returns:
+        int: 本次实际写库的商户数(已去重)。
+    """
+    info_list = []
+    # print(good_list)
+    for item in good_list:
+        shop_id = item.get("merchantAlias")
+        shop_name = item.get("merchantName")
+        data_dict = {"shop_id": shop_id, "shop_name": shop_name}
+
+        if not shop_id or shop_id in seen: continue
+        seen.add(shop_id)
+        info_list.append(data_dict)
+        # print(data_dict)
+
+    if not info_list:
+        log.info("无新商户,不写库")
+        return 0
+
+    # 存在则更新店名并把 is_deleted 刷回 0:能出现在在售列表 = 商户在营业
+    # (曾被标记注销的商户在此复活,重新纳入采集)
+    sql = ("INSERT INTO kj_shop_record (shop_id, shop_name) VALUES (%s, %s) "
+           "ON DUPLICATE KEY UPDATE shop_name = VALUES(shop_name), is_deleted = 0")
+    args_list = [(d["shop_id"], d["shop_name"]) for d in info_list]
+    sql_pool.insert_many(query=sql, args_list=args_list)
+    return len(info_list)
+
+
+def get_shop_list(log, sql_pool) -> int:
+    """遍历在售主列表翻页,发现并写入全部在售商户。
+
+    翻页靠响应的 isFetchEnd 与空页判断,shop_id 唯一键去重兜底翻页游标的不确定性。
+
+    Args:
+        log: 日志对象。
+        sql_pool (MySQLConnectionPool): MySQL 连接池。
+
+    Returns:
+        int: 去重后发现的商户总数。
+    """
+    seen = set()
+    fetch_from = 1
+
+    while fetch_from <= MAX_PAGES * PAGE_SIZE:
+        try:
+            data = get_forsale_page(log, fetch_from, PAGE_SIZE)
+        except Exception as e:
+            log.error(f"在售主列表 fetch_from={fetch_from} 请求失败: {e}")
+            break
+        if not data:
+            break
+
+        good_list = data.get("goodList", [])
+        if not good_list:
+            log.info(f"在售主列表 fetch_from={fetch_from} 无数据,停止翻页")
+            break
+
+        save_shops(log, good_list, sql_pool, seen)
+        log.info(f"在售主列表 fetch_from={fetch_from} 完成,本页 {len(good_list)} 条,累计商户 {len(seen)}")
+
+        if data.get("isFetchEnd") or len(good_list) < PAGE_SIZE:
+            break
+        fetch_from += PAGE_SIZE
+
+    return len(seen)
+
+
+# ==================== 二、商品采集(商户已售列表) ====================
+def get_merchant_sold_page(log, merchant_code: str, page_index: int, page_size: int = PAGE_SIZE):
+    """获取某商户「已售/已完成」列表的一页(merchant?tp=2,公开接口免签名)。
+
+    Args:
+        log: 日志对象。
+        merchant_code (str): 商户编码(kj_shop_record.shop_id,如 MCT7550827)。
+        page_index (int): 页码,从 1 开始。
+        page_size (int, optional): 每页条数。Defaults to PAGE_SIZE。
+
+    Returns:
+        dict | None: 响应 JSON(含 list / totalPage);失败返回 None。
+    """
+    return kj_request(
+        log, f"merchant/1/goodlist/{merchant_code}", base=PAGE_BASE,
+        params={"pageIndex": str(page_index), "pageSize": str(page_size), "tp": "2"},
+        need_auth=False,
+    )
+
+
+def get_good_detail(log, good_code: str) -> dict:
+    """获取商品详情,返回 detail.good(含时间 / 商户 publisher / 规格)。
+
+    Args:
+        log: 日志对象。
+        good_code (str): 商品编码,如 ZC7653941。
+
+    Returns:
+        dict: detail 中的 good 字典;无数据时返回空字典。
+    """
+    j = kj_request(
+        log, f"good/{good_code}/1/detail", base=PAGE_BASE,
+        params={"referer": "MerchantList"}, need_auth=True,
+    )
+    return (j or {}).get("good", {}) or {}
+
+
+def get_video(log, good_code: str, play_code: str) -> str | None:
+    """获取商品「拆卡回放」视频地址。
+
+    play_code 取自 detail.good.broadcast.playCode(为空表示回放未就绪);body 的 sign 由
+    kj_auth.video_sign 生成,videoPlay 接口本身不需要 Authorization。响应的 media_url 即视频地址。
+
+    Args:
+        log: 日志对象。
+        good_code (str): 商品编码。
+        play_code (str): 回放播放码,来自 detail.good.broadcast.playCode。
+
+    Returns:
+        str | None: 回放视频 URL;play_code 为空或接口无视频时返回 None。
+    """
+    if not play_code:
+        return None
+    ts = int(time.time())
+    body = {"playCode": play_code, "sign": kj_auth.video_sign(ts, good_code, play_code), "ts": ts}
+    j = kj_request(log, f"good/videoPlay/{good_code}", method="POST", base=SERVER_BASE,
+                   json_body=body, need_auth=False, extra_headers={"Content-Type": "application/json"})
+    return (j or {}).get("media_url") if j and j.get("code") == 0 else None
+
+
+def build_product(good_item: dict, detail_good: dict, merchant_code: str, merchant_name: str, video_url: str) -> dict:
+    """把商户列表项 + 商品详情组装成 kj_product_record 一行。
+
+    卡集为拼团模型,zc 表中的直播 / 存储 / 多价字段无对应,统一置 None。
+
+    Args:
+        good_item (dict): merchant?tp=2 列表里的商品项(goodCode/title/pic/price/totalNum/currentNum/status)。
+        detail_good (dict): detail.good,补 startAt/overAt/state/规格;detail 失败时为空字典。
+        merchant_code (str): 商户编码。
+        merchant_name (str): 商户名。
+        video_url (str): 回放视频 URL。
+
+    Returns:
+        dict: 与 kj_product_record 列对应的数据字典。
+    """
+    # 时间戳转文本:ts_to_dt(detail_good.get("startAt"))
+    imgs = detail_good.get("pic", {}).get("carousel", [])
+    imgs = ','.join(imgs) if imgs else None
+
+    # 拼团系列 / 拼团时间(天):从 detail.good.desc 解码提取
+    group_series, group_days = parse_desc(detail_good.get("desc", ""))
+
+    row = {
+        "shop_id": merchant_code,  # 商户编码
+        "shop_name": merchant_name,  # 商户名
+        "pid": good_item.get("goodCode"),  #
+        "title": good_item.get("title"),  # 标题
+        "price": good_item.get("price"),  # 价格
+        "total_num": good_item.get("totalNum"),  # 总数量
+        "current_num": good_item.get("currentNum"),  # 当前数量
+        "status": good_item.get("status"),  # 状态
+        "imgs": imgs,  # 图片链接, ','分割
+        "start_at": ts_to_dt(detail_good.get("startAt")),  # 开始时间
+        "over_at": ts_to_dt(detail_good.get("overAt")),  # 结束时间
+        "spec_name": detail_good.get("spec", {}).get("name"),  # 规格配置
+        "spec_content": detail_good.get("spec", {}).get("content"),  # 产品规格
+        "group_series": group_series,  # 拼团系列,如 "27-28"
+        "group_days": group_days,  # 拼团时间(天),如 15
+        # "publisher_sale": detail_good.get("publisher", {}).get("sale"),  # 在售
+        # "publisher_fans": detail_good.get("publisher", {}).get("fans"),  # 出售者粉丝
+        "video_url": video_url,
+    }
+    return row
+
+
+def get_shop_stop_pid(shop_id: str, sql_pool) -> str | None:
+    """取该商户已入库商品里 start_at 最新那条的 pid,作为增量翻页的停止线。
+
+    merchant?tp=2 结果按 start_at 倒序(最新在前),daily 翻页碰到这个 pid 即可早停:
+    它及其之后的商品都是上一轮已采过的。依赖 (shop_id, start_at) 复合索引,单值查询很快。
+
+    Args:
+        shop_id (str): 商户编码。
+        sql_pool (MySQLConnectionPool): MySQL 连接池。
+
+    Returns:
+        str | None: 该商户最新商品的 pid(goodCode);该商户暂无记录时返回 None(触发全量采集)。
+    """
+    row = sql_pool.select_one(
+        "SELECT pid FROM kj_product_record WHERE shop_id = %s ORDER BY start_at DESC LIMIT 1",
+        (shop_id,),
+    )
+    return row[0] if row else None
+
+
+def get_sold_list(log, shop_id: str, shop_name: str, sql_pool, incremental: bool = True) -> int:
+    """遍历某商户已售列表翻页,逐商品补详情后写入 kj_product_record。
+
+    Args:
+        log: 日志对象。
+        shop_id (str): 商户编码(shop_id)。
+        shop_name (str): 商户名。
+        sql_pool (MySQLConnectionPool): MySQL 连接池。
+        incremental (bool, optional): True(daily) 按 start_at 最新 pid 早停,只采新商品;
+            False(history) 全量深翻所有页。Defaults to True。
+
+    Returns:
+        int: 本商户写入的商品数。
+    """
+    page_index = 1
+    saved = 0
+    stopped = False
+    # 增量停止线:该商户上次采到的最新商品 pid(接口按 start_at 新在前,翻页碰到它即早停)
+    stop_pid = get_shop_stop_pid(shop_id, sql_pool) if incremental else None
+
+    while page_index <= MAX_PAGES:
+        try:
+            data = get_merchant_sold_page(log, shop_id, page_index)
+        except Exception as e:
+            log.error(f"商户 {shop_id} 已售列表第 {page_index} 页请求失败: {e}")
+            break
+
+        if not data:
+            log.info(f"商户 {shop_id} 已售列表第 {page_index} 页无数据,停止翻页")
+            break
+
+        # 商户无效(已注销 / 不存在):merchant?tp=2 返回 code=1、msg='无效商家'。
+        # 首页即无效 → 标记 is_deleted=1;之后 kj_main 的「WHERE is_deleted = 0」会自动跳过它,不再浪费请求。
+        if data.get("code") != 0:
+            if page_index == 1:
+                log.info(f"商户 {shop_id} 无效({data.get('msg')}),标记 is_deleted=1")
+                sql_pool.update_one("UPDATE kj_shop_record SET is_deleted = 1 WHERE shop_id = %s", (shop_id,))
+            break
+
+        good_list = data.get("list", [])
+        if not good_list:
+            log.info(f"商户 {shop_id} 已售列表第 {page_index} 页无数据,停止翻页")
+            break
+        total_page = data.get("totalPage", 1)
+
+        batch = []
+        for gi in good_list:
+            code = gi.get("goodCode")
+            # 增量早停:碰到上次采到的最新商品,它及之后都是已采过的旧数据
+            if stop_pid and code == stop_pid:
+                log.info(f"商户 {shop_id} 碰到停止线 pid={stop_pid},增量早停")
+                stopped = True
+                break
+            detail_good = {}
+            if FETCH_DETAIL:
+                try:
+                    detail_good = get_good_detail(log, code)
+                except Exception as e:
+                    log.error(f"商品 {code} detail 请求失败: {e}")
+
+            # 如需回放视频地址(media_url):playCode 在 detail_good["broadcast"]["playCode"](空=未就绪)
+            # 首次拿不到不影响主数据 —— refill_videos 会在拼团完成后兜底补采
+            try:
+                video_url = get_video(log, code, (detail_good.get("broadcast") or {}).get("playCode"))
+            except Exception as e:
+                log.error(f"商品 {code} 视频获取失败: {e}")
+                video_url = None
+
+            row = build_product(gi, detail_good, shop_id, shop_name, video_url)
+            # print(row)
+            if row:
+                batch.append(row)
+
+        if batch:
+            # 已存在(pid 唯一)则跳过:已完成商品为终态,无需覆盖
+            sql_pool.insert_many(table="kj_product_record", data_list=batch, ignore=True)
+            saved += len(batch)
+
+        if stopped:
+            log.info(f"商户 {shop_id} 增量早停,停止翻页")
+            break
+
+        log.info(f"商户 {shop_id} 已售列表第 {page_index}/{total_page} 页完成,本页 {len(good_list)} 商品")
+        page_index += 1
+        if page_index > total_page:
+            log.info(f"商户 {shop_id} 已售列表翻页完成")
+            break
+
+    return saved
+
+
+# ==================== 三、玩家采集(中奖/参与名单) ====================
+def get_player_page(log, good_code: str, fetch_from: int, fetch_size: int = PLAYER_PAGE_SIZE):
+    """获取某商品玩家名单的一页 good/{code}/result/merge。
+
+    Args:
+        log: 日志对象。
+        good_code (str): 商品编码。
+        fetch_from (int): 游标起始位置(首页为 1,逐页按 fetch_size 递增)。
+        fetch_size (int, optional): 每页条数。Defaults to PLAYER_PAGE_SIZE。
+
+    Returns:
+        dict | None: 响应 JSON(含 code / list / isFetchEnd);失败返回 None。
+    """
+    return kj_request(
+        log, f"good/{good_code}/result/merge", base=PAGE_BASE,
+        params={"fetchFrom": str(fetch_from), "fetchSize": str(fetch_size), "q": ""},
+        need_auth=True,
+    )
+
+
+def save_players(log, good_code: str, player_list: list, sql_pool) -> int:
+    """解析玩家名单并写入 kj_player_record。
+
+    匿名玩家(userId=0、userName 为空)用 anonymousCode 兜底作为标识。
+
+    Args:
+        log: 日志对象。
+        good_code (str): 商品编码(写入 pid 列)。
+        player_list (list[dict]): result/merge 的 list,每项含 userId/userName/total/anonymousCode。
+        sql_pool (MySQLConnectionPool): MySQL 连接池。
+
+    Returns:
+        int: 本次写入的玩家记录数。
+    """
+    log.info(f"开始保存商品 {good_code} 的玩家名单")
+    info_list = []
+    for item in player_list:
+        anon = item.get("anonymousCode")
+        user_id = item.get("userId")
+        user_name = item.get("userName")
+        if not user_name:  # 匿名玩家 userName 为空
+            user_name = f"匿名_{anon}" if anon else None
+        if not user_id:  # 匿名玩家 userId=0,用匿名码兜底
+            user_id = anon
+        info_list.append({
+            "pid": good_code,  # 商品编码,来自函数参数(item 里没有商品标识)
+            "give_number": item.get("total"),  # 该玩家份数:卡集字段是 total
+            "user_id": str(user_id) if user_id is not None else None,
+            "user_name": user_name,
+        })
+
+    if info_list:
+        sql_pool.insert_many(table="kj_player_record", data_list=info_list)
+    return len(info_list)
+
+
+def get_player_list(log, good_code: str, sql_pool) -> bool:
+    """遍历某商品玩家名单翻页并写库。
+
+    Args:
+        log: 日志对象。
+        good_code (str): 商品编码。
+        sql_pool (MySQLConnectionPool): MySQL 连接池。
+
+    Returns:
+        bool: True 表示抓到玩家数据,False 表示无数据(如拼团未完成)。
+    """
+    fetch_from = 1
+    has_data = False
+
+    while fetch_from <= MAX_PAGES * PLAYER_PAGE_SIZE:
+        try:
+            data = get_player_page(log, good_code, fetch_from, PLAYER_PAGE_SIZE)
+        except Exception as e:
+            log.error(f"商品 {good_code} 玩家名单 fetch_from={fetch_from} 请求失败: {e}")
+            break
+        if not data:
+            log.info(f"商品 {good_code} 玩家名单 fetch_from={fetch_from} 无数据,停止翻页")
+            break
+
+        if data.get("code") != 0:
+            # code=1 常见于"拼团未完成的商品",视为暂无玩家
+            log.info(f"商品 {good_code} 暂无玩家: {data.get('msg')}")
+            break
+
+        plist = data.get("list", [])
+        if not plist:
+            log.info(f"商品 {good_code} 玩家名单翻页完成")
+            break
+
+        has_data = True
+        save_players(log, good_code, plist, sql_pool)
+
+        if data.get("isFetchEnd") or len(plist) < PLAYER_PAGE_SIZE:
+            log.info(f"商品 {good_code} 玩家名单翻页完成")
+            break
+        fetch_from += PLAYER_PAGE_SIZE
+
+    return has_data
+
+
+# ==================== 四、视频补采 ====================
+def refill_videos(log, sql_pool) -> int:
+    """补采视频地址:对拼团已完成(player_state=1)但 video_url 仍空的商品,
+    重新拉 detail 取 broadcast.playCode → get_video → 更新 video_url。
+
+    覆盖场景:首次入库时该商品还处于「即将拆卡 / 正在拆卡」等中间状态,
+    broadcast.playCode 为空、视频拿不到;拼团完成、回放就绪后由本函数补回。
+    playCode 仍空则跳过,留到下一轮再试,直到 video_url 有值。
+
+    Args:
+        log: 日志对象。
+        sql_pool (MySQLConnectionPool): MySQL 连接池。
+
+    Returns:
+        int: 本轮成功补采的视频数。
+    """
+    rows = sql_pool.select_all(
+        "SELECT pid FROM kj_product_record WHERE video_url IS NULL AND player_state = 1"
+    )
+    pids = [r[0] for r in rows] if rows else []
+    log.info(f"待补视频商品 {len(pids)} 个")
+    filled = 0
+    for pid in pids:
+        try:
+            detail = get_good_detail(log, pid)
+            play_code = (detail.get("broadcast") or {}).get("playCode")
+            if not play_code:
+                # 回放仍未就绪(正在拆卡 / 即将拆卡等),留到下一轮
+                continue
+            url = get_video(log, pid, play_code)
+            if url:
+                sql_pool.update_one(
+                    "UPDATE kj_product_record SET video_url = %s WHERE pid = %s",
+                    (url, pid),
+                )
+                filled += 1
+                log.info(f"商品 {pid} 视频已补采 → {url}")
+        except Exception as e:
+            log.error(f"商品 {pid} 视频补采失败: {e}")
+    return filled
+
+
+# ==================== 主流程 ====================
+@retry(stop=stop_after_attempt(100), wait=wait_fixed(3600), after=after_log)
+def kj_main(log):
+    """卡集每日采集主函数:商户发现 → 商品采集 → 玩家采集。
+
+    Args:
+        log: 日志对象。
+
+    Raises:
+        RuntimeError: 数据库连接池异常时抛出以触发重试。
+    """
+    log.info(f"开始运行 {sys._getframe().f_code.co_name} 卡集采集任务" + "." * 40)
+
+    sql_pool = MySQLConnectionPool(log=log)
+    if not sql_pool.check_pool_health():
+        log.error("数据库连接池异常")
+        raise RuntimeError("数据库连接池异常")
+
+    try:
+        # 1) 商户发现
+        try:
+            n = get_shop_list(log, sql_pool)
+            log.info(f"商户发现完成,去重商户 {n} 个")
+        except Exception as e:
+            log.error(f"get_shop_list error: {e}")
+
+        time.sleep(5)
+
+        # 2) 商品采集:遍历库内所有商户的已售列表
+        try:
+            shop_rows = sql_pool.select_all("SELECT shop_id, shop_name FROM kj_shop_record WHERE is_deleted = 0")
+            log.info(f"待采集商户 {len(shop_rows)} 个")
+            for shop_id, shop_name in shop_rows:
+                try:
+                    cnt = get_sold_list(log, shop_id, shop_name, sql_pool)
+                    log.info(f"商户 {shop_id} {shop_name} 商品采集完成,写入 {cnt} 个")
+                except Exception as e:
+                    log.error(f"get_sold_list error(商户 {shop_id}): {e}")
+        except Exception as e:
+            log.error(f"iterate_shop_list error: {e}")
+
+        time.sleep(5)
+
+        # 3) 玩家采集:遍历尚未成功采集玩家的商品
+        try:
+            prod_rows = sql_pool.select_all("SELECT pid FROM kj_product_record WHERE player_state != 1")
+            pids = [row[0] for row in prod_rows] if prod_rows else []
+            log.info(f"待采集玩家的商品 {len(pids)} 个")
+            for pid in pids:
+                try:
+                    # 先置 1 表示开始采集(对齐 zc 断点标记)
+                    sql_pool.update_one("UPDATE kj_product_record SET player_state = 1 WHERE pid = %s", (pid,))
+                    has_data = get_player_list(log, pid, sql_pool)
+                    if not has_data:
+                        # 无玩家(如拼团未完成)置 2,下轮仍会重试
+                        sql_pool.update_one("UPDATE kj_product_record SET player_state = 2 WHERE pid = %s", (pid,))
+                except Exception as pid_error:
+                    log.error(f"商品 {pid} 玩家采集失败: {pid_error}")
+                    try:
+                        sql_pool.update_one("UPDATE kj_product_record SET player_state = 3 WHERE pid = %s", (pid,))
+                    except Exception as update_error:
+                        log.error(f"更新商品 {pid} 状态失败: {update_error}")
+        except Exception as e:
+            log.error(f"iterate_player_list error: {e}")
+
+        # 4) 视频补采:玩家采集之后,对拼团已完成(player_state=1)但 video_url 仍空的商品重取视频
+        try:
+            n = refill_videos(log, sql_pool)
+            log.info(f"视频补采完成,本轮 {n} 条")
+        except Exception as e:
+            log.error(f"refill_videos error: {e}")
+    except Exception as e:
+        log.error(f"{sys._getframe().f_code.co_name} error: {e}")
+    finally:
+        log.info(f"卡集采集 {sys._getframe().f_code.co_name} 运行结束,等待下一轮" + "." * 20)
+
+
+def schedule_task():
+    """定时任务入口:每天 00:01 运行一次 kj_main。"""
+    # 立即运行一次(调试时取消注释)
+    kj_main(log=logger)
+
+    schedule.every().day.at("00:01").do(kj_main, log=logger)
+    while True:
+        schedule.run_pending()
+        time.sleep(1)
+
+
+if __name__ == "__main__":
+    # kj_main(logger)
+    schedule_task()

+ 100 - 0
kaji_spider/kj_history_spider.py

@@ -0,0 +1,100 @@
+# -*- coding: utf-8 -*-
+# Author : Charley
+# Python : 3.12.10
+# Date   : 2026/07/02
+"""卡集历史全量采集脚本(一次性)。
+
+与 kj_daily_spider(每天增量)的区别:本脚本对每个商户的已售列表「全量深翻」所有页
+(调 get_sold_list(incremental=False),不走 start_at 停止线早停),把历史商品 / 玩家 / 视频
+一次性灌满。跑一次即止,不带 schedule。
+
+公共函数(请求 / 签名 / 解析 / 翻页 / 补采)全部从 kj_daily_spider 复用,避免重复维护。
+
+运行:python kj_history_spider.py
+"""
+import sys
+
+from loguru import logger
+from mysql_pool import MySQLConnectionPool
+
+from kj_daily_spider import (
+    get_shop_list,
+    get_sold_list,
+    get_player_list,
+    refill_videos,
+)
+
+
+def history_main(log):
+    """历史全量采集主函数:商户发现 → 全量深翻商品 → 全量采玩家 → 补视频。
+
+    Args:
+        log: 日志对象。
+
+    Raises:
+        RuntimeError: 数据库连接池异常时抛出。
+    """
+    log.info(f"开始运行 {sys._getframe().f_code.co_name} 卡集历史全量采集" + "." * 40)
+
+    sql_pool = MySQLConnectionPool(log=log)
+    if not sql_pool.check_pool_health():
+        log.error("数据库连接池异常")
+        raise RuntimeError("数据库连接池异常")
+
+    try:
+        # 1) 商户发现:先用在售主列表把当前活跃商户补进 kj_shop_record(历史脚本也要完整商户清单)
+        try:
+            n = get_shop_list(log, sql_pool)
+            log.info(f"商户发现完成,去重商户 {n} 个")
+        except Exception as e:
+            log.error(f"get_shop_list error: {e}")
+
+        # 2) 商品全量:遍历所有商户,深翻已售列表全部页(incremental=False,不早停)
+        try:
+            shop_rows = sql_pool.select_all("SELECT shop_id, shop_name FROM kj_shop_record WHERE is_deleted = 0")
+            log.info(f"待全量采集商户 {len(shop_rows)} 个")
+            for shop_id, shop_name in shop_rows:
+                try:
+                    cnt = get_sold_list(log, shop_id, shop_name, sql_pool, incremental=False)
+                    log.info(f"商户 {shop_id} {shop_name} 全量商品采集完成,写入 {cnt} 个")
+                except Exception as e:
+                    log.error(f"get_sold_list error(商户 {shop_id}): {e}")
+        except Exception as e:
+            log.error(f"iterate_shop_list error: {e}")
+
+        # 3) 玩家全量:遍历尚未成功采集玩家的商品(player_state != 1)
+        try:
+            prod_rows = sql_pool.select_all("SELECT pid FROM kj_product_record WHERE player_state != 1")
+            pids = [row[0] for row in prod_rows] if prod_rows else []
+            log.info(f"待采集玩家的商品 {len(pids)} 个")
+            for pid in pids:
+                try:
+                    # 先置 1 表示开始采集(对齐 zc 断点标记)
+                    sql_pool.update_one("UPDATE kj_product_record SET player_state = 1 WHERE pid = %s", (pid,))
+                    has_data = get_player_list(log, pid, sql_pool)
+                    if not has_data:
+                        # 无玩家(如拼团未完成)置 2,下次仍会重试
+                        sql_pool.update_one("UPDATE kj_product_record SET player_state = 2 WHERE pid = %s", (pid,))
+                except Exception as pid_error:
+                    log.error(f"商品 {pid} 玩家采集失败: {pid_error}")
+                    try:
+                        sql_pool.update_one("UPDATE kj_product_record SET player_state = 3 WHERE pid = %s", (pid,))
+                    except Exception as update_error:
+                        log.error(f"更新商品 {pid} 状态失败: {update_error}")
+        except Exception as e:
+            log.error(f"iterate_player_list error: {e}")
+
+        # 4) 视频补采:拼团已完成但 video_url 仍空的商品,重取回放地址
+        try:
+            n = refill_videos(log, sql_pool)
+            log.info(f"视频补采完成,本轮 {n} 条")
+        except Exception as e:
+            log.error(f"refill_videos error: {e}")
+    except Exception as e:
+        log.error(f"{sys._getframe().f_code.co_name} error: {e}")
+    finally:
+        log.info(f"卡集历史全量采集 {sys._getframe().f_code.co_name} 运行结束")
+
+
+if __name__ == "__main__":
+    history_main(logger)

+ 671 - 0
kaji_spider/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)

+ 8 - 0
kaji_spider/requirements.txt

@@ -0,0 +1,8 @@
+-i https://mirrors.aliyun.com/pypi/simple/
+DBUtils==3.1.2
+loguru==0.7.3
+PyMySQL==1.1.2
+PyYAML==6.0.3
+requests==2.33.1
+schedule==1.2.2
+tenacity==9.1.4