| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764 |
- # -*- coding: utf-8 -*-
- # Author : Charley
- # Python : 3.12.10
- # Date : 2026/07/06
- """潮谷街(吃谷商城)每日增量采集爬虫。
- 三级采集管道(结构对齐 kaji 的 kj_daily_spider):
- 1. 商户发现:翻页遍历「shop 列表页」group_buying_products(product_types/category_id 模式),
- 从每个商品里提取 merchant_id / merchant_name,写入 cgj_shop_record。
- 2. 商品采集:遍历库内每个商户的「历史成交」group_buying_products(merchant_id + status 模式),
- 逐商品补 product/info 详情后写入 cgj_product_record。
- 3. 玩家采集:遍历 cgj_product_record 中未采集过玩家的商品,抓 announcement_by_user
- 的买家名单写入 cgj_player_record。
- 鉴权:潮谷街这几个接口 access-token 为空、无签名/加密,属公开(游客)接口,
- 只需带固定的 devicecode / version 等请求头即可,适合长期无人值守。
- 翻页约定(三个列表接口统一):请求体/参数带 cursor(首页为空串),响应返回
- data.list(本页数据)、data.cursor(下一页游标)、data.done(true=到底)。
- 循环条件:done 为 True 或 cursor 为空即停止。
- 【字段解析】save_shops / build_product / save_players 里的 `for item in items:` 循环中,
- """
- import sys
- import time
- from datetime import datetime
- import requests
- import schedule
- from loguru import logger
- from tenacity import retry, stop_after_attempt, wait_fixed
- from mysql_pool import MySQLConnectionPool
- 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")
- # ==================== 基础配置 ====================
- # 业务域名(来源:抓包,四个接口都走这个域)
- BASE = "https://5fdc92.chaogujieapp.com"
- PAGE_SIZE = 20 # 列表接口每页条数(抓包实测 20)
- PLAYER_PAGE_SIZE = 20 # 玩家名单每页条数(抓包实测 20)
- # MAX_PAGES = 500 # 单个列表翻页保护上限,防异常时无限翻页
- # 增量早停阈值:某商户历史成交连续这么多页「全是已入库商品(0 新)」即停止翻页。
- # 用「连续多页」而非「单页」,兼容历史成交列表里新商品跨页穿插的情况;调大更稳、调小更省。
- STOP_AFTER_DUPE_PAGES = 2
- MAX_PAGES = 5 # 单个列表翻页保护上限,防异常时无限翻页 每日任务时设置更小限制
- # 「历史成交」过滤值(来源:抓包固定 status=1796)。九尾等仅有在售无历史成交的商户返回 0 条,符合预期。
- HISTORY_STATUS = 1796
- # 设备码 / App 版本(来源:抓包,公开接口无需登录 token)
- DEVICE_CODE = "8420C51411ECF16C802600868F2D214A"
- APP_VERSION = "1.0.15"
- # 是否为每个商品补拉 product/info 详情(拿系列/规格/简介/视频/上下架时间等)。
- # 关闭可大幅减少请求量,但详情维度字段会缺失。
- FETCH_DETAIL = True
- # 是否为每个商品补拉直播回放(live/user_live/detail,room_id 取自详情的 live_room_id)。
- # 依赖 FETCH_DETAIL=True——要先有详情才能拿到 live_room_id。
- FETCH_REPLAY = True
- # 是否使用代理。潮谷街实测可直连;若遇 IP 风控再置 True 并配置 get_proxys。
- USE_PROXY = False
- # 固定 UA(来源:抓包,潮谷街为 Flutter 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")
- # 基础请求头(来源:抓包)。POST 的 content-type 由 requests 的 json= 自动补,不放这里。
- BASE_HEADERS = {
- "user-agent": UA,
- "accept-encoding": "gzip",
- "channel-version": APP_VERSION, # 抓包 channel-version
- "access-token": "", # 抓包为空:公开/游客接口,无需登录
- "source-id": "1", # 抓包固定
- "language": "zh", # 抓包固定
- "version": APP_VERSION, # 抓包 App 版本
- "devicecode": DEVICE_CODE, # 抓包设备码
- "platform-id": "1", # 抓包固定(1=Android)
- "channel-name": "huawei", # 抓包渠道
- }
- 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 | str | None): 秒级时间戳;为 None / 0 / 空时返回 None。
- Returns:
- str | None: 格式化时间字符串;无有效时间戳时返回 None。
- """
- if not ts:
- return None
- try:
- return datetime.fromtimestamp(int(ts)).strftime("%Y-%m-%d %H:%M:%S")
- except (ValueError, TypeError, OSError):
- return None
- @retry(stop=stop_after_attempt(5), wait=wait_fixed(1), after=after_log)
- def cgj_request(log, path: str, method: str = "GET",
- params: dict = None, json_body: dict = None,
- extra_headers: dict = None) -> dict | None:
- """潮谷街通用请求函数(带重试)。
- 潮谷街接口无签名/加密,仅需固定请求头,故本函数只负责拼 URL、发请求、判状态码。
- Args:
- log: 日志对象。
- path (str): 接口相对路径(以 / 开头,不含域名),如 "/api/home/group_buying_products"。
- method (str, optional): 请求方法,"GET" 或 "POST"。Defaults to "GET"。
- params (dict, optional): URL query 参数。Defaults to None。
- json_body (dict, optional): POST 的 JSON body。Defaults to None。
- extra_headers (dict, optional): 追加 / 覆盖的请求头。Defaults to None。
- Returns:
- dict | None: 响应 JSON;非 200 时抛异常触发重试。
- Raises:
- RuntimeError: HTTP 状态码非 200 时抛出。
- """
- url = f"{BASE}{path}"
- req_headers = BASE_HEADERS.copy()
- 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_group_buying_page(log, cursor: str = "", merchant_id=None) -> dict | None:
- """获取 group_buying_products 一页(发现模式 / 商户历史成交模式二合一)。
- 该接口一站两用:
- - 不传 merchant_id:全局「shop 列表页」,用 product_types/category_id 过滤,用于商户发现。
- - 传 merchant_id:某商户「历史成交」,用 status=HISTORY_STATUS 过滤。
- 两种模式响应结构一致(data.list / data.cursor / data.done)。
- Args:
- log: 日志对象。
- cursor (str, optional): 翻页游标,首页传空串。Defaults to ""。
- merchant_id (int | str | None, optional): 商户 id;None=商户发现模式。Defaults to None。
- Returns:
- dict | None: 响应 JSON;失败时由 cgj_request 抛异常触发重试。
- """
- if merchant_id is not None:
- body = {
- "cursor": cursor,
- "size": PAGE_SIZE,
- "image_size": "68",
- "status": HISTORY_STATUS,
- "merchant_id": merchant_id,
- }
- else:
- body = {
- "cursor": cursor,
- "size": PAGE_SIZE,
- "image_size": "178",
- "product_types": [1],
- "category_id": 1,
- }
- return cgj_request(log, "/api/home/group_buying_products", method="POST", json_body=body)
- # ==================== 一、商户发现(shop 列表页) ====================
- def save_shops(log, items: list, sql_pool, seen: set) -> int:
- """从 shop 列表页的商品项提取商户并写入 cgj_shop_record(存在则更新店名并复活)。
- ------------------------------------------------------------------
- 商户维度目前只取驱动管道必需的 merchant_id / merchant_name,
- 如需更多商户字段(如 merchant_image_url),在下面 for 循环里往 data_dict 加即可,
- 并同步 schema.sql 的 cgj_shop_record 列与下方 upsert 的 SQL。
- ------------------------------------------------------------------
- Args:
- log: 日志对象。
- items (list[dict]): group_buying_products 返回的 data.list,每项含 merchant_id / merchant_name。
- sql_pool (MySQLConnectionPool): MySQL 连接池。
- seen (set): 跨页去重的 merchant_id 集合,避免重复写库。
- Returns:
- int: 本次实际写库的商户数(已去重)。
- """
- info_list = []
- for item in items:
- merchant_id = item.get("merchant_id") # 商户 id(管道必需)
- merchant_name = item.get("merchant_name") # 商户名(管道必需)
- data_dict = {"shop_id": merchant_id, "shop_name": merchant_name}
- if not merchant_id or merchant_id in seen:
- log.info(f"商户 {merchant_id} 已存在 seen,跳过")
- continue
- seen.add(merchant_id)
- info_list.append(data_dict)
- if not info_list:
- log.info("本页无新商户,不写库")
- return 0
- # 存在则更新店名并把 is_deleted 刷回 0:能出现在列表 = 商户在营业(曾注销的在此复活)
- sql = ("INSERT INTO cgj_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]
- # print(args_list)
- sql_pool.insert_many(query=sql, args_list=args_list)
- return len(info_list)
- def get_shop_list(log, sql_pool) -> int:
- """翻页遍历 shop 列表页,发现并写入全部商户。
- 翻页靠响应的 done / cursor 判断,merchant_id 唯一键去重兜底。
- Args:
- log: 日志对象。
- sql_pool (MySQLConnectionPool): MySQL 连接池。
- Returns:
- int: 去重后发现的商户总数。
- """
- seen = set()
- cursor = ""
- page = 0
- while page < MAX_PAGES:
- try:
- resp = get_group_buying_page(log, cursor, merchant_id=None)
- except Exception as e:
- log.error(f"shop 列表页 cursor={cursor!r} 请求失败: {e}")
- break
- if not resp or resp.get("code") != 0:
- log.info(f"shop 列表页返回异常: {resp.get('msg') if resp else None}")
- break
- data = resp.get("data") or {}
- items = data.get("list") or []
- if not items:
- log.info(f"shop 列表页第 {page + 1} 页无数据,停止翻页")
- break
- save_shops(log, items, sql_pool, seen)
- log.info(f"shop 列表页第 {page + 1} 页完成,本页 {len(items)} 条,累计商户 {len(seen)}")
- if data.get("done") or not data.get("cursor"):
- break
- cursor = data["cursor"]
- page += 1
- return len(seen)
- # ==================== 二、商品采集(商户历史成交 + 详情) ====================
- def get_product_detail(log, product_id) -> dict:
- """获取商品详情 product/info,返回 data 节点(含 product_info / merchant_info 等)。
- Args:
- log: 日志对象。
- product_id (int | str): 商品 id。
- Returns:
- dict: 详情 data 字典;无数据时返回空字典。
- """
- resp = cgj_request(log, "/api/product/info", method="GET", params={"id": product_id})
- return (resp or {}).get("data") or {}
- def get_live_detail(log, room_id: str) -> dict:
- """获取直播回放详情 live/user_live/detail,返回 data 节点。
- room_id 取自商品详情的 data.live_room_id。回放视频地址在 data.offline_address_info[].url
- (腾讯云 vod,name="默认"),另带 status / start_time / end_time / duration 等直播信息。
- Args:
- log: 日志对象。
- room_id (str): 直播间 id(商品详情的 live_room_id)。
- Returns:
- dict: 回放详情 data 字典;room_id 为空或无数据时返回空字典。
- """
- if not room_id:
- return {}
- resp = cgj_request(log, "/api/live/user_live/detail", method="POST", json_body={"room_id": room_id})
- return (resp or {}).get("data") or {}
- def build_product(item: dict, detail: dict, shop_id, shop_name: str, live_detail: dict = None) -> dict:
- """把历史成交列表项 + 商品详情组装成 cgj_product_record 一行。
- ------------------------------------------------------------------
- 列表项字段(item.*)来自 group_buying_products;详情字段(detail.*)来自 product/info。
- ------------------------------------------------------------------
- Args:
- item (dict): 历史成交列表里的商品项(group_buying_products data.list 项)。
- detail (dict): product/info 的 data 节点;FETCH_DETAIL 关闭或失败时为空字典。
- shop_id (int | str): 商户 id。
- shop_name (str): 商户名。
- live_detail (dict, optional): live/user_live/detail 的 data 节点(回放信息);
- FETCH_REPLAY 关闭 / 无 live_room_id / 失败时为空字典。Defaults to None。
- Returns:
- dict: 与 cgj_product_record 列对应的数据字典。
- """
- # 详情里的嵌套节点(按需取用)
- product_info = detail.get("product_info") or {}
- basic_info = product_info.get("basic_info") or {}
- # print(product_info)
- # print('-------------------')
- # print(basic_info)
- # 直播回放(待确认):回放地址在 offline_address_info[].url(腾讯云 vod),取第一个(name="默认")
- live_detail = live_detail or {}
- offline_list = live_detail.get("offline_address_info") or []
- replay_url = offline_list[0].get("url") if offline_list else None
- row = {
- # ---- 管道必需(勿删)----
- "pid": item.get("id"), # 商品 id
- "shop_id": shop_id, # 商户 id
- "shop_name": shop_name, # 商户名
- # ---- 列表项字段(item.*,待确认)----
- "title": item.get("name"), # 商品名
- "price": item.get("price"), # 价格
- "imgs": item.get("image_url"), # 主图
- "stock_quantity": item.get("stock_quantity"), # 库存
- "sold_quantity": item.get("sold_quantity"), # 已售
- "remaining_quantity": item.get("remaining_quantity"), # 剩余
- "card_edition_name": item.get("card_edition_name"), # 卡版名
- # ---- 详情字段(detail.*,待确认;FETCH_DETAIL 关闭时为空)----
- "specs": basic_info.get("specs"), # 规格
- "status": product_info.get("status"), # 商品状态
- "listing_beg_time": product_info.get("listing_beg_time"), # 上架时间(详情直接给文本)
- "listing_end_time": product_info.get("listing_end_time"), # 下架时间(详情文本)
- # ---- 直播回放字段(来自 live/user_live/detail;FETCH_REPLAY 关闭时为空)----
- "live_room_id": detail.get("live_room_id"), # 直播间 id(补采回放用,勿删)
- "replay_url": replay_url, # 回放视频地址(offline_address_info[0].url,多个只取下标0)
- "live_start_time": live_detail.get("start_time"), # 开播时间
- "live_end_time": live_detail.get("end_time"), # 结束时间
- }
- return row
- def filter_new_pids(shop_id, page_pids: list, sql_pool) -> set:
- """在「本页」的 pid 里查出哪些是新的(库中不存在),只查这一页、不整表拉取。
- 为什么按页查而非一次性把该商户全部 pid 拉成集合:某些商户历史成交可能几十万条,
- 整表拉取会 SELECT 出几十万行并在内存堆一个几十 MB 的集合,每轮每店都来一遍不划算。
- 改为每页用 `pid IN (...)` 批量查已存在的(凭 pid 唯一索引很快),内存只占一页的量,
- 商户体量再大也不受影响;增量本就只翻头几页(连续旧页即早停),每店只几次小查询。
- 历史成交列表顺序不稳定、新商品会跨页穿插,所以判重用「集合」而非「单个停止线精确匹配」,
- 再配合调用方的「连续 STOP_AFTER_DUPE_PAGES 页全旧」早停。
- Args:
- shop_id (int | str): 商户 id。
- page_pids (list): 本页商品 pid 列表(原始类型,可能为 int)。
- sql_pool (MySQLConnectionPool): MySQL 连接池。
- Returns:
- set[int]: 本页中「库里还没有」的 pid 集合;本页为空时返回空集合。
- pid 库内与 API 均为 int,直接按 int 比较(索引最优)。
- """
- pids = [int(p) for p in page_pids if p is not None]
- if not pids:
- return set()
- placeholders = ",".join(["%s"] * len(pids))
- rows = sql_pool.select_all(
- f"SELECT pid FROM cgj_product_record WHERE shop_id = %s AND pid IN ({placeholders})",
- (shop_id, *pids),
- )
- existing = {r[0] for r in rows} if rows else set()
- return {p for p in pids if p not in existing}
- def get_sold_list(log, shop_id, shop_name: str, sql_pool, incremental: bool = True) -> int:
- """遍历某商户历史成交翻页,逐商品补详情后写入 cgj_product_record。
- Args:
- log: 日志对象。
- shop_id (int | str): 商户 id。
- shop_name (str): 商户名。
- sql_pool (MySQLConnectionPool): MySQL 连接池。
- incremental (bool, optional): True(daily) 碰到最新已采 pid 早停,只采新商品;
- False(history) 全量深翻所有页。Defaults to True。
- Returns:
- int: 本商户写入的商品数。
- """
- cursor = ""
- page = 0
- saved = 0
- dupe_pages = 0 # 连续「全是已入库商品(0 新)」的页数,达到阈值即增量早停
- while page < MAX_PAGES:
- try:
- resp = get_group_buying_page(log, cursor, merchant_id=shop_id)
- # print(resp)
- except Exception as e:
- log.error(f"商户 {shop_id} 历史成交 cursor={cursor!r} 请求失败: {e}")
- break
- if not resp or resp.get("code") != 0:
- log.info(f"商户 {shop_id} 历史成交返回异常: {resp.get('msg') if resp else None}")
- break
- data = resp.get("data") or {}
- items = data.get("list") or []
- if not items:
- log.info(f"商户 {shop_id} 历史成交第 {page + 1} 页无数据,停止翻页")
- break
- # 增量判重:只查本页 pid 里哪些是新的(不整表拉取,适配商户体量增长)
- new_pids = filter_new_pids(shop_id, [it.get("id") for it in items], sql_pool) if incremental else None
- batch = []
- new_on_page = 0 # 本页「新商品」数(用于连续全旧页早停判断)
- for item in items:
- pid = item.get("id")
- # 已入库的商品直接跳过,连详情/回放都不再请求(省大量无用请求)
- if incremental and int(pid) not in new_pids:
- continue
- new_on_page += 1
- detail = {}
- if FETCH_DETAIL:
- try:
- detail = get_product_detail(log, pid)
- except Exception as e:
- log.error(f"商品 {pid} 详情请求失败: {e}")
- # 直播回放:room_id 取自详情的 live_room_id,再单独请求 live/user_live/detail
- live_detail = {}
- if FETCH_REPLAY and detail:
- room_id = detail.get("live_room_id")
- if room_id:
- try:
- live_detail = get_live_detail(log, room_id)
- except Exception as e:
- log.error(f"商品 {pid} 回放详情请求失败: {e}")
- row = build_product(item, detail, shop_id, shop_name, live_detail)
- if row:
- batch.append(row)
- if batch:
- # 已存在(pid 唯一)则跳过:历史成交为终态,无需覆盖
- sql_pool.insert_many(table="cgj_product_record", data_list=batch, ignore=True)
- saved += len(batch)
- log.info(f"商户 {shop_id} 历史成交第 {page + 1} 页完成,本页 {len(items)} 商品,新增 {new_on_page} 个")
- # 增量早停:连续 STOP_AFTER_DUPE_PAGES 页都没有新商品 → 已翻到旧数据区,停止翻页
- if incremental:
- if new_on_page == 0:
- dupe_pages += 1
- if dupe_pages >= STOP_AFTER_DUPE_PAGES:
- log.info(f"商户 {shop_id} 连续 {dupe_pages} 页无新商品,增量早停")
- break
- else:
- dupe_pages = 0
- if data.get("done") or not data.get("cursor"):
- break
- cursor = data["cursor"]
- page += 1
- return saved
- # ==================== 三、玩家采集(买家名单) ====================
- def get_player_page(log, product_id, cursor: str = "") -> dict | None:
- """获取某商品玩家(买家)名单的一页 announcement_by_user。
- Args:
- log: 日志对象。
- product_id (int | str): 商品 id。
- cursor (str, optional): 翻页游标,首页传空串。Defaults to ""。
- Returns:
- dict | None: 响应 JSON(含 data.list / cursor / done);失败由 cgj_request 抛异常触发重试。
- """
- return cgj_request(
- log, "/api/product/announcement_by_user", method="GET",
- params={"id": product_id, "cursor": cursor, "size": PLAYER_PAGE_SIZE},
- )
- def save_players(product_id, items: list, sql_pool) -> int:
- """解析玩家(买家)名单并写入 cgj_player_record。
- ------------------------------------------------------------------
- 并同步 schema.sql 的 cgj_player_record 列。
- ------------------------------------------------------------------
- Args:
- product_id (int | str): 商品 id(写入 pid 列)。
- items (list[dict]): announcement_by_user 的 data.list,每项含 view_key/total/buyer_name/buyer_image_url。
- sql_pool (MySQLConnectionPool): MySQL 连接池。
- Returns:
- int: 本次写入的玩家记录数。
- """
- info_list = []
- for item in items:
- # print(item)
- info_list.append({
- "pid": product_id, # 商品 id(管道必需,item 里没有)
- "view_key": item.get("view_key"), # 唯一标识
- "total": item.get("total"), # 份数
- "buyer_name": item.get("buyer_name") # 买家昵称(脱敏)
- })
- # print(info_list)
- if info_list:
- sql_pool.insert_many(table="cgj_player_record", data_list=info_list, ignore=True)
- return len(info_list)
- def get_player_list(log, product_id, sql_pool) -> bool:
- """遍历某商品玩家名单翻页并写库。
- Args:
- log: 日志对象。
- product_id (int | str): 商品 id。
- sql_pool (MySQLConnectionPool): MySQL 连接池。
- Returns:
- bool: True 表示抓到玩家数据,False 表示无数据。
- """
- cursor = ""
- page = 0
- has_data = False
- while page < MAX_PAGES:
- try:
- resp = get_player_page(log, product_id, cursor)
- except Exception as e:
- log.error(f"商品 {product_id} 玩家名单 cursor={cursor!r} 请求失败: {e}")
- break
- if not resp or resp.get("code") != 0:
- log.info(f"商品 {product_id} 暂无玩家: {resp.get('msg') if resp else None}")
- break
- data = resp.get("data") or {}
- items = data.get("list") or []
- if not items:
- log.info(f"商品 {product_id} 玩家名单翻页完成")
- break
- has_data = True
- save_players(product_id, items, sql_pool)
- if data.get("done") or not data.get("cursor"):
- log.info(f"商品 {product_id} 玩家名单翻页完成")
- break
- cursor = data["cursor"]
- page += 1
- return has_data
- # ==================== 四、回放补采 ====================
- def refill_replays(log, sql_pool) -> int:
- """补采回放地址:对已采玩家(player_state=1)但 replay_url 仍空的商品,
- 重取 live_room_id → live/user_live/detail → offline_address_info[0].url → 更新 replay_url。
- 覆盖场景:首次入库时直播还没结束 / 回放 VOD 还没生成,offline_address_info 为空、
- replay_url 拿不到;直播结束、回放就绪后由本函数补回。库里已存 live_room_id 则直接用,
- 没有(如首采详情失败)则回退重取商品详情拿 room_id。replay_url 仍空则留到下一轮再试。
- Args:
- log: 日志对象。
- sql_pool (MySQLConnectionPool): MySQL 连接池。
- Returns:
- int: 本轮成功补采的回放数。
- """
- rows = sql_pool.select_all(
- "SELECT pid, live_room_id FROM cgj_product_record WHERE replay_url IS NULL AND player_state = 1"
- )
- rows = rows or []
- log.info(f"待补回放商品 {len(rows)} 个")
- filled = 0
- for pid, room_id in rows:
- try:
- # 库里没存 room_id(首采详情失败等)→ 回退重取详情
- if not room_id:
- detail = get_product_detail(log, pid)
- room_id = detail.get("live_room_id")
- if not room_id:
- continue
- live_detail = get_live_detail(log, room_id)
- offline_list = live_detail.get("offline_address_info") or []
- url = offline_list[0].get("url") if offline_list else None
- live_start_time = live_detail.get("start_time") # 开播时间
- live_end_time = live_detail.get("end_time") # 结束时间
- if url:
- sql_pool.update_one(
- "UPDATE cgj_product_record SET replay_url = %s, live_start_time = %s, live_end_time = %s WHERE pid = %s",
- (url, live_start_time, live_end_time, 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 cgj_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 cgj_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 cgj_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 表示开始采集
- sql_pool.update_one("UPDATE cgj_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 cgj_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 cgj_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)但 replay_url 仍空的商品重取回放
- try:
- n = refill_replays(log, sql_pool)
- log.info(f"回放补采完成,本轮 {n} 条")
- except Exception as e:
- log.error(f"refill_replays 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 运行一次 cgj_main。"""
- # 立即运行一次(调试时取消注释)
- # cgj_main(log=logger)
- schedule.every().day.at("00:01").do(cgj_main, log=logger)
- while True:
- schedule.run_pending()
- time.sleep(1)
- if __name__ == "__main__":
- # cgj_main(logger)
- schedule_task()
|