crawler.py 9.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260
  1. """帖子详情拉取:多平台 detail 接口(小红书/抖音/快手/B站)。
  2. crawler.aiddit.com 各平台 detail 返回字段已归一化(channel_content_id/title/content_type/
  3. body_text/image_url_list/video_url_list/channel_account_*),所以 parse 一套复用;差异只在
  4. detail path、URL→content_id 解析、平台名/id 前缀;正式环境 host/超时从 .env 读取。
  5. 拆成两层:parse_detail_response 纯函数(可离线用 fixture 测)+ fetch_post_detail 负责 HTTP。
  6. """
  7. from __future__ import annotations
  8. import random
  9. import re
  10. import time
  11. from typing import Any, Callable, Optional
  12. from urllib.parse import urljoin
  13. import httpx
  14. from core.config import Settings
  15. from core.models import Card, Post
  16. RATE_LIMIT_SECONDS = 15.0
  17. CRAWLER_MIN_INTERVAL_SECONDS = 10.0
  18. CRAWLER_MAX_INTERVAL_SECONDS = 12.0
  19. # 平台路由:detail path + Post.id 前缀
  20. PLATFORMS = {
  21. "xiaohongshu": {"path": "/crawler/xiao_hong_shu/detail", "prefix": "xhs"},
  22. "douyin": {"path": "/crawler/dou_yin/detail", "prefix": "dy"},
  23. "kuaishou": {"path": "/crawler/kuai_shou/detail", "prefix": "ks"},
  24. "bilibili": {"path": "/crawler/bilibili/detail", "prefix": "bili"},
  25. }
  26. DOUYIN_PROVIDERS = {"piaoquantv", "aiddit"}
  27. def _last_seg(s: str) -> str:
  28. return s.split("?", 1)[0].rstrip("/").split("/")[-1]
  29. def detect_platform_and_id(content_id_or_url: str) -> tuple[str, str]:
  30. """从 URL(或裸 id)识别平台 + content_id。不支持短链(v.douyin.com 等需先跟随重定向)。"""
  31. s = content_id_or_url.strip()
  32. low = s.lower()
  33. if "xiaohongshu.com" in low or "/explore/" in low or "/discovery/item/" in low:
  34. return "xiaohongshu", parse_content_id(s)
  35. if "douyin.com" in low:
  36. m = (re.search(r"modal_id=(\d+)", s) or re.search(r"/video/(\d+)", s)
  37. or re.search(r"/note/(\d+)", s) or re.search(r"(\d{15,})", s))
  38. return "douyin", m.group(1) if m else _last_seg(s)
  39. if "kuaishou.com" in low or "gifshow.com" in low:
  40. m = re.search(r"/(?:fw/photo|short-video|photo|f)/([A-Za-z0-9_-]+)", s)
  41. return "kuaishou", m.group(1) if m else _last_seg(s)
  42. if "bilibili.com" in low or "b23.tv" in low:
  43. m = re.search(r"(BV[0-9A-Za-z]+)", s)
  44. return "bilibili", m.group(1) if m else _last_seg(s)
  45. # 裸 id 兜底:BV→B站;纯长数字→抖音;其余→小红书(向后兼容)
  46. if s.startswith("BV"):
  47. return "bilibili", s
  48. if re.fullmatch(r"\d{15,}", s):
  49. return "douyin", s
  50. return "xiaohongshu", parse_content_id(s)
  51. class CrawlerError(RuntimeError):
  52. pass
  53. class RateLimiter:
  54. """同一 bucket 两次调用间隔 ≥ min_interval。对齐 CFA RateLimiter。
  55. 给 max_interval_seconds 时,每次间隔在 [min, max] 随机取(搜索接口加随机抖动躲限流)。"""
  56. def __init__(
  57. self,
  58. min_interval_seconds: float = RATE_LIMIT_SECONDS,
  59. now_fn: Callable[[], float] = time.monotonic,
  60. sleep_fn: Callable[[float], None] = time.sleep,
  61. max_interval_seconds: Optional[float] = None,
  62. ) -> None:
  63. self.min_interval_seconds = min_interval_seconds
  64. self.max_interval_seconds = max_interval_seconds
  65. self.now_fn = now_fn
  66. self.sleep_fn = sleep_fn
  67. self._last: dict[str, float] = {}
  68. def _interval(self) -> float:
  69. if self.max_interval_seconds is not None:
  70. return random.uniform(self.min_interval_seconds, self.max_interval_seconds)
  71. return self.min_interval_seconds
  72. def wait(self, bucket: str) -> None:
  73. last = self._last.get(bucket)
  74. if last is not None:
  75. remaining = self._interval() - (self.now_fn() - last)
  76. if remaining > 0:
  77. self.sleep_fn(remaining)
  78. self._last[bucket] = self.now_fn()
  79. def parse_content_id(content_id_or_url: str) -> str:
  80. """接受裸 content_id 或 https://www.xiaohongshu.com/explore/<id>?... 链接。"""
  81. s = content_id_or_url.strip()
  82. if "/explore/" in s:
  83. s = s.split("/explore/", 1)[1]
  84. if "/discovery/item/" in content_id_or_url:
  85. s = content_id_or_url.split("/discovery/item/", 1)[1]
  86. return s.split("?", 1)[0].strip("/").strip()
  87. def _image_urls(inner: dict) -> list[str]:
  88. out: list[str] = []
  89. seen: set[str] = set()
  90. for x in inner.get("image_url_list") or []:
  91. u = x.get("image_url") if isinstance(x, dict) else x
  92. if u and u not in seen:
  93. seen.add(u)
  94. out.append(u)
  95. return out
  96. def _video_urls(inner: dict) -> list[str]:
  97. out: list[str] = []
  98. for x in inner.get("video_url_list") or []:
  99. u = x.get("video_url") if isinstance(x, dict) else x
  100. if u:
  101. out.append(u)
  102. return out
  103. def _douyin_detail_base_url(settings: Settings, provider: str) -> str:
  104. if provider == "piaoquantv":
  105. return settings.piaoquantv_douyin_base_url
  106. if provider == "aiddit":
  107. return settings.aiddit_crawler_base_url
  108. raise CrawlerError(f"unsupported douyin provider: {provider}")
  109. def _douyin_detail_body(content_id: str, *, settings: Settings, provider: str) -> dict[str, Any]:
  110. body: dict[str, Any] = {"content_id": content_id}
  111. if provider == "piaoquantv":
  112. body["account_id"] = settings.piaoquantv_douyin_account_id
  113. body["cookie_batch"] = settings.piaoquantv_douyin_cookie_batch
  114. elif provider != "aiddit":
  115. raise CrawlerError(f"unsupported douyin provider: {provider}")
  116. return body
  117. def parse_detail_response(response: dict, *, platform: str = "xiaohongshu",
  118. provider: str = "",
  119. fallback_content_id: str = "") -> Post:
  120. """把 detail 接口返回(信封 {code,msg,data:{data:{...}}})解析成 Post。字段各平台归一。"""
  121. if not isinstance(response, dict):
  122. raise CrawlerError("bad_response: not a dict")
  123. code = response.get("code")
  124. if code not in (0, "0"):
  125. raise CrawlerError(f"business_error: code={code} msg={response.get('msg')}")
  126. inner = ((response.get("data") or {}).get("data")) or {}
  127. if not inner:
  128. raise CrawlerError("empty_detail: data.data is empty")
  129. prefix = PLATFORMS.get(platform, PLATFORMS["xiaohongshu"])["prefix"]
  130. content_id = inner.get("channel_content_id") or fallback_content_id
  131. link = inner.get("content_link") or ""
  132. images = _image_urls(inner)
  133. # 图文帖:每张图一张卡片(1-based);视频帖的段卡由 extract_video 在提取时写入,覆盖封面卡
  134. cards = [Card(index=i, kind="image", url=u) for i, u in enumerate(images, start=1)]
  135. raw = dict(response)
  136. if provider:
  137. raw.setdefault("detail_provider", provider)
  138. return Post(
  139. id=f"{prefix}_{content_id}",
  140. platform=platform,
  141. provider=provider,
  142. url=link,
  143. content_id=content_id,
  144. title=inner.get("title") or "",
  145. content_type=inner.get("content_type") or "",
  146. body_text=inner.get("body_text") or "",
  147. topic_list=list(inner.get("topic_list") or []),
  148. image_urls=images,
  149. video_urls=_video_urls(inner),
  150. cards=cards,
  151. author_id=inner.get("channel_account_id"),
  152. author_name=inner.get("channel_account_name"),
  153. raw=raw,
  154. )
  155. def fetch_post_detail(
  156. content_id_or_url: str,
  157. *,
  158. platform: str | None = None,
  159. provider: str | None = None,
  160. settings: Optional[Settings] = None,
  161. http_client: Any = None,
  162. rate_limiter: Optional[RateLimiter] = None,
  163. env_file: str = ".env",
  164. ) -> Post:
  165. """真实拉取一条帖子详情,自动识别平台(小红书/抖音/快手/B站),返回 Post。"""
  166. settings = settings or Settings.from_env(env_file)
  167. detected_platform, content_id = detect_platform_and_id(content_id_or_url)
  168. platform = platform or detected_platform
  169. if not content_id:
  170. raise CrawlerError(f"cannot parse content_id from: {content_id_or_url}")
  171. cfg = PLATFORMS[platform]
  172. is_douyin = platform == "douyin"
  173. if is_douyin:
  174. provider = provider or "piaoquantv"
  175. if provider not in DOUYIN_PROVIDERS:
  176. raise CrawlerError(f"unsupported douyin provider: {provider}")
  177. rate_limiter = rate_limiter or RateLimiter(
  178. min_interval_seconds=CRAWLER_MIN_INTERVAL_SECONDS,
  179. max_interval_seconds=CRAWLER_MAX_INTERVAL_SECONDS,
  180. )
  181. rate_limiter.wait("douyin")
  182. else:
  183. provider = provider or ""
  184. rate_limiter = rate_limiter or RateLimiter(
  185. min_interval_seconds=CRAWLER_MIN_INTERVAL_SECONDS,
  186. max_interval_seconds=CRAWLER_MAX_INTERVAL_SECONDS,
  187. )
  188. rate_limiter.wait(f"{platform}_detail")
  189. owns_client = http_client is None
  190. client = http_client or httpx.Client()
  191. try:
  192. base_url = (
  193. _douyin_detail_base_url(settings, provider)
  194. if is_douyin
  195. else settings.aiddit_crawler_base_url
  196. )
  197. body = (
  198. _douyin_detail_body(content_id, settings=settings, provider=provider)
  199. if is_douyin
  200. else {"content_id": content_id}
  201. )
  202. url = urljoin(base_url, cfg["path"])
  203. resp = client.post(
  204. url,
  205. json=body,
  206. headers={"Content-Type": "application/json"},
  207. timeout=settings.crawler_timeout,
  208. )
  209. resp.raise_for_status()
  210. data = resp.json()
  211. except httpx.HTTPError as exc:
  212. raise CrawlerError(f"http_error: {exc}") from exc
  213. except ValueError as exc:
  214. raise CrawlerError("bad_json") from exc
  215. finally:
  216. if owns_client:
  217. client.close()
  218. return parse_detail_response(
  219. data,
  220. platform=platform,
  221. provider=provider or "",
  222. fallback_content_id=content_id,
  223. )