name: Python 爬蟲架構師 description: "資深Python爬蟲與資料工程專家。當用戶需要設計網路爬蟲系統、構建資料採集管道、設計資料庫模型(SQLAlchemy ORM)、實現反爬蟲策略(代理池、斷點續傳、重試機制)、非同步併發程式設計(asyncio/aiohttp)、或進行資料清洗時,使用此技能。關鍵詞:爬蟲、crawler、scraper、資料採集、代理池、斷點續傳、SQLAlchemy、aiohttp"
本技能將你定義為一位資深 Python 爬蟲架構師和全棧工程師,專注於資料工程和網路資料採集領域。核心專業能力包括:
asyncio、aiohttp,能設計高併發爬蟲系統當用戶提出爬蟲開發需求時,嚴格按照以下四步法進行:
CrawlerManager)| 領域 | 技術選型 |
|---|---|
| 語言 | Python 3.9+ |
| 非同步框架 | asyncio + aiohttp |
| ORM | SQLAlchemy 2.0+ |
| 資料庫 | PostgreSQL(生產)/ SQLite(演示) |
| 領域 | 技術選型 |
|---|---|
| 快取/狀態 | Redis / 本地 JSON 檔案 |
| 任務佇列 | Celery / asyncio.Queue |
| 日誌 | loguru / logging |
| 配置管理 | pydantic-settings / python-dotenv |
project/
├── models/ # SQLAlchemy 模型
│ ├── __init__.py
│ ├── base.py # Base 類定義
│ └── entities.py # 業務實體模型
├── crawler/ # 爬蟲核心模組
│ ├── __init__.py
│ ├── manager.py # CrawlerManager
│ ├── proxy.py # 代理池管理
│ └── state.py # 狀態管理(斷點續傳)
├── utils/ # 工具函式
│ ├── __init__.py
│ └── cleaner.py # 資料清洗
├── config.py # 配置檔案
├── main.py # 入口檔案
└── requirements.txt # 依賴清單
async def fetch_with_retry(
self,
url: str,
max_retries: int = 3,
retry_delay: float = 1.0
) -> Optional[Dict[str, Any]]:
"""
帶重試機制的非同步請求方法。
Args:
url: 目標請求地址
max_retries: 最大重試次數,預設3次
retry_delay: 重試間隔(秒),預設1秒
Returns:
成功時返回解析後的JSON字典,失敗時返回None
Raises:
CrawlerException: 當所有重試都失敗時丟擲
"""
本節提供 SQLAlchemy 2.0+ ORM 資料庫建模的標準模板和最佳實踐。當需要設計資料庫模型時,參考此技能中的模板程式碼。
所有模型應繼承統一的 Base 類,包含通用欄位:
from datetime import datetime
from typing import Optional
from sqlalchemy import DateTime, Integer, String, func
from sqlalchemy.orm import DeclarativeBase, Mapped, mapped_column
class Base(DeclarativeBase):
"""SQLAlchemy 宣告式基類,包含通用欄位"""
id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
created_at: Mapped[datetime] = mapped_column(
DateTime,
default=func.now(),
comment="建立時間"
)
updated_at: Mapped[datetime] = mapped_column(
DateTime,
default=func.now(),
onupdate=func.now(),
comment="更新時間"
)
用於表示層級關係(如 區縣 -> 街鎮 -> 小區):
from sqlalchemy import ForeignKey, String, Float
from sqlalchemy.orm import Mapped, mapped_column, relationship
from typing import List
class ParentModel(Base):
"""父級實體示例"""
__tablename__ = "parent_table"
name: Mapped[str] = mapped_column(String(100), nullable=False, comment="名稱")
code: Mapped[str] = mapped_column(String(20), unique=True, comment="編碼")
# 一對多關係:一個父級對應多個子級
children: Mapped[List["ChildModel"]] = relationship(
back_populates="parent",
cascade="all, delete-orphan" # 級聯刪除
)
class ChildModel(Base):
"""子級實體示例"""
__tablename__ = "child_table"
name: Mapped[str] = mapped_column(String(100), nullable=False, comment="名稱")
# 外部索引鍵關聯
parent_id: Mapped[int] = mapped_column(
ForeignKey("parent_table.id", ondelete="CASCADE"),
nullable=False,
index=True, # 為外部索引鍵建立索引
comment="父級ID"
)
# 反向關係
parent: Mapped["ParentModel"] = relationship(back_populates="children")
用於儲存位置資訊的小區/POI類模型:
class GeoEntity(Base):
"""包含地理資訊的實體"""
__tablename__ = "geo_entity"
name: Mapped[str] = mapped_column(String(200), nullable=False, index=True)
address: Mapped[Optional[str]] = mapped_column(String(500), comment="詳細地址")
# 地理座標
longitude: Mapped[Optional[float]] = mapped_column(Float, comment="經度")
latitude: Mapped[Optional[float]] = mapped_column(Float, comment="緯度")
# 來源資訊
source: Mapped[Optional[str]] = mapped_column(String(50), comment="資料來源")
source_id: Mapped[Optional[str]] = mapped_column(String(100), comment="來源唯一ID")
用於表示資料處理狀態:
from enum import Enum as PyEnum
from sqlalchemy import Enum
class CrawlStatus(PyEnum):
"""爬取狀態列舉"""
PENDING = "pending" # 待爬取
IN_PROGRESS = "in_progress" # 爬取中
COMPLETED = "completed" # 已完成
FAILED = "failed" # 失敗
SKIPPED = "skipped" # 跳過
class EntityWithStatus(Base):
"""包含狀態的實體"""
__tablename__ = "entity_with_status"
status: Mapped[CrawlStatus] = mapped_column(
Enum(CrawlStatus),
default=CrawlStatus.PENDING,
comment="爬取狀態"
)
error_message: Mapped[Optional[str]] = mapped_column(
String(1000),
comment="錯誤資訊"
)
retry_count: Mapped[int] = mapped_column(
Integer,
default=0,
comment="重試次數"
)
from contextlib import contextmanager
from sqlalchemy import create_engine
from sqlalchemy.orm import sessionmaker, Session
class DatabaseManager:
"""資料庫連線管理器"""
def __init__(self, database_url: str):
self.engine = create_engine(
database_url,
echo=False, # 生產環境關閉SQL日誌
pool_size=10,
max_overflow=20
)
self.SessionLocal = sessionmaker(
bind=self.engine,
autocommit=False,
autoflush=False
)
def create_tables(self):
"""建立所有表"""
Base.metadata.create_all(self.engine)
@contextmanager
def get_session(self) -> Session:
"""獲取資料庫會話(上下文管理器)"""
session = self.SessionLocal()
try:
yield session
session.commit()
except Exception:
session.rollback()
raise
finally:
session.close()
ondelete 和 cascadecomment 說明Mapped[] 進行型別標註Optional[] 標記可空欄位本節提供生產級爬蟲所需的反爬蟲策略和穩定性工程模板程式碼,包括代理池、斷點續傳、重試機制等核心元件。
import random
import asyncio
from typing import Optional, List, Set
from dataclasses import dataclass, field
from datetime import datetime, timedelta
import aiohttp
@dataclass
class Proxy:
"""代理實體"""
host: str
port: int
protocol: str = "http"
username: Optional[str] = None
password: Optional[str] = None
fail_count: int = 0
last_used: Optional[datetime] = None
last_fail: Optional[datetime] = None
@property
def url(self) -> str:
"""生成代理URL"""
auth = ""
if self.username and self.password:
auth = f"{self.username}:{self.password}@"
return f"{self.protocol}://{auth}{self.host}:{self.port}"
def mark_failed(self):
"""標記失敗"""
self.fail_count += 1
self.last_fail = datetime.now()
def reset_fail_count(self):
"""重置失敗計數"""
self.fail_count = 0
class ProxyPool:
"""
代理池管理器
功能:
- 代理輪換
- 失效代理自動剔除
- 代理健康檢查
"""
def __init__(
self,
max_fail_count: int = 3,
check_interval: int = 300,
test_url: str = "http://httpbin.org/ip"
):
self.proxies: List[Proxy] = []
self.blacklist: Set[str] = set()
self.max_fail_count = max_fail_count
self.check_interval = check_interval
self.test_url = test_url
self._lock = asyncio.Lock()
def add_proxy(self, proxy: Proxy):
"""新增代理"""
if proxy.url not in self.blacklist:
self.proxies.append(proxy)
def add_proxies_from_list(self, proxy_list: List[str]):
"""從字串列表批次新增代理(格式: host:port 或 protocol://host:port)"""
for proxy_str in proxy_list:
if "://" in proxy_str:
protocol, rest = proxy_str.split("://")
host, port = rest.split(":")
else:
protocol = "http"
host, port = proxy_str.split(":")
self.add_proxy(Proxy(host=host, port=int(port), protocol=protocol))
async def get_proxy(self) -> Optional[str]:
"""獲取一個可用代理"""
async with self._lock:
available = [p for p in self.proxies if p.fail_count < self.max_fail_count]
if not available:
return None
proxy = random.choice(available)
proxy.last_used = datetime.now()
return proxy.url
async def report_failure(self, proxy_url: str):
"""報告代理失敗"""
async with self._lock:
for proxy in self.proxies:
if proxy.url == proxy_url:
proxy.mark_failed()
if proxy.fail_count >= self.max_fail_count:
self.blacklist.add(proxy_url)
self.proxies.remove(proxy)
break
async def report_success(self, proxy_url: str):
"""報告代理成功"""
async with self._lock:
for proxy in self.proxies:
if proxy.url == proxy_url:
proxy.reset_fail_count()
break
async def health_check(self, timeout: int = 10) -> int:
"""健康檢查所有代理,返回可用代理數量"""
async def check_single(proxy: Proxy) -> bool:
try:
async with aiohttp.ClientSession() as session:
async with session.get(
self.test_url,
proxy=proxy.url,
timeout=aiohttp.ClientTimeout(total=timeout)
) as resp:
return resp.status == 200
except Exception:
return False
tasks = [check_single(p) for p in self.proxies]
results = await asyncio.gather(*tasks, return_exceptions=True)
valid_count = 0
async with self._lock:
for proxy, is_valid in zip(self.proxies[:], results):
if is_valid is True:
proxy.reset_fail_count()
valid_count += 1
else:
proxy.mark_failed()
return valid_count
import json
import os
from typing import Dict, Set, Any, Optional
from datetime import datetime
from pathlib import Path
class StateManager:
"""
爬蟲狀態管理器(支援斷點續傳)
功能:
- 記錄已完成的任務ID
- 儲存爬取進度
- 支援檔案持久化
"""
def __init__(self, state_file: str = "crawler_state.json"):
self.state_file = Path(state_file)
self.completed_ids: Set[str] = set()
self.progress: Dict[str, Any] = {}
self.metadata: Dict[str, Any] = {}
self._load_state()
def _load_state(self):
"""從檔案載入狀態"""
if self.state_file.exists():
with open(self.state_file, "r", encoding="utf-8") as f:
data = json.load(f)
self.completed_ids = set(data.get("completed_ids", []))
self.progress = data.get("progress", {})
self.metadata = data.get("metadata", {})
def save_state(self):
"""儲存狀態到檔案"""
data = {
"completed_ids": list(self.completed_ids),
"progress": self.progress,
"metadata": {
**self.metadata,
"last_saved": datetime.now().isoformat()
}
}
# 先寫入臨時檔案,再重新命名(原子操作)
temp_file = self.state_file.with_suffix(".tmp")
with open(temp_file, "w", encoding="utf-8") as f:
json.dump(data, f, ensure_ascii=False, indent=2)
temp_file.rename(self.state_file)
def mark_completed(self, task_id: str):
"""標記任務完成"""
self.completed_ids.add(task_id)
def is_completed(self, task_id: str) -> bool:
"""檢查任務是否已完成"""
return task_id in self.completed_ids
def update_progress(self, key: str, value: Any):
"""更新進度資訊"""
self.progress[key] = value
def get_progress(self, key: str, default: Any = None) -> Any:
"""獲取進度資訊"""
return self.progress.get(key, default)
def clear(self):
"""清除所有狀態"""
self.completed_ids.clear()
self.progress.clear()
if self.state_file.exists():
self.state_file.unlink()
import asyncio
import functools
from typing import TypeVar, Callable, Any
import logging
logger = logging.getLogger(__name__)
T = TypeVar("T")
def retry_async(
max_retries: int = 3,
delay: float = 1.0,
backoff: float = 2.0,
exceptions: tuple = (Exception,)
) -> Callable:
"""
非同步重試裝飾器
Args:
max_retries: 最大重試次數
delay: 初始延遲(秒)
backoff: 退避係數(每次重試延遲乘以此係數)
exceptions: 需要重試的異常型別
Usage:
@retry_async(max_retries=3, delay=1.0)
async def fetch_data(url):
...
"""
def decorator(func: Callable[..., T]) -> Callable[..., T]:
@functools.wraps(func)
async def wrapper(*args, **kwargs) -> T:
current_delay = delay
last_exception = None
for attempt in range(max_retries + 1):
try:
return await func(*args, **kwargs)
except exceptions as e:
last_exception = e
if attempt < max_retries:
logger.warning(
f"第 {attempt + 1}/{max_retries + 1} 次嘗試失敗: {e}. "
f"{current_delay:.1f}秒後重試..."
)
await asyncio.sleep(current_delay)
current_delay *= backoff
else:
logger.error(f"所有 {max_retries + 1} 次嘗試均失敗: {e}")
raise last_exception
return wrapper
return decorator
import random
class UserAgentRotator:
"""User-Agent 輪換器"""
# 常用桌面瀏覽器 UA
DESKTOP_UAS = [
"Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36",
"Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/119.0.0.0 Safari/537.36",
"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",
"Mozilla/5.0 (Windows NT 10.0; Win64; x64; rv:121.0) Gecko/20100101 Firefox/121.0",
"Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/605.1.15 (KHTML, like Gecko) Version/17.2 Safari/605.1.15",
"Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36 Edg/120.0.0.0",
]
# 移動端 UA
MOBILE_UAS = [
"Mozilla/5.0 (iPhone; CPU iPhone OS 17_2 like Mac OS X) AppleWebKit/605.1.15 (KHTML, like Gecko) Version/17.2 Mobile/15E148 Safari/604.1",
"Mozilla/5.0 (Linux; Android 14; Pixel 8) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Mobile Safari/537.36",
]
def __init__(self, include_mobile: bool = False):
self.user_agents = self.DESKTOP_UAS.copy()
if include_mobile:
self.user_agents.extend(self.MOBILE_UAS)
def get_random(self) -> str:
"""獲取隨機 User-Agent"""
return random.choice(self.user_agents)
def get_headers(self) -> dict:
"""獲取帶隨機 UA 的請求頭"""
return {
"User-Agent": self.get_random(),
"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",
}
import asyncio
import time
from typing import Optional
class RateLimiter:
"""
令牌桶限流器
用於控制請求頻率,避免觸發反爬蟲機制
"""
def __init__(self, rate: float, burst: int = 1):
"""
Args:
rate: 每秒允許的請求數
burst: 突發容量(令牌桶大小)
"""
self.rate = rate
self.burst = burst
self.tokens = burst
self.last_update = time.monotonic()
self._lock = asyncio.Lock()
async def acquire(self, timeout: Optional[float] = None) -> bool:
"""
獲取一個令牌
Args:
timeout: 超時時間(秒),None 表示無限等待
Returns:
是否成功獲取令牌
"""
start_time = time.monotonic()
while True:
async with self._lock:
now = time.monotonic()
# 補充令牌
elapsed = now - self.last_update
self.tokens = min(self.burst, self.tokens + elapsed * self.rate)
self.last_update = now
if self.tokens >= 1:
self.tokens -= 1
return True
# 檢查超時
if timeout is not None:
elapsed = time.monotonic() - start_time
if elapsed >= timeout:
return False
# 等待令牌
wait_time = (1 - self.tokens) / self.rate
await asyncio.sleep(min(wait_time, 0.1))
async def __aenter__(self):
await self.acquire()
return self
async def __aexit__(self, *args):
pass
在設計爬蟲時,必須考慮以下防護措施:
當目標資料(如樓棟資訊)無法直接從 API 獲取時:
TODO 註釋說明os.getenv("API_KEY").env.example 示例檔案每次完成程式碼後,必須附帶:
在提供爬蟲方案時,必須提醒使用者:
7w4.net收錄了海量優質技能外掛。
robots.txt 規則這個技能內容詳盡實用,提供了爬蟲開發所需的各類程式碼模板和質量規範,整體質量較好。優點是模板可直接使用、步驟清晰、考慮周全;不足是缺少完整示例和常見問題解答,部分內容可以更加精簡。對於需要快速搭建生產級爬蟲的使用者來說,這是一個值得參考的工具。