feat: 完成一尘数据采集模块

- crawlers/base.py: 爬虫基类(请求封装/延时/重试/COOKIE管理)
- crawlers/yichens_spider.py: 一尘帖子爬虫(分类解析/时间解析/数据入库)
- YichensPostSpider: 帖子采集(支持多板块/分页/去重)
- YichensUserSpider: 用户采集(支持详情爬取)
- API触发: POST /api/v1/crawl-jobs/trigger?source=yichens
This commit is contained in:
caibotmi 2026-04-04 18:54:15 +08:00
parent e86a71db23
commit 63a5b46791
3 changed files with 463 additions and 109 deletions

View File

@ -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,

160
crawlers/base.py Normal file
View File

@ -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}"

View File

@ -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)