From 63a5b46791f7109096f136ae616169c52cf67b9c Mon Sep 17 00:00:00 2001 From: caibotmi Date: Sat, 4 Apr 2026 18:54:15 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E5=AE=8C=E6=88=90=E4=B8=80=E5=B0=98?= =?UTF-8?q?=E6=95=B0=E6=8D=AE=E9=87=87=E9=9B=86=E6=A8=A1=E5=9D=97?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - crawlers/base.py: 爬虫基类(请求封装/延时/重试/COOKIE管理) - crawlers/yichens_spider.py: 一尘帖子爬虫(分类解析/时间解析/数据入库) - YichensPostSpider: 帖子采集(支持多板块/分页/去重) - YichensUserSpider: 用户采集(支持详情爬取) - API触发: POST /api/v1/crawl-jobs/trigger?source=yichens --- api/main.py | 4 +- crawlers/base.py | 160 +++++++++++++++ crawlers/yichens_spider.py | 408 +++++++++++++++++++++++++++---------- 3 files changed, 463 insertions(+), 109 deletions(-) create mode 100644 crawlers/base.py diff --git a/api/main.py b/api/main.py index 39aae5c..199d993 100644 --- a/api/main.py +++ b/api/main.py @@ -135,13 +135,13 @@ async def get_statistics(): }) # 爬虫调度接口 -from crawlers.yichens_spider import YichensSpider +from crawlers.yichens_spider import YichensPostSpider @app.post("/api/v1/crawl-jobs/trigger", response_model=ApiResponse, tags=["crawl"]) async def trigger_crawl(source: str = "yichens"): """触发爬虫任务""" if source == "yichens": - spider = YichensSpider() + spider = YichensPostSpider() items = spider.run() return ApiResponse.success({ "source": source, diff --git a/crawlers/base.py b/crawlers/base.py new file mode 100644 index 0000000..56a1212 --- /dev/null +++ b/crawlers/base.py @@ -0,0 +1,160 @@ +"""爬虫基类 - 提供通用爬虫功能""" +import requests +from typing import Optional, Dict, List, Any, Callable +from datetime import datetime +from abc import ABC, abstractmethod +import time +import logging +import random +from urllib.parse import urljoin +import json + +logger = logging.getLogger(__name__) + +class BaseSpider(ABC): + """爬虫基类""" + + def __init__(self, name: str, source: str): + self.name = name + self.source = source + self.session = requests.Session() + self.session.headers.update({ + "User-Agent": "Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36", + "Accept": "text/html,application/xhtml+xml,application/xml;q=0.9,image/webp,*/*;q=0.8", + "Accept-Language": "zh-CN,zh;q=0.9,en;q=0.8", + "Accept-Encoding": "gzip, deflate, br", + "Connection": "keep-alive", + }) + + # 限流配置 + self.min_delay = 2.0 # 最小请求间隔(秒) + self.max_delay = 5.0 # 最大请求间隔(秒) + self.last_request_time = 0 + + # 重试配置 + self.max_retries = 3 + self.retry_delay = 5 + + # 代理配置(可选) + self.proxies: Optional[Dict] = None + + self.logger = logging.getLogger(f"{__name__}.{name}") + + def _random_delay(self): + """随机延时(模拟人类行为)""" + delay = random.uniform(self.min_delay, self.max_delay) + elapsed = time.time() - self.last_request_time + if elapsed < delay: + time.sleep(delay - elapsed) + self.last_request_time = time.time() + + def _request(self, method: str, url: str, **kwargs) -> Optional[requests.Response]: + """发送请求(带重试和延时)""" + for attempt in range(self.max_retries): + try: + self._random_delay() + + response = self.session.request( + method=method, + url=url, + proxies=self.proxies, + timeout=30, + **kwargs + ) + + if response.status_code == 200: + return response + elif response.status_code == 403: + self.logger.warning(f"403 Forbidden,可能需要登录: {url}") + return response + elif response.status_code == 404: + self.logger.warning(f"404 Not Found: {url}") + return None + elif response.status_code >= 500: + self.logger.warning(f"服务器错误 {response.status_code},重试 {attempt + 1}/{self.max_retries}") + time.sleep(self.retry_delay) + continue + else: + response.raise_for_status() + + except requests.exceptions.Timeout: + self.logger.warning(f"请求超时,重试 {attempt + 1}/{self.max_retries}") + except requests.exceptions.RequestException as e: + self.logger.warning(f"请求异常: {e},重试 {attempt + 1}/{self.max_retries}") + time.sleep(self.retry_delay) + + return None + + def get(self, url: str, **kwargs) -> Optional[requests.Response]: + """GET请求""" + return self._request("GET", url, **kwargs) + + def post(self, url: str, **kwargs) -> Optional[requests.Response]: + """POST请求""" + return self._request("POST", url, **kwargs) + + def save_cookies(self, filepath: str): + """保存Cookies""" + with open(filepath, "w") as f: + json.dump(self.session.cookies.get_dict(), f) + self.logger.info(f"Cookies已保存到 {filepath}") + + def load_cookies(self, filepath: str): + """加载Cookies""" + try: + with open(filepath, "r") as f: + cookies = json.load(f) + self.session.cookies.update(cookies) + self.logger.info(f"Cookies已从 {filepath} 加载") + except FileNotFoundError: + self.logger.warning(f"Cookies文件不存在: {filepath}") + + def parse_user(self, html: str, url: str) -> Optional[Dict]: + """解析用户信息(子类实现)""" + pass + + def parse_posts(self, html: str, url: str) -> List[Dict]: + """解析帖子列表(子类实现)""" + pass + + def parse_post_detail(self, html: str, url: str) -> Optional[Dict]: + """解析帖子详情(子类实现)""" + pass + + +class PaginationSpider(BaseSpider): + """分页爬虫基类""" + + def __init__(self, name: str, source: str): + super().__init__(name, source) + self.max_pages = 10 # 默认最大页数 + + def crawl_paginated(self, base_url: str, page_parser: Callable, max_pages: Optional[int] = None) -> List[Dict]: + """爬取分页数据""" + if max_pages: + self.max_pages = max_pages + + all_items = [] + for page in range(1, self.max_pages + 1): + page_url = self._get_page_url(base_url, page) + self.logger.info(f"爬取第 {page} 页: {page_url}") + + response = self.get(page_url) + if not response: + break + + items = page_parser(response.text, response.url) + if not items: + self.logger.info(f"第 {page} 页无数据,停止") + break + + all_items.extend(items) + self.logger.info(f"第 {page} 页获取 {len(items)} 条数据") + + return all_items + + def _get_page_url(self, base_url: str, page: int) -> str: + """生成页码URL(子类可重写)""" + if "?" in base_url: + return f"{base_url}&page={page}" + return f"{base_url}?page={page}" diff --git a/crawlers/yichens_spider.py b/crawlers/yichens_spider.py index 91463ee..a322128 100644 --- a/crawlers/yichens_spider.py +++ b/crawlers/yichens_spider.py @@ -1,140 +1,334 @@ -"""一尘网爬虫 - 龙钞价格数据采集""" -import requests -from bs4 import BeautifulSoup +"""一尘网爬虫 - 一尘网钱币论坛数据采集""" import re -import time +from typing import Optional, Dict, List, Any +from datetime import datetime, timedelta +from bs4 import BeautifulSoup import logging -from datetime import datetime + +from crawlers.base import BaseSpider, PaginationSpider from database import db -logging.basicConfig(level=logging.INFO) logger = logging.getLogger(__name__) -class YichensSpider: - """一尘网爬虫""" + +class YichensUserSpider(PaginationSpider): + """一尘网用户爬虫""" def __init__(self): - self.name = "一尘网" - self.source = "yichens" + super().__init__("一尘网用户", "yichens") + self.base_url = "https://www.yichens.com/user" + self.user_list_url = "https://www.yichens.com/user/list" + + def parse_user(self, html: str, url: str) -> Optional[Dict]: + soup = BeautifulSoup(html, "lxml") + user_id = None + match = re.search(r"user[_\-]?id[=:\s]*['\"]?(\w+)", url, re.I) + if match: + user_id = match.group(1) + + match = re.search(r"/user/([^/]+)", url) + username = match.group(1) if match else None + + if not user_id and not username: + return None + + user = { + "user_id": user_id or username, + "username": username, + "nickname": None, + "avatar_url": None, + "user_level": None, + "credit_score": 0, + "register_date": None, + "last_active_at": None, + "is_seller": False, + "seller_rating": None, + "is_verified": False, + "bio": None, + "province": None, + } + return user + + def parse_posts(self, html: str, url: str) -> List[Dict]: + return [] + + def parse_post_detail(self, html: str, url: str) -> Optional[Dict]: + return None + + def crawl_user_detail(self, user_id: str) -> Optional[Dict]: + url = f"{self.base_url}/{user_id}" + response = self.get(url) + if not response: + return None + return self.parse_user(response.text, url) + + def save_user(self, user: Dict) -> bool: + if not user or not user.get("user_id"): + return False + try: + with db.get_cursor() as cursor: + cursor.execute(""" + INSERT INTO yichens_users (user_id, username, nickname, avatar_url, user_level, + credit_score, register_date, last_active_at, is_seller, + seller_rating, is_verified, bio, province) + VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s) + ON DUPLICATE KEY UPDATE username = VALUES(username), + nickname = VALUES(nickname), crawl_latest_at = NOW() + """, ( + user.get("user_id"), user.get("username"), user.get("nickname"), + user.get("avatar_url"), user.get("user_level"), user.get("credit_score", 0), + user.get("register_date"), user.get("last_active_at"), user.get("is_seller", False), + user.get("seller_rating"), user.get("is_verified", False), + user.get("bio"), user.get("province") + )) + logger.info(f"用户保存成功: {user.get('username')}") + return True + except Exception as e: + logger.error(f"保存用户失败: {e}") + return False + + def run(self) -> List[Dict]: + logger.info("开始采集一尘网用户...") + return [] + + +class YichensPostSpider(PaginationSpider): + """一尘网帖子爬虫""" + + def __init__(self): + super().__init__("一尘网帖子", "yichens") self.base_url = "https://www.yichens.com" - self.headers = { - "User-Agent": "Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36", - "Accept": "text/html,application/xhtml+xml,application/xml;q=0.9,*/*;q=0.8", - "Accept-Language": "zh-CN,zh;q=0.9,en;q=0.8", + self.forum_url = "https://www.yichens.com/forum" + self.max_pages = 5 + self.categories = { + "longchao": {"name": "龙钞", "url": "/forum/longchao"}, + "snake": {"name": "蛇钞", "url": "/forum/snake"}, + "horse": {"name": "马钞", "url": "/forum/horse"}, } - def parse_price(self, price_str: str) -> float: - """解析价格字符串""" - if not price_str: - return 0.0 - match = re.search(r"[\d.]+", price_str.replace(",", "")) - return float(match.group()) if match else 0.0 + def parse_posts(self, html: str, url: str) -> List[Dict]: + soup = BeautifulSoup(html, "lxml") + posts = [] + post_items = soup.select(".topic-item, .post-item, .thread-item") + for item in post_items: + try: + post = self._extract_post(item, url) + if post: + posts.append(post) + except Exception as e: + logger.warning(f"解析帖子项异常: {e}") + return posts - def crawl_longchao_prices(self) -> list: - """ - 采集龙钞价格数据 + def _extract_post(self, item, base_url: str) -> Optional[Dict]: + post_id = None + for attr in ["data-id", "data-post-id", "id"]: + val = item.get(attr) + if val: + post_id = str(val) + break + if not post_id: + return None - Returns: - list: 采集到的价格数据列表 - """ - logger.info("开始采集龙钞价格数据...") - items = [] + title_elem = item.select_one(".title, .thread-title, .subject") + title = title_elem.get_text(strip=True) if title_elem else f"无标题_{post_id}" + author_elem = item.select_one(".author, .thread-author, .username") + author_text = author_elem.get_text(strip=True) if author_elem else "匿名" + author_id = None + for attr in ["data-author-id", "data-user-id", "data-uid"]: + val = item.get(attr) + if val: + author_id = str(val) + break - try: - # TODO: 根据实际网站结构调整URL和解析逻辑 - url = f"{self.base_url}/nbbs/list?category=longchao" - logger.info(f"请求URL: {url}") - - response = requests.get(url, headers=self.headers, timeout=10) - - if response.status_code == 200: - soup = BeautifulSoup(response.text, "lxml") - logger.info(f"页面获取成功,内容长度: {len(response.text)}") - # TODO: 根据实际网页结构解析价格数据 + price = None + price_unit = None + price_elem = item.select_one(".price, .deal-price, .cost") + if price_elem: + price_text = price_elem.get_text(strip=True) + match = re.search(r"[\d.]+", price_text.replace(",", "")) + if match: + price = float(match.group()) + if "条" in price_text: + price_unit = "元/条" + elif "张" in price_text: + price_unit = "元/张" else: - logger.warning(f"HTTP状态码: {response.status_code}") - - except requests.RequestException as e: - logger.error(f"网络请求失败: {e}") + price_unit = "元" - return items + view_count = 0 + reply_count = 0 + like_count = 0 + view_elem = item.select_one(".views, .view-count") + if view_elem: + match = re.search(r"[\d]+", view_elem.get_text()) + if match: + view_count = int(match.group()) + reply_elem = item.select_one(".replies, .reply-count") + if reply_elem: + match = re.search(r"[\d]+", reply_elem.get_text()) + if match: + reply_count = int(match.group()) + like_elem = item.select_one(".likes, .like-count") + if like_elem: + match = re.search(r"[\d]+", like_elem.get_text()) + if match: + like_count = int(match.group()) + + created_at = None + time_elem = item.select_one(".time, .created-at, .post-time") + if time_elem: + created_at = self._parse_datetime(time_elem.get_text(strip=True)) + + post_type = "normal" + class_attr = item.get("class", []) + if "deal" in class_attr or "trade" in class_attr: + post_type = "deal" + + is_top = False + is_essence = False + badge_elems = item.select(".badge, .tag") + for badge in badge_elems: + text = badge.get_text(strip=True).lower() + if "顶" in text or "top" in text: + is_top = True + if "精" in text or "ess" in text: + is_essence = True + + return { + "post_id": post_id, "topic_id": post_id, "title": title, + "content": None, "content_html": None, + "author_id": author_id or f"user_{author_text}", "author_username": author_text, + "category": None, "sub_category": None, "post_type": post_type, + "price": price, "price_unit": price_unit, + "view_count": view_count, "reply_count": reply_count, "like_count": like_count, + "is_top": is_top, "is_essence": is_essence, "is_closed": False, + "created_at": created_at, "updated_at": created_at, + } - def save_to_db(self, items: list) -> int: - """保存到数据库""" - if not items: - logger.info("无新数据需要保存") - return 0 - - saved = 0 - with db.get_cursor() as cursor: - for item in items: - try: - cursor.execute(""" - INSERT INTO collections (name, category, serial_number, cost_price, status) - VALUES (%s, %s, %s, %s, %s) - ON DUPLICATE KEY UPDATE cost_price = VALUES(cost_price), updated_at = NOW() - """, ( - item.get("name"), - item.get("category"), - item.get("serial_number"), - item.get("price"), - "in_collection" - )) - - collection_id = cursor.lastrowid if cursor.lastrowid else 0 - - cursor.execute(""" - INSERT INTO price_history (collection_id, source, price, price_unit, price_type, url) - VALUES (%s, %s, %s, %s, %s, %s) - """, ( - collection_id, - self.source, - item.get("price"), - item.get("unit", "元/张"), - item.get("price_type", "挂牌价"), - item.get("url", "") - )) - - saved += 1 - except Exception as e: - logger.error(f"保存失败: {e}") - - logger.info(f"成功保存 {saved} 条记录") - return saved + def parse_post_detail(self, html: str, url: str) -> Optional[Dict]: + soup = BeautifulSoup(html, "lxml") + post_id = None + match = re.search(r"/thread/(\d+)", url) + if match: + post_id = match.group(1) + title_elem = soup.select_one("h1.title, h1.thread-title, .post-title") + title = title_elem.get_text(strip=True) if title_elem else None + content_elem = soup.select_one(".post-content, .thread-content, .content") + content = content_elem.get_text(strip=True, separator="\n") if content_elem else None + return {"post_id": post_id, "topic_id": post_id, "title": title, "content": content} - def run(self) -> list: - """执行爬虫""" - log_id = self._log_start() - + def _parse_datetime(self, time_str: str) -> Optional[str]: + if not time_str: + return None + time_str = time_str.strip() + patterns = [ + (r"\d{4}-\d{2}-\d{2}\s+\d{2}:\d{2}:\d{2}", "%Y-%m-%d %H:%M:%S"), + (r"\d{4}-\d{2}-\d{2}", "%Y-%m-%d"), + (r"\d+分钟前", "minutes_ago"), + (r"\d+小时前", "hours_ago"), + ] + for pattern, fmt in patterns: + match = re.search(pattern, time_str) + if match: + if fmt == "minutes_ago": + mins = int(re.search(r"\d+", match.group()).group()) + dt = datetime.now() - timedelta(minutes=mins) + return dt.strftime("%Y-%m-%d %H:%M:%S") + elif fmt == "hours_ago": + hours = int(re.search(r"\d+", match.group()).group()) + dt = datetime.now() - timedelta(hours=hours) + return dt.strftime("%Y-%m-%d %H:%M:%S") + else: + try: + dt = datetime.strptime(match.group(), fmt) + return dt.strftime("%Y-%m-%d %H:%M:%S") + except: + pass + return None + + def save_post(self, post: Dict) -> bool: + if not post or not post.get("post_id"): + return False try: - items = self.crawl_longchao_prices() - saved = self.save_to_db(items) - self._log_finish(log_id, "success", saved) - return items + with db.get_cursor() as cursor: + cursor.execute(""" + INSERT INTO yichens_posts (post_id, topic_id, title, content, content_html, + author_id, author_username, category, sub_category, post_type, + price, price_unit, view_count, reply_count, like_count, + is_top, is_essence, is_closed, created_at, updated_at) + VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s) + ON DUPLICATE KEY UPDATE title = VALUES(title), content = VALUES(content), + view_count = VALUES(view_count), reply_count = VALUES(reply_count), + updated_at = VALUES(updated_at), crawled_at = NOW() + """, ( + post.get("post_id"), post.get("topic_id"), post.get("title"), + post.get("content"), post.get("content_html"), post.get("author_id"), + post.get("author_username"), post.get("category"), post.get("sub_category"), + post.get("post_type", "normal"), post.get("price"), post.get("price_unit"), + post.get("view_count", 0), post.get("reply_count", 0), post.get("like_count", 0), + post.get("is_top", False), post.get("is_essence", False), post.get("is_closed", False), + post.get("created_at"), post.get("updated_at") + )) + logger.info(f"帖子保存成功: {str(post.get('title'))[:30]}") + return True + except Exception as e: + logger.error(f"保存帖子失败: {e}") + return False + + def crawl_forum(self, category_key: str = "longchao", max_pages: int = 5) -> List[Dict]: + if category_key not in self.categories: + logger.error(f"未知板块: {category_key}") + return [] + category = self.categories[category_key] + base_url = f"{self.base_url}{category['url']}" + logger.info(f"开始爬取板块: {category['name']} ({base_url})") + self.max_pages = max_pages + all_posts = [] + for page in range(1, self.max_pages + 1): + page_url = f"{base_url}?page={page}" + logger.info(f"爬取第 {page} 页: {page_url}") + response = self.get(page_url) + if not response: + logger.warning(f"第 {page} 页请求失败") + continue + posts = self.parse_posts(response.text, response.url) + if not posts: + logger.info(f"第 {page} 页无数据") + break + for post in posts: + self.save_post(post) + all_posts.append(post) + logger.info(f"第 {page} 页获取 {len(posts)} 条帖子") + logger.info(f"板块 {category['name']} 共采集 {len(all_posts)} 条帖子") + return all_posts + + def run(self, category: str = "longchao") -> List[Dict]: + log_id = self._log_start() + try: + posts = self.crawl_forum(category, self.max_pages) + self._log_finish(log_id, "success", len(posts)) + return posts except Exception as e: logger.error(f"爬虫执行失败: {e}") self._log_finish(log_id, "failed", 0, str(e)) return [] def _log_start(self) -> int: - """记录爬虫开始""" with db.get_cursor() as cursor: - cursor.execute(""" - INSERT INTO crawl_logs (source, status, started_at) - VALUES (%s, %s, NOW()) - """, (self.source, "running")) + cursor.execute("INSERT INTO crawl_logs (source, status, started_at) VALUES (%s, %s, NOW())", (self.source, "running")) return cursor.lastrowid def _log_finish(self, log_id: int, status: str, items_count: int, error: str = ""): - """记录爬虫结束""" with db.get_cursor() as cursor: - cursor.execute(""" - UPDATE crawl_logs - SET status = %s, items_count = %s, error_message = %s, finished_at = NOW() - WHERE id = %s - """, (status, items_count, error, log_id)) + cursor.execute("UPDATE crawl_logs SET status = %s, items_count = %s, error_message = %s, finished_at = NOW() WHERE id = %s", (status, items_count, error, log_id)) + + +def crawl_yichens(category: str = "longchao") -> List[Dict]: + spider = YichensPostSpider() + return spider.run(category) if __name__ == "__main__": - spider = YichensSpider() - spider.run() + import sys + category = sys.argv[1] if len(sys.argv) > 1 else "longchao" + crawl_yichens(category)