Python爬虫第12课:机器学习在爬虫中的应用
·
目录
专栏导读
🌸 欢迎来到Python办公自动化专栏—Python处理办公问题,解放您的双手
🏳️🌈 个人博客主页:请点击——> 个人的博客主页 求收藏
🏳️🌈 Github主页:请点击——> Github主页 求Star⭐
🏳️🌈 知乎主页:请点击——> 知乎主页 求关注
🏳️🌈 CSDN博客主页:请点击——> CSDN的博客主页 求关注
👍 该系列文章专栏:请点击——>Python办公自动化专栏 求订阅
🕷 此外还有爬虫专栏:请点击——>Python爬虫基础专栏 求订阅
📕 此外还有python基础专栏:请点击——>Python基础学习专栏 求订阅
文章作者技术和水平有限,如果文中出现错误,希望大家能指正🙏
❤️ 欢迎各位佬关注! ❤️
Python爬虫第12课:机器学习在爬虫中的应用
课程目标
- 掌握机器学习在爬虫中的应用场景
- 学会使用机器学习进行智能内容识别与分类
- 了解反爬虫检测与对抗技术
- 实现数据挖掘与模式识别
- 构建预测性爬虫调度系统
1. 机器学习在爬虫中的应用概述
1.1 应用场景
- 智能内容识别:自动识别和分类网页内容
- 反爬虫对抗:检测和绕过反爬虫机制
- 数据质量评估:评估爬取数据的质量和可信度
- 智能调度:基于历史数据预测最佳爬取时机
- 异常检测:识别异常网页和数据
1.2 技术栈
# 机器学习库
import numpy as np
import pandas as pd
from sklearn.feature_extraction.text import TfidfVectorizer
from sklearn.model_selection import train_test_split
from sklearn.ensemble import RandomForestClassifier
from sklearn.linear_model import LogisticRegression
from sklearn.metrics import classification_report, accuracy_score
from sklearn.cluster import KMeans
import joblib
# 深度学习库
import tensorflow as tf
from tensorflow.keras.models import Sequential
from tensorflow.keras.layers import Dense, LSTM, Embedding
from transformers import pipeline, AutoTokenizer, AutoModel
# 自然语言处理
import jieba
import re
from textblob import TextBlob
import spacy
# 图像处理
from PIL import Image
import cv2
import pytesseract
# 数据处理
import json
import time
import logging
from typing import Dict, List, Any, Optional, Tuple
from dataclasses import dataclass
2. 智能内容识别与分类
2.1 文本分类器
# text_classifier.py
import numpy as np
import pandas as pd
from sklearn.feature_extraction.text import TfidfVectorizer
from sklearn.model_selection import train_test_split
from sklearn.ensemble import RandomForestClassifier
from sklearn.linear_model import LogisticRegression
from sklearn.metrics import classification_report, accuracy_score
import jieba
import re
import joblib
import logging
from typing import List, Dict, Any, Optional, Tuple
class TextClassifier:
"""文本分类器"""
def __init__(self, model_type: str = 'random_forest'):
self.model_type = model_type
self.vectorizer = TfidfVectorizer(
max_features=10000,
stop_words=None,
ngram_range=(1, 2)
)
self.model = None
self.label_encoder = {}
self.reverse_label_encoder = {}
self.logger = logging.getLogger(__name__)
# 初始化模型
if model_type == 'random_forest':
self.model = RandomForestClassifier(
n_estimators=100,
random_state=42,
n_jobs=-1
)
elif model_type == 'logistic_regression':
self.model = LogisticRegression(
random_state=42,
max_iter=1000
)
else:
raise ValueError(f"不支持的模型类型: {model_type}")
def preprocess_text(self, text: str) -> str:
"""文本预处理"""
if not isinstance(text, str):
return ""
# 移除HTML标签
text = re.sub(r'<[^>]+>', '', text)
# 移除特殊字符,保留中文、英文、数字
text = re.sub(r'[^\u4e00-\u9fa5a-zA-Z0-9\s]', '', text)
# 中文分词
words = jieba.cut(text)
# 过滤停用词
stop_words = {'的', '了', '在', '是', '我', '有', '和', '就',
'不', '人', '都', '一', '一个', '上', '也', '很',
'到', '说', '要', '去', '你', '会', '着', '没有',
'看', '好', '自己', '这'}
filtered_words = [word for word in words if word not in stop_words and len(word) > 1]
return ' '.join(filtered_words)
def prepare_data(self, texts: List[str], labels: List[str]) -> Tuple[np.ndarray, np.ndarray]:
"""准备训练数据"""
# 文本预处理
processed_texts = [self.preprocess_text(text) for text in texts]
# 标签编码
unique_labels = list(set(labels))
self.label_encoder = {label: idx for idx, label in enumerate(unique_labels)}
self.reverse_label_encoder = {idx: label for label, idx in self.label_encoder.items()}
encoded_labels = [self.label_encoder[label] for label in labels]
# 文本向量化
X = self.vectorizer.fit_transform(processed_texts)
y = np.array(encoded_labels)
return X, y
def train(self, texts: List[str], labels: List[str], test_size: float = 0.2) -> Dict[str, Any]:
"""训练模型"""
self.logger.info("开始训练文本分类模型")
# 准备数据
X, y = self.prepare_data(texts, labels)
# 分割训练集和测试集
X_train, X_test, y_train, y_test = train_test_split(
X, y, test_size=test_size, random_state=42, stratify=y
)
# 训练模型
self.model.fit(X_train, y_train)
# 评估模型
y_pred = self.model.predict(X_test)
accuracy = accuracy_score(y_test, y_pred)
# 生成分类报告
target_names = [self.reverse_label_encoder[i] for i in range(len(self.label_encoder))]
report = classification_report(y_test, y_pred, target_names=target_names, output_dict=True)
self.logger.info(f"模型训练完成,准确率: {accuracy:.4f}")
return {
'accuracy': accuracy,
'classification_report': report,
'feature_count': X.shape[1],
'sample_count': X.shape[0]
}
def predict(self, texts: List[str]) -> List[Dict[str, Any]]:
"""预测文本分类"""
if self.model is None:
raise ValueError("模型未训练,请先调用train方法")
# 文本预处理
processed_texts = [self.preprocess_text(text) for text in texts]
# 向量化
X = self.vectorizer.transform(processed_texts)
# 预测
predictions = self.model.predict(X)
probabilities = self.model.predict_proba(X)
results = []
for i, (pred, prob) in enumerate(zip(predictions, probabilities)):
result = {
'text': texts[i],
'predicted_label': self.reverse_label_encoder[pred],
'confidence': float(np.max(prob)),
'probabilities': {
self.reverse_label_encoder[j]: float(prob[j])
for j in range(len(prob))
}
}
results.append(result)
return results
def save_model(self, file_path: str):
"""保存模型"""
model_data = {
'model': self.model,
'vectorizer': self.vectorizer,
'label_encoder': self.label_encoder,
'reverse_label_encoder': self.reverse_label_encoder,
'model_type': self.model_type
}
joblib.dump(model_data, file_path)
self.logger.info(f"模型已保存到: {file_path}")
def load_model(self, file_path: str):
"""加载模型"""
model_data = joblib.load(file_path)
self.model = model_data['model']
self.vectorizer = model_data['vectorizer']
self.label_encoder = model_data['label_encoder']
self.reverse_label_encoder = model_data['reverse_label_encoder']
self.model_type = model_data['model_type']
self.logger.info(f"模型已从 {file_path} 加载")
class ContentCategoryClassifier:
"""内容类别分类器"""
def __init__(self):
self.classifier = TextClassifier('random_forest')
self.categories = {
'news': '新闻',
'tech': '科技',
'sports': '体育',
'entertainment': '娱乐',
'finance': '财经',
'education': '教育',
'health': '健康',
'travel': '旅游',
'food': '美食',
'other': '其他'
}
self.logger = logging.getLogger(__name__)
def extract_features(self, html_content: str, url: str = "") -> Dict[str, Any]:
"""提取网页特征"""
from bs4 import BeautifulSoup
soup = BeautifulSoup(html_content, 'html.parser')
# 提取文本内容
title = soup.find('title')
title_text = title.get_text() if title else ""
# 移除脚本和样式
for script in soup(["script", "style"]):
script.decompose()
body_text = soup.get_text()
# 提取关键词
keywords_meta = soup.find('meta', attrs={'name': 'keywords'})
keywords = keywords_meta.get('content', '') if keywords_meta else ""
# 提取描述
description_meta = soup.find('meta', attrs={'name': 'description'})
description = description_meta.get('content', '') if description_meta else ""
# URL特征
url_features = self._extract_url_features(url)
return {
'title': title_text,
'body': body_text,
'keywords': keywords,
'description': description,
'url_features': url_features,
'combined_text': f"{title_text} {description} {keywords} {body_text}"
}
def _extract_url_features(self, url: str) -> Dict[str, Any]:
"""提取URL特征"""
import urllib.parse
parsed = urllib.parse.urlparse(url)
path_parts = parsed.path.split('/')
return {
'domain': parsed.netloc,
'path_depth': len([p for p in path_parts if p]),
'has_news_keyword': any(keyword in url.lower()
for keyword in ['news', 'article', 'post']),
'has_tech_keyword': any(keyword in url.lower()
for keyword in ['tech', 'technology', 'it']),
'has_sports_keyword': any(keyword in url.lower()
for keyword in ['sport', 'football', 'basketball'])
}
def train_with_sample_data(self):
"""使用示例数据训练模型"""
# 示例训练数据
sample_data = [
("苹果公司发布新款iPhone,搭载最新A17芯片", "tech"),
("中国足球队在世界杯预选赛中获胜", "sports"),
("股市今日大涨,科技股领涨", "finance"),
("新冠疫苗接种率达到80%", "health"),
("教育部发布新的高考改革方案", "education"),
("著名演员获得奥斯卡最佳男主角奖", "entertainment"),
("北京故宫博物院推出新的文物展览", "news"),
("马尔代夫成为热门旅游目的地", "travel"),
("川菜成为世界美食文化遗产", "food"),
("人工智能技术在医疗领域的应用", "tech"),
("NBA总决赛即将开始", "sports"),
("央行宣布降准政策", "finance"),
("专家建议每天运动30分钟", "health"),
("在线教育平台用户数量激增", "education"),
("好莱坞新片票房破纪录", "entertainment"),
("政府发布环保新政策", "news"),
("日本樱花季吸引大量游客", "travel"),
("传统中式点心制作工艺", "food")
]
texts = [item[0] for item in sample_data]
labels = [item[1] for item in sample_data]
# 训练模型
results = self.classifier.train(texts, labels)
self.logger.info(f"内容分类模型训练完成: {results}")
return results
def classify_content(self, html_content: str, url: str = "") -> Dict[str, Any]:
"""分类网页内容"""
# 提取特征
features = self.extract_features(html_content, url)
# 使用组合文本进行分类
predictions = self.classifier.predict([features['combined_text']])
if predictions:
prediction = predictions[0]
return {
'category': prediction['predicted_label'],
'category_name': self.categories.get(prediction['predicted_label'], '未知'),
'confidence': prediction['confidence'],
'probabilities': prediction['probabilities'],
'features': features
}
return {
'category': 'other',
'category_name': '其他',
'confidence': 0.0,
'probabilities': {},
'features': features
}
# 使用示例
if __name__ == "__main__":
# 配置日志
logging.basicConfig(level=logging.INFO)
# 创建内容分类器
classifier = ContentCategoryClassifier()
# 训练模型
classifier.train_with_sample_data()
# 测试分类
test_html = """
<html>
<head>
<title>人工智能技术突破:GPT-4在自然语言处理领域的新进展</title>
<meta name="keywords" content="人工智能,GPT-4,自然语言处理,技术">
<meta name="description" content="最新的GPT-4模型在自然语言处理任务中表现出色">
</head>
<body>
<h1>人工智能技术突破</h1>
<p>OpenAI公司最新发布的GPT-4模型在自然语言处理领域取得了重大突破...</p>
</body>
</html>
"""
result = classifier.classify_content(test_html, "https://tech.example.com/ai-news")
print(f"分类结果: {result}")
# 保存模型
classifier.classifier.save_model('content_classifier.pkl')
2.2 图像识别与分类
# image_classifier.py
import cv2
import numpy as np
from PIL import Image
import pytesseract
import tensorflow as tf
from tensorflow.keras.applications import VGG16
from tensorflow.keras.applications.vgg16 import preprocess_input, decode_predictions
from tensorflow.keras.preprocessing import image
import requests
import io
import logging
from typing import List, Dict, Any, Optional, Tuple
class ImageClassifier:
"""图像分类器"""
def __init__(self):
self.logger = logging.getLogger(__name__)
# 加载预训练的VGG16模型
self.model = VGG16(weights='imagenet')
self.logger.info("VGG16模型加载完成")
def download_image(self, url: str) -> Optional[np.ndarray]:
"""下载图像"""
try:
response = requests.get(url, timeout=10)
response.raise_for_status()
# 将字节数据转换为图像
img = Image.open(io.BytesIO(response.content))
img_array = np.array(img)
return img_array
except Exception as e:
self.logger.error(f"下载图像失败 {url}: {e}")
return None
def preprocess_image(self, img_array: np.ndarray) -> np.ndarray:
"""预处理图像"""
# 转换为PIL图像
if len(img_array.shape) == 3 and img_array.shape[2] == 4:
# 移除alpha通道
img_array = img_array[:, :, :3]
img = Image.fromarray(img_array)
# 调整大小到224x224
img = img.resize((224, 224))
# 转换为数组
img_array = image.img_to_array(img)
img_array = np.expand_dims(img_array, axis=0)
img_array = preprocess_input(img_array)
return img_array
def classify_image(self, img_array: np.ndarray, top_k: int = 5) -> List[Dict[str, Any]]:
"""分类图像"""
try:
# 预处理图像
processed_img = self.preprocess_image(img_array)
# 预测
predictions = self.model.predict(processed_img)
# 解码预测结果
decoded_predictions = decode_predictions(predictions, top=top_k)[0]
results = []
for pred in decoded_predictions:
results.append({
'class_id': pred[0],
'class_name': pred[1],
'confidence': float(pred[2])
})
return results
except Exception as e:
self.logger.error(f"图像分类失败: {e}")
return []
def classify_image_from_url(self, url: str, top_k: int = 5) -> List[Dict[str, Any]]:
"""从URL分类图像"""
img_array = self.download_image(url)
if img_array is not None:
return self.classify_image(img_array, top_k)
return []
class CaptchaRecognizer:
"""验证码识别器"""
def __init__(self):
self.logger = logging.getLogger(__name__)
# 配置Tesseract
# 需要安装Tesseract OCR: https://github.com/tesseract-ocr/tesseract
# pytesseract.pytesseract.tesseract_cmd = r'C:\Program Files\Tesseract-OCR\tesseract.exe'
def preprocess_captcha(self, img_array: np.ndarray) -> np.ndarray:
"""预处理验证码图像"""
# 转换为灰度图
if len(img_array.shape) == 3:
gray = cv2.cvtColor(img_array, cv2.COLOR_RGB2GRAY)
else:
gray = img_array
# 二值化
_, binary = cv2.threshold(gray, 127, 255, cv2.THRESH_BINARY)
# 去噪
kernel = np.ones((2, 2), np.uint8)
cleaned = cv2.morphologyEx(binary, cv2.MORPH_CLOSE, kernel)
# 放大图像以提高识别率
height, width = cleaned.shape
enlarged = cv2.resize(cleaned, (width * 3, height * 3), interpolation=cv2.INTER_CUBIC)
return enlarged
def recognize_text_captcha(self, img_array: np.ndarray) -> str:
"""识别文本验证码"""
try:
# 预处理图像
processed_img = self.preprocess_captcha(img_array)
# 使用Tesseract识别
config = '--psm 8 -c tessedit_char_whitelist=0123456789ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz'
text = pytesseract.image_to_string(processed_img, config=config)
# 清理结果
text = text.strip().replace(' ', '').replace('\n', '')
return text
except Exception as e:
self.logger.error(f"验证码识别失败: {e}")
return ""
def recognize_captcha_from_url(self, url: str) -> str:
"""从URL识别验证码"""
try:
response = requests.get(url, timeout=10)
response.raise_for_status()
img = Image.open(io.BytesIO(response.content))
img_array = np.array(img)
return self.recognize_text_captcha(img_array)
except Exception as e:
self.logger.error(f"从URL识别验证码失败 {url}: {e}")
return ""
class ImageContentAnalyzer:
"""图像内容分析器"""
def __init__(self):
self.image_classifier = ImageClassifier()
self.captcha_recognizer = CaptchaRecognizer()
self.logger = logging.getLogger(__name__)
def analyze_webpage_images(self, html_content: str, base_url: str = "") -> Dict[str, Any]:
"""分析网页中的图像"""
from bs4 import BeautifulSoup
import urllib.parse
soup = BeautifulSoup(html_content, 'html.parser')
images = soup.find_all('img')
analysis_results = {
'total_images': len(images),
'classified_images': [],
'captcha_images': [],
'failed_images': []
}
for img in images:
src = img.get('src', '')
alt = img.get('alt', '')
if not src:
continue
# 构建完整URL
if src.startswith('//'):
src = 'https:' + src
elif src.startswith('/'):
src = urllib.parse.urljoin(base_url, src)
elif not src.startswith('http'):
src = urllib.parse.urljoin(base_url, src)
try:
# 判断是否为验证码
if self._is_captcha_image(src, alt):
captcha_text = self.captcha_recognizer.recognize_captcha_from_url(src)
analysis_results['captcha_images'].append({
'url': src,
'alt': alt,
'recognized_text': captcha_text
})
else:
# 分类普通图像
classifications = self.image_classifier.classify_image_from_url(src)
if classifications:
analysis_results['classified_images'].append({
'url': src,
'alt': alt,
'classifications': classifications
})
else:
analysis_results['failed_images'].append({
'url': src,
'alt': alt,
'error': 'Classification failed'
})
except Exception as e:
analysis_results['failed_images'].append({
'url': src,
'alt': alt,
'error': str(e)
})
return analysis_results
def _is_captcha_image(self, src: str, alt: str) -> bool:
"""判断是否为验证码图像"""
captcha_keywords = ['captcha', 'verify', 'code', '验证码', '验证', 'checkcode']
src_lower = src.lower()
alt_lower = alt.lower()
return any(keyword in src_lower or keyword in alt_lower for keyword in captcha_keywords)
# 使用示例
if __name__ == "__main__":
# 配置日志
logging.basicConfig(level=logging.INFO)
# 创建图像分类器
classifier = ImageClassifier()
# 测试图像分类
test_image_url = "https://example.com/test-image.jpg"
results = classifier.classify_image_from_url(test_image_url)
print(f"图像分类结果: {results}")
# 创建验证码识别器
captcha_recognizer = CaptchaRecognizer()
# 测试验证码识别
captcha_url = "https://example.com/captcha.jpg"
captcha_text = captcha_recognizer.recognize_captcha_from_url(captcha_url)
print(f"验证码识别结果: {captcha_text}")
# 创建图像内容分析器
analyzer = ImageContentAnalyzer()
# 分析网页图像
html_content = """
<html>
<body>
<img src="/logo.png" alt="Company Logo">
<img src="/captcha.jpg" alt="验证码">
<img src="https://example.com/product.jpg" alt="Product Image">
</body>
</html>
"""
analysis = analyzer.analyze_webpage_images(html_content, "https://example.com")
print(f"图像分析结果: {analysis}")
3. 反爬虫检测与对抗
3.1 行为模式分析
# behavior_analysis.py
import numpy as np
import pandas as pd
from sklearn.ensemble import IsolationForest
from sklearn.preprocessing import StandardScaler
from sklearn.cluster import DBSCAN
import time
import random
import logging
from typing import Dict, List, Any, Optional, Tuple
from dataclasses import dataclass
from collections import defaultdict, deque
import json
@dataclass
class RequestPattern:
"""请求模式"""
timestamp: float
url: str
user_agent: str
ip_address: str
response_time: float
status_code: int
content_length: int
referer: str = ""
session_id: str = ""
class BehaviorAnalyzer:
"""行为模式分析器"""
def __init__(self, window_size: int = 100):
self.window_size = window_size
self.request_history = deque(maxlen=window_size)
self.user_sessions = defaultdict(list)
self.anomaly_detector = IsolationForest(contamination=0.1, random_state=42)
self.scaler = StandardScaler()
self.logger = logging.getLogger(__name__)
# 行为特征权重
self.feature_weights = {
'request_frequency': 0.3,
'response_time_variance': 0.2,
'url_diversity': 0.2,
'user_agent_consistency': 0.15,
'session_duration': 0.15
}
def add_request(self, pattern: RequestPattern):
"""添加请求记录"""
self.request_history.append(pattern)
# 按会话分组
session_key = f"{pattern.ip_address}_{pattern.session_id}"
self.user_sessions[session_key].append(pattern)
def extract_features(self, patterns: List[RequestPattern]) -> np.ndarray:
"""提取行为特征"""
if not patterns:
return np.array([])
features = []
# 按会话分组分析
sessions = defaultdict(list)
for pattern in patterns:
session_key = f"{pattern.ip_address}_{pattern.session_id}"
sessions[session_key].append(pattern)
for session_id, session_patterns in sessions.items():
if len(session_patterns) < 2:
continue
# 计算特征
session_features = self._calculate_session_features(session_patterns)
features.append(session_features)
return np.array(features) if features else np.array([])
def _calculate_session_features(self, patterns: List[RequestPattern]) -> List[float]:
"""计算会话特征"""
if not patterns:
return [0.0] * 10
# 时间特征
timestamps = [p.timestamp for p in patterns]
time_intervals = np.diff(timestamps)
# 请求频率特征
avg_interval = np.mean(time_intervals) if len(time_intervals) > 0 else 0
interval_variance = np.var(time_intervals) if len(time_intervals) > 0 else 0
# 响应时间特征
response_times = [p.response_time for p in patterns]
avg_response_time = np.mean(response_times)
response_time_variance = np.var(response_times)
# URL多样性
unique_urls = len(set(p.url for p in patterns))
url_diversity = unique_urls / len(patterns)
# User-Agent一致性
user_agents = [p.user_agent for p in patterns]
unique_user_agents = len(set(user_agents))
ua_consistency = 1.0 - (unique_user_agents - 1) / len(patterns)
# 会话持续时间
session_duration = timestamps[-1] - timestamps[0] if len(timestamps) > 1 else 0
# 状态码分布
status_codes = [p.status_code for p in patterns]
success_rate = sum(1 for code in status_codes if 200 <= code < 300) / len(status_codes)
# 内容长度特征
content_lengths = [p.content_length for p in patterns]
avg_content_length = np.mean(content_lengths)
content_length_variance = np.var(content_lengths)
return [
avg_interval,
interval_variance,
avg_response_time,
response_time_variance,
url_diversity,
ua_consistency,
session_duration,
success_rate,
avg_content_length,
content_length_variance
]
def detect_anomalies(self, patterns: List[RequestPattern] = None) -> Dict[str, Any]:
"""检测异常行为"""
if patterns is None:
patterns = list(self.request_history)
if len(patterns) < 10:
return {'anomalies': [], 'normal_sessions': 0, 'anomaly_rate': 0.0}
# 提取特征
features = self.extract_features(patterns)
if len(features) == 0:
return {'anomalies': [], 'normal_sessions': 0, 'anomaly_rate': 0.0}
# 标准化特征
features_scaled = self.scaler.fit_transform(features)
# 异常检测
anomaly_labels = self.anomaly_detector.fit_predict(features_scaled)
# 分析结果
anomaly_indices = np.where(anomaly_labels == -1)[0]
normal_count = np.sum(anomaly_labels == 1)
anomaly_rate = len(anomaly_indices) / len(features)
# 获取异常会话详情
sessions = defaultdict(list)
for pattern in patterns:
session_key = f"{pattern.ip_address}_{pattern.session_id}"
sessions[session_key].append(pattern)
session_list = list(sessions.keys())
anomalous_sessions = []
for idx in anomaly_indices:
if idx < len(session_list):
session_key = session_list[idx]
session_patterns = sessions[session_key]
anomalous_sessions.append({
'session_id': session_key,
'pattern_count': len(session_patterns),
'time_span': session_patterns[-1].timestamp - session_patterns[0].timestamp,
'features': features[idx].tolist(),
'risk_score': self._calculate_risk_score(session_patterns)
})
return {
'anomalies': anomalous_sessions,
'normal_sessions': normal_count,
'anomaly_rate': anomaly_rate,
'total_sessions': len(features)
}
def _calculate_risk_score(self, patterns: List[RequestPattern]) -> float:
"""计算风险评分"""
if not patterns:
return 0.0
score = 0.0
# 高频请求检测
timestamps = [p.timestamp for p in patterns]
if len(timestamps) > 1:
time_intervals = np.diff(timestamps)
avg_interval = np.mean(time_intervals)
if avg_interval < 1.0: # 平均间隔小于1秒
score += 30.0
elif avg_interval < 5.0: # 平均间隔小于5秒
score += 15.0
# User-Agent检测
user_agents = set(p.user_agent for p in patterns)
if len(user_agents) > 1:
score += 20.0 # 频繁更换User-Agent
# 响应时间异常
response_times = [p.response_time for p in patterns]
if response_times:
avg_response_time = np.mean(response_times)
if avg_response_time < 0.1: # 响应时间过快
score += 25.0
# 状态码异常
status_codes = [p.status_code for p in patterns]
error_rate = sum(1 for code in status_codes if code >= 400) / len(status_codes)
if error_rate > 0.5: # 错误率过高
score += 15.0
# URL模式检测
urls = [p.url for p in patterns]
unique_urls = len(set(urls))
if unique_urls / len(urls) < 0.3: # URL重复率过高
score += 10.0
return min(score, 100.0) # 最高100分
class AntiDetectionStrategy:
"""反检测策略"""
def __init__(self):
self.logger = logging.getLogger(__name__)
self.user_agents = [
'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/91.0.4472.124 Safari/537.36',
'Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/91.0.4472.124 Safari/537.36',
'Mozilla/5.0 (Windows NT 10.0; Win64; x64; rv:89.0) Gecko/20100101 Firefox/89.0',
'Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/605.1.15 (KHTML, like Gecko) Version/14.1.1 Safari/605.1.15',
'Mozilla/5.0 (X11; Linux x86_64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/91.0.4472.124 Safari/537.36'
]
self.request_intervals = {
'conservative': (3, 8), # 保守策略:3-8秒
'moderate': (1, 5), # 中等策略:1-5秒
'aggressive': (0.5, 2) # 激进策略:0.5-2秒
}
def get_random_user_agent(self) -> str:
"""获取随机User-Agent"""
return random.choice(self.user_agents)
def get_request_delay(self, strategy: str = 'moderate') -> float:
"""获取请求延迟"""
if strategy not in self.request_intervals:
strategy = 'moderate'
min_delay, max_delay = self.request_intervals[strategy]
return random.uniform(min_delay, max_delay)
def generate_realistic_headers(self) -> Dict[str, str]:
"""生成真实的请求头"""
headers = {
'User-Agent': self.get_random_user_agent(),
'Accept': 'text/html,application/xhtml+xml,application/xml;q=0.9,image/webp,*/*;q=0.8',
'Accept-Language': 'en-US,en;q=0.5',
'Accept-Encoding': 'gzip, deflate',
'Connection': 'keep-alive',
'Upgrade-Insecure-Requests': '1',
}
# 随机添加一些可选头部
optional_headers = {
'Cache-Control': 'max-age=0',
'Sec-Fetch-Dest': 'document',
'Sec-Fetch-Mode': 'navigate',
'Sec-Fetch-Site': 'none',
'Sec-Fetch-User': '?1'
}
for key, value in optional_headers.items():
if random.random() > 0.3: # 70%概率添加
headers[key] = value
return headers
def simulate_human_behavior(self, session_requests: int = None) -> List[Dict[str, Any]]:
"""模拟人类行为模式"""
if session_requests is None:
session_requests = random.randint(5, 20)
behavior_plan = []
current_time = time.time()
for i in range(session_requests):
# 模拟不同类型的请求
if i == 0:
# 首次访问
action = {
'type': 'initial_visit',
'delay': 0,
'headers': self.generate_realistic_headers(),
'timestamp': current_time
}
elif random.random() < 0.3:
# 页面跳转
action = {
'type': 'navigation',
'delay': random.uniform(2, 10),
'headers': self.generate_realistic_headers(),
'timestamp': current_time
}
elif random.random() < 0.2:
# 搜索行为
action = {
'type': 'search',
'delay': random.uniform(5, 15),
'headers': self.generate_realistic_headers(),
'timestamp': current_time
}
else:
# 普通浏览
action = {
'type': 'browse',
'delay': self.get_request_delay('moderate'),
'headers': self.generate_realistic_headers(),
'timestamp': current_time
}
behavior_plan.append(action)
current_time += action['delay']
return behavior_plan
# 使用示例
if __name__ == "__main__":
# 配置日志
logging.basicConfig(level=logging.INFO)
# 创建行为分析器
analyzer = BehaviorAnalyzer()
# 模拟请求数据
current_time = time.time()
# 正常用户行为
for i in range(20):
pattern = RequestPattern(
timestamp=current_time + i * random.uniform(2, 10),
url=f"https://example.com/page{random.randint(1, 10)}",
user_agent="Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36",
ip_address="192.168.1.100",
response_time=random.uniform(0.5, 2.0),
status_code=200,
content_length=random.randint(1000, 10000),
session_id="normal_session"
)
analyzer.add_request(pattern)
# 异常用户行为(爬虫)
for i in range(30):
pattern = RequestPattern(
timestamp=current_time + i * 0.5, # 高频请求
url=f"https://example.com/api/data{i}",
user_agent="Python-requests/2.25.1", # 明显的爬虫标识
ip_address="10.0.0.50",
response_time=0.1, # 响应时间过快
status_code=200,
content_length=500,
session_id="bot_session"
)
analyzer.add_request(pattern)
# 检测异常
results = analyzer.detect_anomalies()
print(f"异常检测结果: {json.dumps(results, indent=2)}")
# 创建反检测策略
strategy = AntiDetectionStrategy()
# 生成人类行为模式
behavior_plan = strategy.simulate_human_behavior(10)
print(f"\n人类行为模拟: {json.dumps(behavior_plan, indent=2)}")
4. 数据挖掘与模式识别
4.1 内容相似度检测
# similarity_detection.py
import numpy as np
import pandas as pd
from sklearn.feature_extraction.text import TfidfVectorizer
from sklearn.metrics.pairwise import cosine_similarity
from sklearn.cluster import KMeans
import jieba
import re
import hashlib
from typing import List, Dict, Any, Tuple, Optional
import logging
from dataclasses import dataclass
from collections import defaultdict
import json
@dataclass
class ContentItem:
"""内容项"""
id: str
title: str
content: str
url: str
timestamp: float
metadata: Dict[str, Any] = None
class ContentSimilarityDetector:
"""内容相似度检测器"""
def __init__(self, similarity_threshold: float = 0.8):
self.similarity_threshold = similarity_threshold
self.vectorizer = TfidfVectorizer(
max_features=5000,
stop_words=None,
ngram_range=(1, 2)
)
self.content_vectors = None
self.content_items = []
self.logger = logging.getLogger(__name__)
def preprocess_text(self, text: str) -> str:
"""文本预处理"""
if not isinstance(text, str):
return ""
# 移除HTML标签
text = re.sub(r'<[^>]+>', '', text)
# 移除特殊字符
text = re.sub(r'[^\u4e00-\u9fa5a-zA-Z0-9\s]', '', text)
# 中文分词
words = jieba.cut(text)
# 过滤停用词
stop_words = {'的', '了', '在', '是', '我', '有', '和', '就',
'不', '人', '都', '一', '一个', '上', '也', '很',
'到', '说', '要', '去', '你', '会', '着', '没有'}
filtered_words = [word for word in words if word not in stop_words and len(word) > 1]
return ' '.join(filtered_words)
def add_content(self, item: ContentItem):
"""添加内容项"""
self.content_items.append(item)
# 重新计算向量(在实际应用中可能需要增量更新)
if len(self.content_items) % 100 == 0: # 每100个重新计算一次
self._update_vectors()
def _update_vectors(self):
"""更新内容向量"""
if not self.content_items:
return
# 预处理所有文本
processed_texts = []
for item in self.content_items:
combined_text = f"{item.title} {item.content}"
processed_text = self.preprocess_text(combined_text)
processed_texts.append(processed_text)
# 计算TF-IDF向量
self.content_vectors = self.vectorizer.fit_transform(processed_texts)
self.logger.info(f"更新了 {len(self.content_items)} 个内容项的向量")
def find_similar_content(self, query_item: ContentItem, top_k: int = 5) -> List[Dict[str, Any]]:
"""查找相似内容"""
if not self.content_items or self.content_vectors is None:
return []
# 预处理查询文本
query_text = f"{query_item.title} {query_item.content}"
processed_query = self.preprocess_text(query_text)
# 计算查询向量
query_vector = self.vectorizer.transform([processed_query])
# 计算相似度
similarities = cosine_similarity(query_vector, self.content_vectors)[0]
# 获取最相似的内容
similar_indices = np.argsort(similarities)[::-1][:top_k]
results = []
for idx in similar_indices:
if similarities[idx] >= self.similarity_threshold:
similar_item = self.content_items[idx]
results.append({
'item': similar_item,
'similarity': float(similarities[idx]),
'match_type': self._classify_similarity(similarities[idx])
})
return results
def _classify_similarity(self, similarity: float) -> str:
"""分类相似度等级"""
if similarity >= 0.95:
return 'identical'
elif similarity >= 0.85:
return 'very_similar'
elif similarity >= 0.7:
return 'similar'
else:
return 'somewhat_similar'
def detect_duplicate_content(self) -> Dict[str, Any]:
"""检测重复内容"""
if not self.content_items or self.content_vectors is None:
self._update_vectors()
if self.content_vectors is None:
return {'duplicates': [], 'total_groups': 0}
# 计算所有内容之间的相似度
similarity_matrix = cosine_similarity(self.content_vectors)
# 查找重复组
duplicate_groups = []
processed_indices = set()
for i in range(len(self.content_items)):
if i in processed_indices:
continue
# 查找与当前项相似的所有项
similar_indices = np.where(similarity_matrix[i] >= self.similarity_threshold)[0]
if len(similar_indices) > 1: # 包括自己,所以要大于1
group = []
for idx in similar_indices:
if idx not in processed_indices:
group.append({
'item': self.content_items[idx],
'similarity_to_first': float(similarity_matrix[i][idx])
})
processed_indices.add(idx)
if len(group) > 1:
duplicate_groups.append({
'group_id': len(duplicate_groups),
'items': group,
'group_size': len(group)
})
return {
'duplicates': duplicate_groups,
'total_groups': len(duplicate_groups),
'total_duplicates': sum(group['group_size'] for group in duplicate_groups)
}
class ContentClusterAnalyzer:
"""内容聚类分析器"""
def __init__(self, n_clusters: int = 10):
self.n_clusters = n_clusters
self.vectorizer = TfidfVectorizer(
max_features=3000,
stop_words=None,
ngram_range=(1, 2)
)
self.kmeans = KMeans(n_clusters=n_clusters, random_state=42)
self.content_items = []
self.cluster_labels = None
self.logger = logging.getLogger(__name__)
def preprocess_text(self, text: str) -> str:
"""文本预处理"""
if not isinstance(text, str):
return ""
# 移除HTML标签
text = re.sub(r'<[^>]+>', '', text)
# 移除特殊字符
text = re.sub(r'[^\u4e00-\u9fa5a-zA-Z0-9\s]', '', text)
# 中文分词
words = jieba.cut(text)
# 过滤停用词
stop_words = {'的', '了', '在', '是', '我', '有', '和', '就',
'不', '人', '都', '一', '一个', '上', '也', '很'}
filtered_words = [word for word in words if word not in stop_words and len(word) > 1]
return ' '.join(filtered_words)
def add_content_batch(self, items: List[ContentItem]):
"""批量添加内容"""
self.content_items.extend(items)
self.logger.info(f"添加了 {len(items)} 个内容项,总计 {len(self.content_items)} 个")
def perform_clustering(self) -> Dict[str, Any]:
"""执行聚类分析"""
if len(self.content_items) < self.n_clusters:
self.logger.warning(f"内容数量 ({len(self.content_items)}) 少于聚类数量 ({self.n_clusters})")
return {'clusters': [], 'total_items': len(self.content_items)}
# 预处理文本
processed_texts = []
for item in self.content_items:
combined_text = f"{item.title} {item.content}"
processed_text = self.preprocess_text(combined_text)
processed_texts.append(processed_text)
# 向量化
vectors = self.vectorizer.fit_transform(processed_texts)
# 聚类
self.cluster_labels = self.kmeans.fit_predict(vectors)
# 分析聚类结果
clusters = defaultdict(list)
for i, label in enumerate(self.cluster_labels):
clusters[label].append(self.content_items[i])
# 生成聚类报告
cluster_analysis = []
for cluster_id, items in clusters.items():
# 提取关键词
cluster_texts = [f"{item.title} {item.content}" for item in items]
cluster_keywords = self._extract_cluster_keywords(cluster_texts, cluster_id)
cluster_analysis.append({
'cluster_id': int(cluster_id),
'size': len(items),
'keywords': cluster_keywords,
'sample_titles': [item.title for item in items[:5]], # 前5个标题作为样本
'items': items
})
# 按聚类大小排序
cluster_analysis.sort(key=lambda x: x['size'], reverse=True)
return {
'clusters': cluster_analysis,
'total_items': len(self.content_items),
'n_clusters': self.n_clusters,
'silhouette_score': self._calculate_silhouette_score(vectors)
}
def _extract_cluster_keywords(self, texts: List[str], cluster_id: int, top_k: int = 10) -> List[str]:
"""提取聚类关键词"""
if not texts:
return []
# 预处理文本
processed_texts = [self.preprocess_text(text) for text in texts]
# 计算TF-IDF
cluster_vectorizer = TfidfVectorizer(
max_features=1000,
stop_words=None,
ngram_range=(1, 2)
)
try:
tfidf_matrix = cluster_vectorizer.fit_transform(processed_texts)
feature_names = cluster_vectorizer.get_feature_names_out()
# 计算平均TF-IDF分数
mean_scores = np.mean(tfidf_matrix.toarray(), axis=0)
# 获取top-k关键词
top_indices = np.argsort(mean_scores)[::-1][:top_k]
keywords = [feature_names[i] for i in top_indices]
return keywords
except Exception as e:
self.logger.error(f"提取聚类 {cluster_id} 关键词失败: {e}")
return []
def _calculate_silhouette_score(self, vectors) -> float:
"""计算轮廓系数"""
try:
from sklearn.metrics import silhouette_score
if self.cluster_labels is not None and len(set(self.cluster_labels)) > 1:
return float(silhouette_score(vectors, self.cluster_labels))
except Exception as e:
self.logger.error(f"计算轮廓系数失败: {e}")
return 0.0
# 使用示例
if __name__ == "__main__":
# 配置日志
logging.basicConfig(level=logging.INFO)
# 创建测试数据
test_items = [
ContentItem(
id="1",
title="人工智能技术发展趋势",
content="人工智能技术在近年来取得了巨大进展,深度学习、机器学习等技术日趋成熟...",
url="https://example.com/ai-trends",
timestamp=time.time()
),
ContentItem(
id="2",
title="AI技术的未来发展方向",
content="随着计算能力的提升,人工智能技术将在更多领域得到应用,包括自动驾驶、医疗诊断...",
url="https://example.com/ai-future",
timestamp=time.time()
),
ContentItem(
id="3",
title="足球世界杯精彩回顾",
content="本届世界杯为球迷们带来了众多精彩瞬间,各国球队展现了高水平的竞技状态...",
url="https://example.com/worldcup",
timestamp=time.time()
),
ContentItem(
id="4",
title="NBA总决赛激战正酣",
content="NBA总决赛进入白热化阶段,两支球队实力相当,比赛异常激烈...",
url="https://example.com/nba-finals",
timestamp=time.time()
),
ContentItem(
id="5",
title="机器学习在医疗领域的应用",
content="机器学习技术在医疗诊断、药物研发等方面展现出巨大潜力,AI辅助诊断准确率不断提升...",
url="https://example.com/ml-medical",
timestamp=time.time()
)
]
# 测试相似度检测
print("=== 相似度检测测试 ===")
similarity_detector = ContentSimilarityDetector(similarity_threshold=0.3)
for item in test_items:
similarity_detector.add_content(item)
# 查找与第一个项目相似的内容
similar_results = similarity_detector.find_similar_content(test_items[0])
print(f"与 '{test_items[0].title}' 相似的内容:")
for result in similar_results:
print(f" - {result['item'].title} (相似度: {result['similarity']:.3f}, 类型: {result['match_type']})")
# 检测重复内容
duplicate_results = similarity_detector.detect_duplicate_content()
print(f"\n重复内容检测: 发现 {duplicate_results['total_groups']} 个重复组")
# 测试聚类分析
print("\n=== 聚类分析测试 ===")
cluster_analyzer = ContentClusterAnalyzer(n_clusters=3)
cluster_analyzer.add_content_batch(test_items)
clustering_results = cluster_analyzer.perform_clustering()
print(f"聚类分析完成,共 {clustering_results['n_clusters']} 个聚类:")
for cluster in clustering_results['clusters']:
print(f"\n聚类 {cluster['cluster_id']} (大小: {cluster['size']}):")
print(f" 关键词: {', '.join(cluster['keywords'][:5])}")
print(f" 样本标题: {cluster['sample_titles']}")
5. 预测性爬虫调度
5.1 智能调度系统
# intelligent_scheduler.py
import numpy as np
import pandas as pd
from sklearn.ensemble import RandomForestRegressor
from sklearn.linear_model import LinearRegression
from sklearn.preprocessing import StandardScaler
from sklearn.model_selection import train_test_split
from sklearn.metrics import mean_squared_error, r2_score
import time
import datetime
import json
import logging
from typing import Dict, List, Any, Optional, Tuple
from dataclasses import dataclass, asdict
from collections import defaultdict, deque
import joblib
import asyncio
import threading
@dataclass
class CrawlTask:
"""爬取任务"""
task_id: str
url: str
priority: int = 1
estimated_duration: float = 60.0
retry_count: int = 0
max_retries: int = 3
created_time: float = None
scheduled_time: float = None
completed_time: float = None
status: str = "pending" # pending, running, completed, failed
metadata: Dict[str, Any] = None
def __post_init__(self):
if self.created_time is None:
self.created_time = time.time()
if self.metadata is None:
self.metadata = {}
@dataclass
class CrawlResult:
"""爬取结果"""
task_id: str
success: bool
duration: float
response_size: int
status_code: int
error_message: str = ""
timestamp: float = None
def __post_init__(self):
if self.timestamp is None:
self.timestamp = time.time()
class PredictiveScheduler:
"""预测性调度器"""
def __init__(self, max_concurrent_tasks: int = 10):
self.max_concurrent_tasks = max_concurrent_tasks
self.task_queue = deque()
self.running_tasks = {}
self.completed_tasks = []
self.failed_tasks = []
# 机器学习模型
self.duration_predictor = RandomForestRegressor(n_estimators=100, random_state=42)
self.success_predictor = RandomForestRegressor(n_estimators=100, random_state=42)
self.scaler = StandardScaler()
# 历史数据
self.historical_data = []
self.model_trained = False
self.logger = logging.getLogger(__name__)
def extract_task_features(self, task: CrawlTask) -> np.ndarray:
"""提取任务特征"""
import urllib.parse
parsed_url = urllib.parse.urlparse(task.url)
features = [
# URL特征
len(task.url),
len(parsed_url.path),
len(parsed_url.query) if parsed_url.query else 0,
1 if parsed_url.scheme == 'https' else 0,
# 任务特征
task.priority,
task.retry_count,
task.estimated_duration,
# 时间特征
datetime.datetime.fromtimestamp(task.created_time).hour,
datetime.datetime.fromtimestamp(task.created_time).weekday(),
# 域名特征
len(parsed_url.netloc),
1 if 'api' in parsed_url.path.lower() else 0,
1 if any(ext in parsed_url.path.lower() for ext in ['.json', '.xml', '.csv']) else 0,
]
return np.array(features)
def add_historical_data(self, task: CrawlTask, result: CrawlResult):
"""添加历史数据"""
features = self.extract_task_features(task)
self.historical_data.append({
'features': features,
'duration': result.duration,
'success': 1 if result.success else 0,
'response_size': result.response_size,
'status_code': result.status_code
})
# 定期重新训练模型
if len(self.historical_data) % 50 == 0:
self.train_models()
def train_models(self):
"""训练预测模型"""
if len(self.historical_data) < 10:
self.logger.warning("历史数据不足,无法训练模型")
return
# 准备训练数据
X = np.array([item['features'] for item in self.historical_data])
y_duration = np.array([item['duration'] for item in self.historical_data])
y_success = np.array([item['success'] for item in self.historical_data])
# 标准化特征
X_scaled = self.scaler.fit_transform(X)
# 训练持续时间预测模型
X_train, X_test, y_train_dur, y_test_dur = train_test_split(
X_scaled, y_duration, test_size=0.2, random_state=42
)
self.duration_predictor.fit(X_train, y_train_dur)
# 评估持续时间预测
y_pred_dur = self.duration_predictor.predict(X_test)
duration_r2 = r2_score(y_test_dur, y_pred_dur)
# 训练成功率预测模型
_, _, y_train_suc, y_test_suc = train_test_split(
X_scaled, y_success, test_size=0.2, random_state=42
)
self.success_predictor.fit(X_train, y_train_suc)
# 评估成功率预测
y_pred_suc = self.success_predictor.predict(X_test)
success_r2 = r2_score(y_test_suc, y_pred_suc)
self.model_trained = True
self.logger.info(f"模型训练完成 - 持续时间R²: {duration_r2:.3f}, 成功率R²: {success_r2:.3f}")
def predict_task_metrics(self, task: CrawlTask) -> Dict[str, float]:
"""预测任务指标"""
if not self.model_trained:
return {
'predicted_duration': task.estimated_duration,
'predicted_success_rate': 0.8,
'confidence': 0.0
}
features = self.extract_task_features(task).reshape(1, -1)
features_scaled = self.scaler.transform(features)
predicted_duration = self.duration_predictor.predict(features_scaled)[0]
predicted_success = self.success_predictor.predict(features_scaled)[0]
# 计算预测置信度(基于模型的特征重要性)
confidence = min(len(self.historical_data) / 100.0, 1.0)
return {
'predicted_duration': max(predicted_duration, 1.0), # 最小1秒
'predicted_success_rate': max(0.0, min(predicted_success, 1.0)),
'confidence': confidence
}
def calculate_task_priority_score(self, task: CrawlTask) -> float:
"""计算任务优先级分数"""
predictions = self.predict_task_metrics(task)
# 基础优先级
base_score = task.priority * 10
# 成功率加权
success_weight = predictions['predicted_success_rate'] * 20
# 持续时间惩罚(时间越长,优先级越低)
duration_penalty = min(predictions['predicted_duration'] / 60.0, 10.0)
# 重试惩罚
retry_penalty = task.retry_count * 5
# 等待时间加权(等待越久,优先级越高)
wait_time = time.time() - task.created_time
wait_bonus = min(wait_time / 3600.0, 5.0) # 最多5分
total_score = base_score + success_weight - duration_penalty - retry_penalty + wait_bonus
return max(total_score, 0.0)
def add_task(self, task: CrawlTask):
"""添加任务到队列"""
self.task_queue.append(task)
self.logger.info(f"添加任务: {task.task_id} - {task.url}")
def get_next_tasks(self, count: int = None) -> List[CrawlTask]:
"""获取下一批要执行的任务"""
if count is None:
count = self.max_concurrent_tasks - len(self.running_tasks)
if count <= 0 or not self.task_queue:
return []
# 计算所有任务的优先级分数
task_scores = []
for task in self.task_queue:
score = self.calculate_task_priority_score(task)
task_scores.append((task, score))
# 按分数排序
task_scores.sort(key=lambda x: x[1], reverse=True)
# 选择前count个任务
selected_tasks = []
for task, score in task_scores[:count]:
selected_tasks.append(task)
self.task_queue.remove(task)
task.status = "running"
task.scheduled_time = time.time()
self.running_tasks[task.task_id] = task
return selected_tasks
def complete_task(self, task_id: str, result: CrawlResult):
"""完成任务"""
if task_id in self.running_tasks:
task = self.running_tasks.pop(task_id)
task.status = "completed" if result.success else "failed"
task.completed_time = result.timestamp
if result.success:
self.completed_tasks.append(task)
else:
# 检查是否需要重试
if task.retry_count < task.max_retries:
task.retry_count += 1
task.status = "pending"
self.task_queue.append(task)
self.logger.info(f"任务重试: {task_id} (第{task.retry_count}次)")
else:
self.failed_tasks.append(task)
self.logger.error(f"任务失败: {task_id}")
# 添加到历史数据
self.add_historical_data(task, result)
def get_scheduler_stats(self) -> Dict[str, Any]:
"""获取调度器统计信息"""
total_tasks = len(self.completed_tasks) + len(self.failed_tasks) + len(self.running_tasks) + len(self.task_queue)
if total_tasks == 0:
return {
'total_tasks': 0,
'success_rate': 0.0,
'average_duration': 0.0,
'queue_size': 0,
'running_tasks': 0
}
success_rate = len(self.completed_tasks) / (len(self.completed_tasks) + len(self.failed_tasks)) if (len(self.completed_tasks) + len(self.failed_tasks)) > 0 else 0.0
completed_durations = []
for task in self.completed_tasks:
if task.completed_time and task.scheduled_time:
completed_durations.append(task.completed_time - task.scheduled_time)
average_duration = np.mean(completed_durations) if completed_durations else 0.0
return {
'total_tasks': total_tasks,
'completed_tasks': len(self.completed_tasks),
'failed_tasks': len(self.failed_tasks),
'running_tasks': len(self.running_tasks),
'queue_size': len(self.task_queue),
'success_rate': success_rate,
'average_duration': average_duration,
'model_trained': self.model_trained,
'historical_data_size': len(self.historical_data)
}
class AdaptiveCrawlManager:
"""自适应爬虫管理器"""
def __init__(self, max_concurrent_tasks: int = 10):
self.scheduler = PredictiveScheduler(max_concurrent_tasks)
self.running = False
self.worker_thread = None
self.logger = logging.getLogger(__name__)
# 性能监控
self.performance_history = deque(maxlen=100)
self.last_adjustment_time = time.time()
self.adjustment_interval = 300 # 5分钟调整一次
async def crawl_url(self, task: CrawlTask) -> CrawlResult:
"""爬取URL(模拟实现)"""
import aiohttp
start_time = time.time()
try:
async with aiohttp.ClientSession() as session:
async with session.get(task.url, timeout=30) as response:
content = await response.read()
duration = time.time() - start_time
return CrawlResult(
task_id=task.task_id,
success=response.status == 200,
duration=duration,
response_size=len(content),
status_code=response.status
)
except Exception as e:
duration = time.time() - start_time
return CrawlResult(
task_id=task.task_id,
success=False,
duration=duration,
response_size=0,
status_code=0,
error_message=str(e)
)
def monitor_performance(self):
"""监控性能并自适应调整"""
current_time = time.time()
if current_time - self.last_adjustment_time < self.adjustment_interval:
return
stats = self.scheduler.get_scheduler_stats()
# 记录性能指标
self.performance_history.append({
'timestamp': current_time,
'success_rate': stats['success_rate'],
'average_duration': stats['average_duration'],
'queue_size': stats['queue_size'],
'running_tasks': stats['running_tasks']
})
# 自适应调整
if len(self.performance_history) >= 3:
recent_performance = list(self.performance_history)[-3:]
# 分析趋势
success_rates = [p['success_rate'] for p in recent_performance]
avg_success_rate = np.mean(success_rates)
queue_sizes = [p['queue_size'] for p in recent_performance]
avg_queue_size = np.mean(queue_sizes)
# 调整并发数
if avg_success_rate > 0.9 and avg_queue_size > 10:
# 成功率高且队列积压,增加并发
new_concurrent = min(self.scheduler.max_concurrent_tasks + 2, 20)
self.scheduler.max_concurrent_tasks = new_concurrent
self.logger.info(f"增加并发数到 {new_concurrent}")
elif avg_success_rate < 0.7:
# 成功率低,减少并发
new_concurrent = max(self.scheduler.max_concurrent_tasks - 1, 1)
self.scheduler.max_concurrent_tasks = new_concurrent
self.logger.info(f"减少并发数到 {new_concurrent}")
self.last_adjustment_time = current_time
async def worker_loop(self):
"""工作循环"""
while self.running:
try:
# 获取下一批任务
tasks = self.scheduler.get_next_tasks()
if not tasks:
await asyncio.sleep(1)
continue
# 并发执行任务
crawl_coroutines = [self.crawl_url(task) for task in tasks]
results = await asyncio.gather(*crawl_coroutines, return_exceptions=True)
# 处理结果
for task, result in zip(tasks, results):
if isinstance(result, Exception):
result = CrawlResult(
task_id=task.task_id,
success=False,
duration=0.0,
response_size=0,
status_code=0,
error_message=str(result)
)
self.scheduler.complete_task(task.task_id, result)
# 性能监控
self.monitor_performance()
except Exception as e:
self.logger.error(f"工作循环错误: {e}")
await asyncio.sleep(5)
def start(self):
"""启动爬虫管理器"""
if self.running:
return
self.running = True
self.logger.info("启动自适应爬虫管理器")
# 在新线程中运行异步循环
def run_async_loop():
loop = asyncio.new_event_loop()
asyncio.set_event_loop(loop)
loop.run_until_complete(self.worker_loop())
self.worker_thread = threading.Thread(target=run_async_loop)
self.worker_thread.start()
def stop(self):
"""停止爬虫管理器"""
self.running = False
if self.worker_thread:
self.worker_thread.join()
self.logger.info("停止自适应爬虫管理器")
def add_crawl_task(self, url: str, priority: int = 1, estimated_duration: float = 60.0) -> str:
"""添加爬取任务"""
task_id = f"task_{int(time.time() * 1000)}"
task = CrawlTask(
task_id=task_id,
url=url,
priority=priority,
estimated_duration=estimated_duration
)
self.scheduler.add_task(task)
return task_id
def get_status(self) -> Dict[str, Any]:
"""获取管理器状态"""
stats = self.scheduler.get_scheduler_stats()
return {
'running': self.running,
'scheduler_stats': stats,
'performance_history_size': len(self.performance_history),
'last_adjustment_time': self.last_adjustment_time
}
# 使用示例
if __name__ == "__main__":
import asyncio
# 配置日志
logging.basicConfig(level=logging.INFO)
# 创建自适应爬虫管理器
manager = AdaptiveCrawlManager(max_concurrent_tasks=5)
# 启动管理器
manager.start()
# 添加一些测试任务
test_urls = [
"https://httpbin.org/delay/1",
"https://httpbin.org/delay/2",
"https://httpbin.org/delay/3",
"https://httpbin.org/status/200",
"https://httpbin.org/status/404",
"https://httpbin.org/json",
"https://httpbin.org/xml",
"https://httpbin.org/html"
]
for i, url in enumerate(test_urls):
task_id = manager.add_crawl_task(
url=url,
priority=random.randint(1, 5),
estimated_duration=random.uniform(30, 120)
)
print(f"添加任务 {task_id}: {url}")
# 运行一段时间
try:
time.sleep(30) # 运行30秒
# 获取状态
status = manager.get_status()
print(f"\n管理器状态: {json.dumps(status, indent=2)}")
finally:
# 停止管理器
manager.stop()
6. 实战案例:智能新闻爬虫
6.1 综合应用示例
# intelligent_news_crawler.py
import asyncio
import aiohttp
from bs4 import BeautifulSoup
import time
import json
import logging
from typing import Dict, List, Any, Optional
from dataclasses import dataclass, asdict
import hashlib
# 导入之前定义的类
from text_classifier import ContentCategoryClassifier
from similarity_detection import ContentSimilarityDetector, ContentItem
from intelligent_scheduler import AdaptiveCrawlManager, CrawlTask
@dataclass
class NewsArticle:
"""新闻文章"""
title: str
content: str
url: str
category: str
publish_time: str
author: str = ""
summary: str = ""
keywords: List[str] = None
similarity_score: float = 0.0
def __post_init__(self):
if self.keywords is None:
self.keywords = []
class IntelligentNewsCrawler:
"""智能新闻爬虫"""
def __init__(self):
self.category_classifier = ContentCategoryClassifier()
self.similarity_detector = ContentSimilarityDetector(similarity_threshold=0.85)
self.crawl_manager = AdaptiveCrawlManager(max_concurrent_tasks=8)
self.articles = []
self.processed_urls = set()
self.logger = logging.getLogger(__name__)
# 初始化分类器
self.category_classifier.train_with_sample_data()
async def extract_article_content(self, html_content: str, url: str) -> Optional[NewsArticle]:
"""提取文章内容"""
try:
soup = BeautifulSoup(html_content, 'html.parser')
# 提取标题
title_selectors = ['h1', 'title', '.title', '.headline', '[class*="title"]']
title = ""
for selector in title_selectors:
title_elem = soup.select_one(selector)
if title_elem:
title = title_elem.get_text().strip()
break
# 提取正文内容
content_selectors = [
'.content', '.article-content', '.post-content',
'[class*="content"]', 'article', '.article-body'
]
content = ""
for selector in content_selectors:
content_elem = soup.select_one(selector)
if content_elem:
# 移除脚本和样式
for script in content_elem(["script", "style"]):
script.decompose()
content = content_elem.get_text().strip()
break
# 如果没有找到专门的内容区域,使用整个body
if not content:
body = soup.find('body')
if body:
for script in body(["script", "style", "nav", "header", "footer"]):
script.decompose()
content = body.get_text().strip()
# 提取发布时间
time_selectors = [
'[class*="time"]', '[class*="date"]', 'time',
'.publish-time', '.post-date'
]
publish_time = ""
for selector in time_selectors:
time_elem = soup.select_one(selector)
if time_elem:
publish_time = time_elem.get_text().strip()
break
# 提取作者
author_selectors = [
'.author', '[class*="author"]', '.byline', '.writer'
]
author = ""
for selector in author_selectors:
author_elem = soup.select_one(selector)
if author_elem:
author = author_elem.get_text().strip()
break
if not title or not content or len(content) < 100:
return None
# 使用机器学习分类内容
classification_result = self.category_classifier.classify_content(html_content, url)
category = classification_result.get('category', 'other')
# 创建文章对象
article = NewsArticle(
title=title,
content=content,
url=url,
category=category,
publish_time=publish_time,
author=author,
summary=content[:200] + "..." if len(content) > 200 else content
)
return article
except Exception as e:
self.logger.error(f"提取文章内容失败 {url}: {e}")
return None
def check_content_similarity(self, article: NewsArticle) -> bool:
"""检查内容相似度"""
content_item = ContentItem(
id=hashlib.md5(article.url.encode()).hexdigest(),
title=article.title,
content=article.content,
url=article.url,
timestamp=time.time()
)
# 查找相似内容
similar_results = self.similarity_detector.find_similar_content(content_item, top_k=3)
if similar_results:
max_similarity = max(result['similarity'] for result in similar_results)
article.similarity_score = max_similarity
# 如果相似度过高,认为是重复内容
if max_similarity > 0.85:
self.logger.info(f"发现重复内容: {article.title} (相似度: {max_similarity:.3f})")
return True
# 添加到相似度检测器
self.similarity_detector.add_content(content_item)
return False
async def crawl_news_site(self, base_url: str, max_pages: int = 5):
"""爬取新闻网站"""
self.logger.info(f"开始爬取新闻网站: {base_url}")
# 启动爬虫管理器
self.crawl_manager.start()
try:
# 添加初始页面到任务队列
for page in range(1, max_pages + 1):
page_url = f"{base_url}?page={page}"
if page_url not in self.processed_urls:
self.crawl_manager.add_crawl_task(
url=page_url,
priority=5, # 列表页优先级较高
estimated_duration=30.0
)
self.processed_urls.add(page_url)
# 等待任务完成
processed_count = 0
while processed_count < max_pages * 10: # 假设每页最多10篇文章
await asyncio.sleep(2)
status = self.crawl_manager.get_status()
if status['scheduler_stats']['queue_size'] == 0 and status['scheduler_stats']['running_tasks'] == 0:
break
processed_count += 1
# 处理已完成的任务(这里需要实际的结果处理逻辑)
# 在实际应用中,你需要从crawl_manager获取结果并处理
finally:
self.crawl_manager.stop()
def generate_crawl_report(self) -> Dict[str, Any]:
"""生成爬取报告"""
if not self.articles:
return {'total_articles': 0}
# 统计分类分布
category_distribution = {}
for article in self.articles:
category_distribution[article.category] = category_distribution.get(article.category, 0) + 1
# 统计相似度分布
similarity_scores = [article.similarity_score for article in self.articles]
avg_similarity = sum(similarity_scores) / len(similarity_scores) if similarity_scores else 0
# 统计重复内容
duplicate_count = sum(1 for article in self.articles if article.similarity_score > 0.85)
return {
'total_articles': len(self.articles),
'category_distribution': category_distribution,
'average_similarity_score': avg_similarity,
'duplicate_articles': duplicate_count,
'unique_articles': len(self.articles) - duplicate_count,
'processed_urls': len(self.processed_urls),
'crawl_manager_stats': self.crawl_manager.get_status()
}
def save_articles(self, filename: str):
"""保存文章到文件"""
articles_data = [asdict(article) for article in self.articles]
with open(filename, 'w', encoding='utf-8') as f:
json.dump(articles_data, f, ensure_ascii=False, indent=2)
self.logger.info(f"保存了 {len(self.articles)} 篇文章到 {filename}")
# 使用示例
async def main():
# 配置日志
logging.basicConfig(level=logging.INFO)
# 创建智能新闻爬虫
crawler = IntelligentNewsCrawler()
# 爬取新闻网站
await crawler.crawl_news_site("https://news.example.com", max_pages=3)
# 生成报告
report = crawler.generate_crawl_report()
print(f"爬取报告: {json.dumps(report, indent=2, ensure_ascii=False)}")
# 保存文章
crawler.save_articles("crawled_news.json")
if __name__ == "__main__":
asyncio.run(main())
7. 练习与作业
7.1 基础练习
-
文本分类器训练
- 收集不同类别的网页内容
- 训练自己的文本分类模型
- 评估模型性能并优化
-
图像识别应用
- 实现网页图像自动分类
- 集成验证码识别功能
- 处理不同格式的图像
-
行为模式分析
- 分析爬虫访问模式
- 实现异常检测算法
- 设计反检测策略
7.2 高级练习
-
相似度检测系统
- 构建大规模内容去重系统
- 实现增量相似度计算
- 优化检测性能
-
智能调度器
- 设计多目标优化调度算法
- 实现动态负载均衡
- 集成机器学习预测
7.3 实战项目
-
智能新闻聚合器
- 多源新闻采集
- 自动分类和去重
- 热点话题识别
-
电商价格监控系统
- 商品信息智能提取
- 价格趋势预测
- 异常价格检测
8. 常见问题与解决方案
8.1 模型训练问题
- 数据不足:使用数据增强、迁移学习
- 过拟合:正则化、交叉验证
- 性能优化:特征选择、模型集成
8.2 实时处理挑战
- 延迟问题:异步处理、缓存优化
- 内存管理:批处理、流式计算
- 扩展性:分布式部署、负载均衡
8.3 准确性保证
- 模型验证:A/B测试、在线学习
- 数据质量:异常检测、数据清洗
- 持续优化:模型更新、反馈循环
9. 总结与展望
本课程介绍了机器学习在爬虫中的多种应用:
- 智能内容识别:自动分类和理解网页内容
- 反爬虫对抗:检测和规避反爬虫机制
- 数据挖掘:发现内容模式和相似性
- 预测调度:优化爬虫执行策略
- 综合应用:构建智能化爬虫系统
机器学习为爬虫技术带来了新的可能性,使爬虫能够更智能地理解和处理数据,提高效率和准确性。
下节预告
下一课我们将学习**《Python爬虫第13课:爬虫性能优化与监控》**,内容包括:
- 性能瓶颈分析与优化
- 系统监控与告警
- 资源使用优化
- 高可用架构设计
结尾
希望对初学者有帮助;致力于办公自动化的小小程序员一枚
希望能得到大家的【❤️一个免费关注❤️】感谢!
求个 🤞 关注 🤞 +❤️ 喜欢 ❤️ +👍 收藏 👍
此外还有办公自动化专栏,欢迎大家订阅:Python办公自动化专栏
此外还有爬虫专栏,欢迎大家订阅:Python爬虫基础专栏
此外还有Python基础专栏,欢迎大家订阅:Python基础学习专栏
更多推荐

所有评论(0)