引言
在当今数据驱动的科研时代,高质量的科学实验数据集是研究人员进行数据分析、模型训练和科学发现的基础。然而,获取这些数据集往往面临诸多挑战:数据分散在不同网站、格式不统一、需要定期更新、手动下载效率低下等。Python爬虫技术为我们提供了自动化采集网络数据的强大工具,能够显著提升科研工作的效率。
本文将深入探讨如何使用Python爬虫技术从各类科研网站、数据仓库和开放数据库中采集科学实验数据。我们将从基础概念讲起,逐步深入到高级应用,包括动态页面抓取、反爬虫策略应对、分布式爬虫架构等最新技术,并辅以完整的代码示例。无论你是刚接触爬虫的科研人员,还是希望提升数据采集技能的数据科学家,本文都将为你提供有价值的参考。
第一章:Python爬虫基础与科学数据采集概述
1.1 什么是网络爬虫?
网络爬虫(Web Crawler)是一种按照一定规则自动抓取万维网信息的程序或脚本。在科学数据采集场景下,爬虫可以帮助我们自动化地从各类科研数据库、期刊网站、开放数据平台中提取实验数据,大大减少人工收集的工作量。
1.2 科学数据采集的特殊性
与通用爬虫相比,科学数据采集面临一些特殊挑战:
-
数据质量要求高:科研数据对准确性、完整性的要求极高
-
数据结构复杂:可能包含多维数组、时间序列、图像等多种格式
-
采集频率特殊:有些数据集只需一次性采集,有些需要定期更新
-
伦理与合规性:必须遵守数据使用条款和版权规定
1.3 技术栈选择
现代Python爬虫技术栈通常包括:
-
请求库:requests, aiohttp(异步)
-
解析库:BeautifulSoup, lxml, pyquery
-
动态渲染:Selenium, Playwright, Puppeteer
-
框架:Scrapy, pyspider
-
存储:pandas, SQLAlchemy, MongoDB
-
反爬虫应对:代理IP池,请求头伪装,验证码识别
第二章:环境搭建与基础工具介绍
2.1 Python环境配置
首先确保安装了Python 3.8+版本,建议使用虚拟环境管理项目依赖:
bash
# 创建虚拟环境
python -m venv science_crawler_env
# 激活虚拟环境(Windows)
science_crawler_env\\Scripts\\activate
# 激活虚拟环境(Mac/Linux)
source science_crawler_env/bin/activate
# 升级pip
pip install –upgrade pip
2.2 核心库安装
bash
# 基础爬虫库
pip install requests beautifulsoup4 lxml
# 数据解析与处理
pip install pandas numpy
# 动态页面爬取
pip install selenium playwright
playwright install
# 高级框架
pip install scrapy
# 数据库驱动
pip install pymongo psycopg2-binary
# 其他实用工具
pip install retry fake-useragent
2.3 开发工具推荐
-
IDE:VS Code + Python插件,或PyCharm专业版
-
调试工具:浏览器开发者工具(F12),Postman
-
代理工具:Charles,Fiddler(用于抓包分析)
第三章:静态网页数据采集——以NCBI PubMed为例
3.1 网站分析
PubMed是美国国立医学图书馆开发的生物医学文献检索系统,包含超过3000万篇生物医学文献的引用和摘要。我们将演示如何从PubMed采集文献摘要数据。
首先分析目标URL结构:
text
https://pubmed.ncbi.nlm.nih.gov/?term=cancer+genomics&size=200
3.2 基础请求与响应处理
python
import requests
from bs4 import BeautifulSoup
import pandas as pd
import time
from fake_useragent import UserAgent
from retry import retry
import logging
# 配置日志
logging.basicConfig(level=logging.INFO, format='%(asctime)s – %(levelname)s – %(message)s')
class PubMedCrawler:
def __init__(self):
self.base_url = "https://pubmed.ncbi.nlm.nih.gov/"
self.ua = UserAgent()
self.session = requests.Session()
def get_random_headers(self):
"""生成随机请求头"""
return {
'User-Agent': self.ua.random,
'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, br',
'Connection': 'keep-alive',
'Upgrade-Insecure-Requests': '1',
}
@retry(tries=3, delay=2, backoff=2)
def fetch_page(self, url, params=None):
"""获取页面内容,带重试机制"""
try:
response = self.session.get(
url,
params=params,
headers=self.get_random_headers(),
timeout=30
)
response.raise_for_status()
# 检查是否被重定向到验证页面
if "captcha" in response.url.lower():
raise Exception("遇到验证码页面")
return response.text
except requests.RequestException as e:
logging.error(f"请求失败: {e}")
raise
def parse_search_results(self, html):
"""解析搜索结果页面"""
soup = BeautifulSoup(html, 'lxml')
articles = []
# 查找所有文献条目
for article in soup.find_all('article', class_='full-docsum'):
try:
# 提取PMID(唯一标识符)
pmid_elem = article.find('a', class_='docsum-title')
pmid = pmid_elem.get('href', '').split('/')[-2] if pmid_elem else None
# 提取标题
title_elem = article.find('a', class_='docsum-title')
title = title_elem.text.strip() if title_elem else None
# 提取作者信息
author_elem = article.find('span', class_='docsum-authors')
authors = author_elem.text.strip() if author_elem else None
# 提取期刊信息
journal_elem = article.find('span', class_='docsum-journal-citation')
journal = journal_elem.text.strip() if journal_elem else None
# 提取发表日期
date_elem = article.find('span', class_='docsum-pubdate')
pub_date = date_elem.text.strip() if date_elem else None
# 提取摘要预览
abstract_elem = article.find('div', class_='full-view-snippet')
abstract_preview = abstract_elem.text.strip() if abstract_elem else None
article_data = {
'pmid': pmid,
'title': title,
'authors': authors,
'journal': journal,
'publication_date': pub_date,
'abstract_preview': abstract_preview,
'url': f"https://pubmed.ncbi.nlm.nih.gov/{pmid}/" if pmid else None
}
articles.append(article_data)
logging.info(f"解析到文献: {pmid} – {title[:50]}…")
except Exception as e:
logging.error(f"解析文献条目时出错: {e}")
continue
return articles
def crawl(self, keyword, max_pages=10):
"""主爬取方法"""
all_articles = []
params = {
'term': keyword,
'size': 200, # 每页结果数
'page': 1
}
for page in range(1, max_pages + 1):
logging.info(f"正在爬取第 {page} 页,关键词: {keyword}")
params['page'] = page
try:
html = self.fetch_page(self.base_url, params)
articles = self.parse_search_results(html)
if not articles:
logging.info("没有更多数据,停止爬取")
break
all_articles.extend(articles)
# 礼貌性延迟
time.sleep(2)
except Exception as e:
logging.error(f"爬取第 {page} 页失败: {e}")
continue
return all_articles
def save_to_csv(self, data, filename):
"""保存数据到CSV文件"""
df = pd.DataFrame(data)
df.to_csv(filename, index=False, encoding='utf-8-sig')
logging.info(f"数据已保存到 {filename},共 {len(df)} 条记录")
# 使用示例
if __name__ == "__main__":
crawler = PubMedCrawler()
# 爬取癌症基因组学相关文献
articles = crawler.crawl(keyword="cancer+genomics", max_pages=5)
crawler.save_to_csv(articles, "pubmed_cancer_genomics.csv")
3.3 高级功能:增量更新与去重
python
import hashlib
import json
from pathlib import Path
class IncrementalCrawler(PubMedCrawler):
def __init__(self, cache_dir="cache"):
super().__init__()
self.cache_dir = Path(cache_dir)
self.cache_dir.mkdir(exist_ok=True)
self.seen_pmids = self.load_seen_pmids()
def load_seen_pmids(self):
"""加载已爬取的PMID"""
seen_file = self.cache_dir / "seen_pmids.json"
if seen_file.exists():
with open(seen_file, 'r') as f:
return set(json.load(f))
return set()
def save_seen_pmids(self):
"""保存已爬取的PMID"""
seen_file = self.cache_dir / "seen_pmids.json"
with open(seen_file, 'w') as f:
json.dump(list(self.seen_pmids), f)
def is_duplicate(self, pmid):
"""检查是否重复"""
return pmid in self.seen_pmids
def crawl_incremental(self, keyword, max_pages=10):
"""增量爬取"""
new_articles = []
params = {
'term': keyword,
'size': 200,
'page': 1
}
for page in range(1, max_pages + 1):
logging.info(f"增量爬取第 {page} 页")
params['page'] = page
try:
html = self.fetch_page(self.base_url, params)
articles = self.parse_search_results(html)
if not articles:
break
# 过滤重复数据
for article in articles:
pmid = article.get('pmid')
if pmid and not self.is_duplicate(pmid):
new_articles.append(article)
self.seen_pmids.add(pmid)
time.sleep(2)
except Exception as e:
logging.error(f"增量爬取失败: {e}")
continue
# 保存更新后的seen集合
self.save_seen_pmids()
return new_articles
第四章:动态网页数据采集——以全球气候数据网站为例
4.1 Selenium基础应用
许多科学数据网站使用JavaScript动态加载数据,传统requests库无法获取。我们以NASA气候数据网站为例,演示如何使用Selenium采集动态内容。
python
from selenium import webdriver
from selenium.webdriver.common.by import By
from selenium.webdriver.support.ui import WebDriverWait
from selenium.webdriver.support import expected_conditions as EC
from selenium.common.exceptions import TimeoutException, NoSuchElementException
from selenium.webdriver.chrome.options import Options
import pandas as pd
import time
import logging
class ClimateDataCrawler:
def __init__(self, headless=True):
self.logger = logging.getLogger(__name__)
self.driver = self.init_driver(headless)
self.wait = WebDriverWait(self.driver, 10)
def init_driver(self, headless):
"""初始化Chrome驱动"""
chrome_options = Options()
if headless:
chrome_options.add_argument('–headless')
chrome_options.add_argument('–no-sandbox')
chrome_options.add_argument('–disable-dev-shm-usage')
chrome_options.add_argument('–disable-gpu')
chrome_options.add_argument('–window-size=1920,1080')
chrome_options.add_argument('–disable-blink-features=AutomationControlled')
chrome_options.add_experimental_option("excludeSwitches", ["enable-automation"])
chrome_options.add_experimental_option('useAutomationExtension', False)
# 设置用户代理
chrome_options.add_argument('user-agent=Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36')
driver = webdriver.Chrome(options=chrome_options)
# 隐藏WebDriver特征
driver.execute_script("Object.defineProperty(navigator, 'webdriver', {get: () => undefined})")
return driver
def wait_for_element(self, by, value, timeout=10):
"""等待元素出现"""
try:
element = WebDriverWait(self.driver, timeout).until(
EC.presence_of_element_located((by, value))
)
return element
except TimeoutException:
self.logger.error(f"等待元素超时: {by}={value}")
return None
def scroll_to_bottom(self):
"""滚动到页面底部(用于触发懒加载)"""
last_height = self.driver.execute_script("return document.body.scrollHeight")
while True:
self.driver.execute_script("window.scrollTo(0, document.body.scrollHeight);")
time.sleep(2)
new_height = self.driver.execute_script("return document.body.scrollHeight")
if new_height == last_height:
break
last_height = new_height
def extract_table_data(self):
"""提取表格数据"""
data = []
try:
# 找到数据表格
table = self.wait_for_element(By.TAG_NAME, "table")
if not table:
return data
# 提取表头
headers = []
header_row = table.find_element(By.TAG_NAME, "thead")
if header_row:
headers = [th.text.strip() for th in header_row.find_elements(By.TAG_NAME, "th")]
# 提取数据行
tbody = table.find_element(By.TAG_NAME, "tbody")
rows = tbody.find_elements(By.TAG_NAME, "tr")
for row in rows:
cells = row.find_elements(By.TAG_NAME, "td")
if len(cells) == len(headers):
row_data = {}
for i, cell in enumerate(cells):
row_data[headers[i]] = cell.text.strip()
data.append(row_data)
except Exception as e:
self.logger.error(f"提取表格数据失败: {e}")
return data
def click_load_more(self):
"""点击"加载更多"按钮"""
try:
load_more = self.driver.find_element(By.XPATH, "//button[contains(text(), 'Load More')]")
self.driver.execute_script("arguments[0].click();", load_more)
time.sleep(2)
return True
except NoSuchElementException:
return False
def crawl_nasa_climate(self, url="https://data.giss.nasa.gov/gistemp/"):
"""爬取NASA气候数据"""
self.logger.info(f"正在访问: {url}")
self.driver.get(url)
# 等待页面加载
self.wait_for_element(By.TAG_NAME, "body")
# 处理可能的cookie同意弹窗
try:
cookie_btn = self.driver.find_element(By.XPATH, "//button[contains(text(), 'Accept')]")
cookie_btn.click()
except:
pass
# 选择数据参数
try:
# 选择年份范围
year_from = self.driver.find_element(By.NAME, "year_from")
year_from.clear()
year_from.send_keys("1880")
year_to = self.driver.find_element(By.NAME, "year_to")
year_to.clear()
year_to.send_keys("2023")
# 提交表单
submit_btn = self.driver.find_element(By.XPATH, "//input[@type='submit']")
submit_btn.click()
# 等待结果加载
time.sleep(5)
except Exception as e:
self.logger.error(f"设置参数失败: {e}")
# 提取数据
climate_data = self.extract_table_data()
# 尝试加载更多数据(如果有分页)
while self.click_load_more():
new_data = self.extract_table_data()
climate_data.extend(new_data)
return climate_data
def close(self):
"""关闭浏览器"""
if self.driver:
self.driver.quit()
# 使用示例
if __name__ == "__main__":
logging.basicConfig(level=logging.INFO)
crawler = ClimateDataCrawler(headless=False) # 设为False可以看到浏览器操作
try:
data = crawler.crawl_nasa_climate()
df = pd.DataFrame(data)
df.to_csv("nasa_climate_data.csv", index=False)
logging.info(f"成功采集 {len(df)} 条气候数据记录")
finally:
crawler.close()
4.2 更现代的方案:Playwright异步爬虫
Playwright是微软开发的现代浏览器自动化工具,支持异步操作,性能更好:
python
import asyncio
from playwright.async_api import async_playwright
import pandas as pd
import logging
from typing import List, Dict
class AsyncClimateCrawler:
def __init__(self, headless=True):
self.logger = logging.getLogger(__name__)
self.headless = headless
async def crawl_page(self, page, url):
"""单个页面的爬取逻辑"""
await page.goto(url, wait_until="networkidle")
# 等待数据表格加载
await page.wait_for_selector("table", state="attached")
# 模拟人类行为
await page.mouse.move(100, 100)
await page.mouse.wheel(delta_y=500)
await asyncio.sleep(1)
# 提取数据
data = await page.evaluate("""
() => {
const rows = Array.from(document.querySelectorAll('table tbody tr'));
return rows.map(row => {
const cells = row.querySelectorAll('td');
return Array.from(cells).map(cell => cell.textContent.trim());
});
}
""")
# 提取表头
headers = await page.evaluate("""
() => {
const headers = Array.from(document.querySelectorAll('table thead th'));
return headers.map(h => h.textContent.trim());
}
""")
return headers, data
async def crawl_multiple_pages(self, urls: List[str]) -> List[Dict]:
"""并发爬取多个页面"""
async with async_playwright() as p:
browser = await p.chromium.launch(headless=self.headless)
context = await browser.new_context(
viewport={'width': 1920, 'height': 1080},
user_agent='Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36'
)
# 创建并发任务
tasks = []
for url in urls:
page = await context.new_page()
tasks.append(self.crawl_page(page, url))
# 并发执行
results = await asyncio.gather(*tasks, return_exceptions=True)
await browser.close()
# 整理结果
all_data = []
for result in results:
if isinstance(result, Exception):
self.logger.error(f"爬取失败: {result}")
else:
headers, data_rows = result
for row in data_rows:
if len(row) == len(headers):
item = dict(zip(headers, row))
all_data.append(item)
return all_data
async def run(self, urls):
"""运行爬虫"""
return await self.crawl_multiple_pages(urls)
# 使用示例
async def main():
crawler = AsyncClimateCrawler()
urls = [
"https://data.giss.nasa.gov/gistemp/",
# 可以添加更多URL
]
data = await crawler.run(urls)
df = pd.DataFrame(data)
df.to_csv("climate_data_async.csv", index=False)
print(f"成功采集 {len(df)} 条记录")
if __name__ == "__main__":
asyncio.run(main())
第五章:大规模数据采集——Scrapy框架实战
对于需要采集大量科学数据的项目,Scrapy框架提供了更强大的功能和更好的性能。
5.1 Scrapy项目初始化
bash
# 安装Scrapy
pip install scrapy
# 创建项目
scrapy startproject science_data_crawler
cd science_data_crawler
# 创建爬虫
scrapy genspider pubmed_spider pubmed.ncbi.nlm.nih.gov
5.2 完整的Scrapy爬虫实现
python
# spiders/pubmed_spider.py
import scrapy
from scrapy.linkextractors import LinkExtractor
from scrapy.spiders import CrawlSpider, Rule
from science_data_crawler.items import PublicationItem
from scrapy.loader import ItemLoader
from scrapy.loader.processors import MapCompose, Join
import re
def clean_text(text):
"""清理文本"""
if text:
return re.sub(r'\\s+', ' ', text.strip())
return ''
def extract_pmid(url):
"""从URL中提取PMID"""
match = re.search(r'/(\\d+)/?$', url)
return match.group(1) if match else None
class PubMedSpider(CrawlSpider):
name = 'pubmed'
allowed_domains = ['pubmed.ncbi.nlm.nih.gov']
start_urls = ['https://pubmed.ncbi.nlm.nih.gov/?term=cancer+genomics']
# 定义爬取规则
rules = (
# 跟踪分页链接
Rule(LinkExtractor(allow=r'\\?term=cancer\\+genomics&page=\\d+'), follow=True),
# 提取详情页链接
Rule(LinkExtractor(allow=r'/\\d+/$'), callback='parse_article'),
)
def parse_article(self, response):
"""解析文献详情页"""
loader = ItemLoader(item=PublicationItem(), response=response)
# 添加基本信息
loader.add_value('url', response.url)
loader.add_value('pmid', extract_pmid(response.url))
# 提取标题
loader.add_css('title', 'h1.heading-title::text', MapCompose(clean_text))
# 提取作者
loader.add_css('authors', 'div.authors-list span.author-name::text', MapCompose(clean_text))
# 提取摘要
loader.add_css('abstract', 'div.abstract-content p::text', MapCompose(clean_text), Join('\\n'))
# 提取期刊信息
loader.add_css('journal', 'div.journal-info button.context-journal-trigger::text', MapCompose(clean_text))
# 提取发表日期
loader.add_css('publication_date', 'span.citation-part.date::text', MapCompose(clean_text))
# 提取DOI
loader.add_css('doi', 'span.citation-doi a::text', MapCompose(clean_text))
# 提取关键词
loader.add_css('keywords', 'div.keywords-list a.keyword-link::text', MapCompose(clean_text))
return loader.load_item()
5.3 Item定义与数据处理
python
# items.py
import scrapy
from scrapy.item import Field
class PublicationItem(scrapy.Item):
"""文献数据项"""
url = Field()
pmid = Field()
title = Field()
authors = Field()
abstract = Field()
journal = Field()
publication_date = Field()
doi = Field()
keywords = Field()
crawl_time = Field()
# pipelines.py
import pymongo
import logging
from scrapy.exceptions import DropItem
from datetime import datetime
class MongoPipeline:
"""MongoDB存储管道"""
def __init__(self, mongo_uri, mongo_db):
self.mongo_uri = mongo_uri
self.mongo_db = mongo_db
self.logger = logging.getLogger(__name__)
@classmethod
def from_crawler(cls, crawler):
return cls(
mongo_uri=crawler.settings.get('MONGO_URI'),
mongo_db=crawler.settings.get('MONGO_DATABASE')
)
def open_spider(self, spider):
self.client = pymongo.MongoClient(self.mongo_uri)
self.db = self.client[self.mongo_db]
# 创建索引
self.db.publications.create_index('pmid', unique=True)
def close_spider(self, spider):
self.client.close()
def process_item(self, item, spider):
# 添加爬取时间
item['crawl_time'] = datetime.utcnow()
# 检查必填字段
if not item.get('pmid'):
raise DropItem("缺少PMID字段")
# 去重
existing = self.db.publications.find_one({'pmid': item['pmid']})
if existing:
raise DropItem(f"重复文献: {item['pmid']}")
# 插入数据库
self.db.publications.insert_one(dict(item))
self.logger.info(f"保存文献: {item['pmid']}")
return item
class DuplicatesPipeline:
"""内存去重管道"""
def __init__(self):
self.seen_pmids = set()
def process_item(self, item, spider):
pmid = item.get('pmid')
if pmid in self.seen_pmids:
raise DropItem(f"重复PMID: {pmid}")
self.seen_pmids.add(pmid)
return item
# settings.py
# 启用管道
ITEM_PIPELINES = {
'science_data_crawler.pipelines.DuplicatesPipeline': 300,
'science_data_crawler.pipelines.MongoPipeline': 800,
}
# MongoDB配置
MONGO_URI = 'mongodb://localhost:27017'
MONGO_DATABASE = 'science_data'
# 下载延迟
DOWNLOAD_DELAY = 1
RANDOMIZE_DOWNLOAD_DELAY = True
# 并发设置
CONCURRENT_REQUESTS = 16
CONCURRENT_REQUESTS_PER_DOMAIN = 8
# 自动限速
AUTOTHROTTLE_ENABLED = True
AUTOTHROTTLE_START_DELAY = 1.0
AUTOTHROTTLE_MAX_DELAY = 60.0
# 重试设置
RETRY_ENABLED = True
RETRY_TIMES = 3
RETRY_HTTP_CODES = [500, 502, 503, 504, 408, 429]
# 请求头
DEFAULT_REQUEST_HEADERS = {
'Accept': 'text/html,application/xhtml+xml,application/xml;q=0.9,*/*;q=0.8',
'Accept-Language': 'en',
'User-Agent': 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36'
}
5.4 中间件实现
python
# middlewares.py
import random
from scrapy import signals
from scrapy.downloadermiddlewares.retry import RetryMiddleware
from scrapy.utils.response import response_status_message
import logging
class RandomUserAgentMiddleware:
"""随机User-Agent中间件"""
def __init__(self, user_agents):
self.user_agents = user_agents
@classmethod
def from_crawler(cls, crawler):
return cls(
user_agents=crawler.settings.get('USER_AGENTS', [])
)
def process_request(self, request, spider):
request.headers.setdefault('User-Agent', random.choice(self.user_agents))
class ProxyMiddleware:
"""代理中间件"""
def __init__(self, proxy_list):
self.proxy_list = proxy_list
@classmethod
def from_crawler(cls, crawler):
return cls(
proxy_list=crawler.settings.get('PROXY_LIST', [])
)
def process_request(self, request, spider):
if self.proxy_list:
request.meta['proxy'] = random.choice(self.proxy_list)
class TooManyRequestsRetryMiddleware(RetryMiddleware):
"""处理429状态码的重试中间件"""
def __init__(self, crawler):
super().__init__(crawler.settings)
self.crawler = crawler
@classmethod
def from_crawler(cls, crawler):
return cls(crawler)
def process_response(self, request, response, spider):
if response.status == 429:
retryreq = self._retry(request, '429', spider)
if retryreq:
self.logger.info(f"遇到429,正在重试: {request.url}")
return retryreq
return response
5.5 运行Scrapy爬虫
bash
# 运行爬虫并保存为JSON
scrapy crawl pubmed -o publications.json
# 运行爬虫并保存为CSV
scrapy crawl pubmed -o publications.csv
# 使用日志级别
scrapy crawl pubmed -o publications.json -s LOG_LEVEL=INFO
# 指定设置文件
scrapy crawl pubmed -o publications.json –set FEED_EXPORT_ENCODING=utf-8
第六章:反爬虫策略与应对技巧
6.1 常见的反爬虫机制
科学数据网站通常采用以下反爬虫措施:
-
IP封禁:检测短时间内大量请求
-
User-Agent检测:识别非浏览器请求
-
验证码:包括文字验证码、滑动验证码
-
JavaScript渲染:数据动态加载
-
请求频率限制:通过API速率限制
-
Cookie验证:需要维持会话
6.2 应对策略实现
python
# anti_anti_spider.py
import random
import time
from collections import defaultdict
from datetime import datetime, timedelta
import requests
from fake_useragent import UserAgent
class AntiAntiSpider:
"""反反爬虫综合策略"""
def __init__(self):
self.ua = UserAgent()
self.proxy_pool = []
self.request_timestamps = defaultdict(list)
self.max_requests_per_minute = 30
self.max_requests_per_ip = 100
def get_random_delay(self, base_delay=2):
"""获取随机延迟时间"""
return base_delay + random.uniform(0.5, 2.0)
def get_rotating_user_agent(self):
"""获取轮换的User-Agent"""
return self.ua.random
def get_proxy(self):
"""获取代理IP(示例使用免费代理API)"""
try:
# 从代理API获取
response = requests.get(
'http://api.proxy.com/get',
params={'num': 1, 'protocol': 'http'},
timeout=5
)
if response.status_code == 200:
proxy = response.json().get('proxy')
return {'http': f'http://{proxy}', 'https': f'http://{proxy}'}
except:
pass
# 返回本地直连
return None
def check_rate_limit(self, ip):
"""检查IP请求频率"""
now = datetime.now()
one_minute_ago = now – timedelta(minutes=1)
# 清理旧记录
self.request_timestamps[ip] = [
ts for ts in self.request_timestamps[ip]
if ts > one_minute_ago
]
# 检查频率
if len(self.request_timestamps[ip]) >= self.max_requests_per_minute:
wait_time = 60 – (now – self.request_timestamps[ip][0]).seconds
time.sleep(wait_time)
self.request_timestamps[ip].append(now)
def simulate_human_behavior(self, session):
"""模拟人类行为"""
# 添加常见请求头
session.headers.update({
'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, br',
'Connection': 'keep-alive',
'Upgrade-Insecure-Requests': '1',
'Sec-Fetch-Dest': 'document',
'Sec-Fetch-Mode': 'navigate',
'Sec-Fetch-Site': 'none',
'Sec-Fetch-User': '?1',
'Cache-Control': 'max-age=0',
})
# 随机添加一些额外头
if random.random() > 0.5:
session.headers['DNT'] = '1'
def retry_with_backoff(self, func, max_retries=5):
"""指数退避重试"""
for i in range(max_retries):
try:
return func()
except Exception as e:
if i == max_retries – 1:
raise
wait_time = (2 ** i) + random.uniform(0, 1)
print(f"请求失败,{wait_time:.2f}秒后重试…")
time.sleep(wait_time)
# 使用示例
def create_robust_session():
"""创建健壮的请求会话"""
anti = AntiAntiSpider()
session = requests.Session()
# 配置会话
adapter = requests.adapters.HTTPAdapter(
pool_connections=20,
pool_maxsize=20,
max_retries=3,
pool_block=False
)
session.mount('http://', adapter)
session.mount('https://', adapter)
# 设置User-Agent
session.headers.update({
'User-Agent': anti.get_rotating_user_agent()
})
# 模拟人类行为
anti.simulate_human_behavior(session)
return session, anti
6.3 验证码识别方案
python
# captcha_solver.py
import base64
import requests
from io import BytesIO
from PIL import Image
import pytesseract
import ddddocr # 通用验证码识别库
class CaptchaSolver:
"""验证码识别器"""
def __init__(self):
# 初始化OCR
self.ocr = ddddocr.DdddOcr(show_ad=False)
# 初始化Tesseract(备选)
pytesseract.pytesseract.tesseract_cmd = r'C:\\Program Files\\Tesseract-OCR\\tesseract.exe'
def solve_image_captcha(self, image_data):
"""识别图片验证码"""
try:
# 使用ddddocr
result = self.ocr.classification(image_data)
return result
except:
# 备选方案
image = Image.open(BytesIO(image_data))
# 图片预处理
image = image.convert('L') # 灰度化
image = image.point(lambda x: 0 if x < 128 else 255) # 二值化
result = pytesseract.image_to_string(image, config='–psm 8')
return result.strip()
def solve_slider_captcha(self, bg_image, slider_image):
"""解决滑块验证码"""
# 这里可以使用OpenCV进行图像匹配
import cv2
import numpy as np
# 读取图像
bg = cv2.imdecode(np.frombuffer(bg_image, np.uint8), cv2.IMREAD_COLOR)
slider = cv2.imdecode(np.frombuffer(slider_image, np.uint8), cv2.IMREAD_COLOR)
# 转换为灰度图
bg_gray = cv2.cvtColor(bg, cv2.COLOR_BGR2GRAY)
slider_gray = cv2.cvtColor(slider, cv2.COLOR_BGR2GRAY)
# 模板匹配
result = cv2.matchTemplate(bg_gray, slider_gray, cv2.TM_CCOEFF_NORMED)
_, max_val, _, max_loc = cv2.minMaxLoc(result)
# 返回缺口位置
return max_loc[0] if max_val > 0.5 else None
def solve_recaptcha(self, site_key, url):
"""解决reCAPTCHA(使用第三方服务)"""
# 这里可以使用2captcha等服务
api_key = 'your_2captcha_api_key'
response = requests.post(
'http://2captcha.com/in.php',
data={
'key': api_key,
'method': 'userrecaptcha',
'googlekey': site_key,
'pageurl': url,
'json': 1
}
)
if response.json()['status'] == 1:
request_id = response.json()['request']
# 轮询结果
for _ in range(30):
time.sleep(5)
result = requests.get(
'http://2captcha.com/res.php',
params={
'key': api_key,
'action': 'get',
'id': request_id,
'json': 1
}
)
if result.json()['status'] == 1:
return result.json()['request']
return None
第七章:数据清洗与存储优化
7.1 数据清洗流程
python
# data_cleaning.py
import pandas as pd
import numpy as np
import re
from datetime import datetime
import hashlib
class ScienceDataCleaner:
"""科学数据清洗器"""
def __init__(self, df):
self.df = df.copy()
def remove_duplicates(self, subset=None):
"""去除重复记录"""
if subset:
initial_len = len(self.df)
self.df = self.df.drop_duplicates(subset=subset, keep='first')
print(f"去重前: {initial_len} 条, 去重后: {len(self.df)} 条")
return self
def handle_missing_values(self, strategy='drop'):
"""处理缺失值"""
if strategy == 'drop':
initial_len = len(self.df)
self.df = self.df.dropna()
print(f"删除缺失值: {initial_len} -> {len(self.df)}")
elif strategy == 'fill':
# 数值列用中位数填充
numeric_cols = self.df.select_dtypes(include=[np.number]).columns
for col in numeric_cols:
self.df[col].fillna(self.df[col].median(), inplace=True)
# 文本列用未知填充
text_cols = self.df.select_dtypes(include=['object']).columns
for col in text_cols:
self.df[col].fillna('Unknown', inplace=True)
return self
def clean_text_fields(self):
"""清理文本字段"""
text_cols = self.df.select_dtypes(include=['object']).columns
for col in text_cols:
# 去除多余空格
self.df[col] = self.df[col].astype(str).str.strip()
# 去除特殊字符
self.df[col] = self.df[col].apply(
lambda x: re.sub(r'[^\\w\\s\\-\\.\\(\\)\\[\\]]', '', x)
)
# 合并多个空格
self.df[col] = self.df[col].apply(
lambda x: re.sub(r'\\s+', ' ', x)
)
return self
def normalize_dates(self, date_column):
"""标准化日期格式"""
def parse_date(date_str):
try:
# 尝试多种格式
for fmt in ['%Y-%m-%d', '%Y/%m/%d', '%d-%b-%Y', '%B %d, %Y']:
try:
return datetime.strptime(date_str, fmt)
except:
continue
# 提取年份
year_match = re.search(r'\\b(19|20)\\d{2}\\b', str(date_str))
if year_match:
return datetime(int(year_match.group()), 1, 1)
except:
pass
return None
self.df[date_column] = self.df[date_column].apply(parse_date)
return self
def standardize_units(self):
"""标准化单位"""
# 温度单位转换
if 'temperature' in self.df.columns:
# 假设原始数据可能是华氏度
self.df['temperature_celsius'] = self.df.apply(
lambda row: (row['temperature'] – 32) * 5/9
if row.get('temp_unit') == 'F' else row['temperature'],
axis=1
)
return self
def add_metadata(self):
"""添加元数据"""
# 添加数据ID
self.df['record_id'] = self.df.apply(
lambda row: hashlib.md5(
str(row.to_dict()).encode()
).hexdigest()[:16],
axis=1
)
# 添加处理时间
self.df['processed_at'] = datetime.now()
return self
def validate_data(self):
"""数据验证"""
# 检查数值范围
numeric_cols = self.df.select_dtypes(include=[np.number]).columns
for col in numeric_cols:
# 检测异常值
mean = self.df[col].mean()
std = self.df[col].std()
outliers = self.df[
(self.df[col] < mean – 3*std) |
(self.df[col] > mean + 3*std)
]
if len(outliers) > 0:
print(f"列 {col} 发现 {len(outliers)} 个异常值")
return self
def get_cleaned_data(self):
"""获取清洗后的数据"""
return self.df
# 使用示例
def clean_science_data(input_file, output_file):
"""完整的清洗流程"""
# 读取数据
df = pd.read_csv(input_file)
# 创建清洗器
cleaner = ScienceDataCleaner(df)
# 执行清洗步骤
cleaned_df = (cleaner
.remove_duplicates()
.handle_missing_values(strategy='fill')
.clean_text_fields()
.normalize_dates('publication_date')
.add_metadata()
.validate_data()
.get_cleaned_data()
)
# 保存清洗后的数据
cleaned_df.to_csv(output_file, index=False, encoding='utf-8-sig')
print(f"清洗完成,保存至 {output_file}")
return cleaned_df
# 运行清洗
cleaned_data = clean_science_data(
'raw_data.csv',
'cleaned_science_data.csv'
)
7.2 多种存储方案
python
# data_storage.py
import sqlite3
import pymongo
from sqlalchemy import create_engine, Table, Column, Integer, String, Float, DateTime, MetaData
import pandas as pd
import pickle
import json
from pathlib import Path
class ScienceDataStorage:
"""科学数据存储管理器"""
def __init__(self, storage_type='csv', connection_string=None):
self.storage_type = storage_type
self.connection_string = connection_string
self.engine = None
self.connection = None
def connect(self):
"""建立连接"""
if self.storage_type == 'sqlite':
self.engine = create_engine(f'sqlite:///{self.connection_string}')
self.connection = self.engine.connect()
elif self.storage_type == 'postgresql':
self.engine = create_engine(self.connection_string)
self.connection = self.engine.connect()
elif self.storage_type == 'mongodb':
self.client = pymongo.MongoClient(self.connection_string)
self.db = self.client.get_default_database()
def save_dataframe(self, df, table_name, if_exists='replace'):
"""保存DataFrame到数据库"""
if self.storage_type in ['sqlite', 'postgresql']:
df.to_sql(
table_name,
self.connection,
if_exists=if_exists,
index=False,
chunksize=1000
)
print(f"数据已保存到 {table_name} 表,共 {len(df)} 条记录")
elif self.storage_type == 'mongodb':
collection = self.db[table_name]
# 转换为字典列表
records = df.to_dict('records')
if if_exists == 'replace':
collection.delete_many({})
collection.insert_many(records)
print(f"数据已保存到 MongoDB 集合 {table_name}")
def save_csv(self, df, filepath, **kwargs):
"""保存为CSV文件"""
df.to_csv(filepath, index=False, encoding='utf-8-sig', **kwargs)
print(f"数据已保存到 {filepath}")
def save_json(self, df, filepath, orient='records'):
"""保存为JSON文件"""
if orient == 'records':
df.to_json(filepath, orient='records', force_ascii=False, indent=2)
else:
df.to_json(filepath, orient=orient, force_ascii=False)
print(f"数据已保存到 {filepath}")
def save_parquet(self, df, filepath):
"""保存为Parquet格式(高效压缩)"""
df.to_parquet(filepath, compression='snappy')
print(f"数据已保存到 {filepath}")
def save_hdf5(self, df, filepath, key='data'):
"""保存为HDF5格式(适合大规模科学数据)"""
df.to_hdf(filepath, key=key, mode='w', complevel=6, complib='blosc')
print(f"数据已保存到 {filepath}")
def save_pickle(self, df, filepath):
"""保存为Pickle格式"""
with open(filepath, 'wb') as f:
pickle.dump(df, f)
print(f"数据已保存到 {filepath}")
def partition_by_date(self, df, date_column, base_path):
"""按日期分区存储"""
Path(base_path).mkdir(parents=True, exist_ok=True)
# 提取年份和月份
df['year'] = pd.to_datetime(df[date_column]).dt.year
df['month'] = pd.to_datetime(df[date_column]).dt.month
# 按年月分组存储
for (year, month), group in df.groupby(['year', 'month']):
filename = f"{base_path}/year={year}/month={month:02d}/data.parquet"
Path(filename).parent.mkdir(parents=True, exist_ok=True)
group.to_parquet(filename, compression='snappy')
print(f"数据已按日期分区存储到 {base_path}")
def close(self):
"""关闭连接"""
if self.connection:
self.connection.close()
if hasattr(self, 'client'):
self.client.close()
# 使用示例
def demonstrate_storage_options():
# 创建示例数据
df = pd.DataFrame({
'pmid': ['12345', '12346', '12347'],
'title': ['Study A', 'Study B', 'Study C'],
'temperature': [23.5, 24.1, 22.8],
'publication_date': ['2023-01-15', '2023-02-20', '2023-03-10']
})
# CSV存储
storage = ScienceDataStorage('csv')
storage.save_csv(df, 'science_data.csv')
# SQLite存储
storage = ScienceDataStorage('sqlite', 'science_data.db')
storage.connect()
storage.save_dataframe(df, 'publications')
storage.close()
# Parquet存储(推荐用于科学数据)
storage = ScienceDataStorage('parquet')
storage.save_parquet(df, 'science_data.parquet')
# 分区存储
storage = ScienceDataStorage('parquet')
storage.partition_by_date(df, 'publication_date', 'partitioned_data')
第八章:爬虫监控与错误处理
8.1 日志系统配置
python
# logging_config.py
import logging
import logging.handlers
import sys
from pathlib import Path
def setup_logging(log_dir='logs', log_level=logging.INFO):
"""配置日志系统"""
# 创建日志目录
Path(log_dir).mkdir(exist_ok=True)
# 创建日志器
logger = logging.getLogger('science_crawler')
logger.setLevel(log_level)
# 日志格式
formatter = logging.Formatter(
'%(asctime)s – %(name)s – %(levelname)s – %(filename)s:%(lineno)d – %(message)s'
)
# 文件处理器(轮转)
file_handler = logging.handlers.RotatingFileHandler(
f'{log_dir}/crawler.log',
maxBytes=10*1024*1024, # 10MB
backupCount=5,
encoding='utf-8'
)
file_handler.setFormatter(formatter)
logger.addHandler(file_handler)
# 错误日志专用文件
error_handler = logging.handlers.RotatingFileHandler(
f'{log_dir}/error.log',
maxBytes=10*1024*1024,
backupCount=5,
encoding='utf-8'
)
error_handler.setLevel(logging.ERROR)
error_handler.setFormatter(formatter)
logger.addHandler(error_handler)
# 控制台处理器
console_handler = logging.StreamHandler(sys.stdout)
console_handler.setFormatter(formatter)
logger.addHandler(console_handler)
return logger
8.2 爬虫状态监控
python
# crawler_monitor.py
import time
import psutil
import json
from datetime import datetime
from collections import deque
import threading
import smtplib
from email.mime.text import MIMEText
from email.mime.multipart import MIMEMultipart
class CrawlerMonitor:
"""爬虫监控器"""
def __init__(self, check_interval=60, alert_threshold=100):
self.check_interval = check_interval
self.alert_threshold = alert_threshold
self.stats = {
'start_time': datetime.now(),
'requests_count': 0,
'success_count': 0,
'error_count': 0,
'retry_count': 0,
'items_scraped': 0,
'items_dropped': 0,
'response_times': deque(maxlen=100),
'errors': deque(maxlen=50)
}
self.running = False
self.monitor_thread = None
def start(self):
"""启动监控"""
self.running = True
self.monitor_thread = threading.Thread(target=self._monitor_loop)
self.monitor_thread.daemon = True
self.monitor_thread.start()
print("监控器已启动")
def stop(self):
"""停止监控"""
self.running = False
if self.monitor_thread:
self.monitor_thread.join()
self._save_stats()
print("监控器已停止")
def record_request(self, success=True, response_time=None, error=None):
"""记录请求"""
self.stats['requests_count'] += 1
if success:
self.stats['success_count'] += 1
else:
self.stats['error_count'] += 1
if error:
self.stats['errors'].append({
'time': datetime.now(),
'error': str(error)
})
if response_time:
self.stats['response_times'].append(response_time)
def record_item(self, scraped=True):
"""记录数据项"""
if scraped:
self.stats['items_scraped'] += 1
else:
self.stats['items_dropped'] += 1
def get_system_stats(self):
"""获取系统状态"""
return {
'cpu_percent': psutil.cpu_percent(interval=1),
'memory_percent': psutil.virtual_memory().percent,
'disk_usage': psutil.disk_usage('/').percent
}
def get_crawler_stats(self):
"""获取爬虫状态"""
runtime = (datetime.now() – self.stats['start_time']).total_seconds()
success_rate = (self.stats['success_count'] / self.stats['requests_count'] * 100) \\
if self.stats['requests_count'] > 0 else 0
avg_response_time = sum(self.stats['response_times']) / len(self.stats['response_times']) \\
if self.stats['response_times'] else 0
return {
'runtime': runtime,
'requests_per_second': self.stats['requests_count'] / runtime if runtime > 0 else 0,
'success_rate': success_rate,
'avg_response_time': avg_response_time,
'items_scraped': self.stats['items_scraped'],
'items_dropped': self.stats['items_dropped'],
'error_count': self.stats['error_count']
}
def _monitor_loop(self):
"""监控循环"""
while self.running:
try:
crawler_stats = self.get_crawler_stats()
system_stats = self.get_system_stats()
# 检查是否需要告警
alerts = self._check_alerts(crawler_stats, system_stats)
if alerts:
self._send_alert(alerts)
# 记录状态
self._log_stats(crawler_stats, system_stats)
time.sleep(self.check_interval)
except Exception as e:
print(f"监控循环出错: {e}")
def _check_alerts(self, crawler_stats, system_stats):
"""检查告警条件"""
alerts = []
# 成功率过低
if crawler_stats['success_rate'] < 80:
alerts.append(f"成功率过低: {crawler_stats['success_rate']:.1f}%")
# 错误数过多
if crawler_stats['error_count'] > self.alert_threshold:
alerts.append(f"错误数过多: {crawler_stats['error_count']}")
# 系统资源过高
if system_stats['cpu_percent'] > 90:
alerts.append(f"CPU使用率过高: {system_stats['cpu_percent']}%")
if system_stats['memory_percent'] > 90:
alerts.append(f"内存使用率过高: {system_stats['memory_percent']}%")
return alerts
def _send_alert(self, alerts):
"""发送告警"""
# 这里可以集成邮件、短信、钉钉等通知方式
alert_message = "\\n".join(alerts)
print(f"!!! 告警: {alert_message}")
# 示例:发送邮件告警
# self._send_email_alert(alert_message)
def _log_stats(self, crawler_stats, system_stats):
"""记录状态到文件"""
stats = {
'timestamp': datetime.now().isoformat(),
'crawler': crawler_stats,
'system': system_stats
}
with open('crawler_stats.json', 'a') as f:
f.write(json.dumps(stats) + '\\n')
def _save_stats(self):
"""保存最终统计"""
final_stats = {
'end_time': datetime.now().isoformat(),
'final_stats': self.stats
}
with open('final_stats.json', 'w') as f:
json.dump(final_stats, f, indent=2, default=str)
# 使用示例
monitor = CrawlerMonitor(check_interval=30)
monitor.start()
# 在爬虫代码中记录
monitor.record_request(success=True, response_time=1.5)
monitor.record_item(scraped=True)
# 停止监控
monitor.stop()
第九章:实战项目——多源科学数据聚合系统
9.1 系统架构设计
我们将构建一个完整的科学数据聚合系统,从多个数据源采集数据并进行整合:
python
# science_data_aggregator.py
import asyncio
from concurrent.futures import ThreadPoolExecutor
import pandas as pd
from typing import Dict, List, Any
import json
from datetime import datetime
import hashlib
class ScienceDataAggregator:
"""多源科学数据聚合器"""
def __init__(self, config_file='sources_config.json'):
self.sources = self.load_sources(config_file)
self.data_cache = {}
self.executor = ThreadPoolExecutor(max_workers=5)
def load_sources(self, config_file):
"""加载数据源配置"""
with open(config_file, 'r') as f:
return json.load(f)
async def fetch_from_source(self, source_config):
"""从单个数据源获取数据"""
source_name = source_config['name']
source_type = source_config['type']
print(f"正在从 {source_name} 采集数据…")
try:
if source_type == 'pubmed':
data = await self.fetch_pubmed(source_config)
elif source_type == 'nasa_climate':
data = await self.fetch_nasa_climate(source_config)
elif source_type == 'arxiv':
data = await self.fetch_arxiv(source_config)
elif source_type == 'dryad':
data = await self.fetch_dryad(source_config)
else:
data = await self.fetch_generic(source_config)
# 标准化数据
normalized_data = self.normalize_data(data, source_config)
print(f"✓ {source_name} 采集完成,获得 {len(normalized_data)} 条记录")
return normalized_data
except Exception as e:
print(f"✗ {source_name} 采集失败: {e}")
return []
async def fetch_pubmed(self, config):
"""从PubMed采集"""
# 使用之前实现的PubMed爬虫
from pubmed_crawler import PubMedCrawler
crawler = PubMedCrawler()
articles = crawler.crawl(
keyword=config['query'],
max_pages=config.get('max_pages', 5)
)
return articles
async def fetch_nasa_climate(self, config):
"""从NASA气候数据网站采集"""
# 使用Playwright异步爬虫
from climate_crawler import AsyncClimateCrawler
crawler = AsyncClimateCrawler()
data = await crawler.run(config['urls'])
return data
async def fetch_arxiv(self, config):
"""从arXiv采集"""
import arxiv
client = arxiv.Client()
search = arxiv.Search(
query=config['query'],
max_results=config.get('max_results', 100),
sort_by=arxiv.SortCriterion.SubmittedDate
)
results = []
for paper in client.results(search):
results.append({
'title': paper.title,
'authors': [author.name for author in paper.authors],
'abstract': paper.summary,
'published': paper.published.isoformat(),
'doi': paper.doi,
'url': paper.entry_id,
'categories': paper.categories
})
return results
async def fetch_dryad(self, config):
"""从Dryad数据仓库采集"""
# 使用Dryad API
import requests
base_url = "https://datadryad.org/api/v2"
headers = {'Accept': 'application/json'}
response = requests.get(
f"{base_url}/search",
params={'q': config['query'], 'per_page': config.get('per_page', 50)},
headers=headers
)
if response.status_code == 200:
return response.json().get('_embedded', {}).get('datasets', [])
return []
async def fetch_generic(self, config):
"""通用采集方法"""
import aiohttp
async with aiohttp.ClientSession() as session:
async with session.get(
config['url'],
params=config.get('params', {}),
headers=config.get('headers', {})
) as response:
if response.status == 200:
return await response.json()
return []
def normalize_data(self, data: List[Dict], source_config) -> List[Dict]:
"""标准化数据格式"""
normalized = []
mapping = source_config.get('field_mapping', {})
for item in data:
normalized_item = {
'source': source_config['name'],
'source_type': source_config['type'],
'collection_time': datetime.now().isoformat(),
'data_id': hashlib.md5(str(item).encode()).hexdigest()[:16]
}
# 映射字段
for target_field, source_field in mapping.items():
if isinstance(source_field, str):
normalized_item[target_field] = item.get(source_field)
elif callable(source_field):
normalized_item[target_field] = source_field(item)
else:
normalized_item[target_field] = None
normalized.append(normalized_item)
return normalized
async def aggregate_all(self):
"""聚合所有数据源"""
# 创建采集任务
tasks = [self.fetch_from_source(source) for source in self.sources]
# 并发执行
results = await asyncio.gather(*tasks)
# 合并所有数据
all_data = []
for result in results:
all_data.extend(result)
return all_data
def deduplicate(self, data: List[Dict]) -> List[Dict]:
"""去重"""
seen = set()
unique_data = []
for item in data:
# 基于关键字段生成唯一标识
key_fields = ['title', 'doi', 'pmid']
for field in key_fields:
if field in item and item[field]:
key = f"{field}:{item[field]}"
if key not in seen:
seen.add(key)
unique_data.append(item)
break
else:
# 如果没有关键字段,使用data_id
if item['data_id'] not in seen:
seen.add(item['data_id'])
unique_data.append(item)
print(f"去重前: {len(data)} 条, 去重后: {len(unique_data)} 条")
return unique_data
def create_merged_dataset(self, data: List[Dict]) -> pd.DataFrame:
"""创建合并的数据集"""
df = pd.DataFrame(data)
# 添加统计信息
print(f"\\n=== 数据集统计 ===")
print(f"总记录数: {len(df)}")
print(f"数据来源: {df['source'].value_counts().to_dict()}")
print(f"时间范围: {df['collection_time'].min()} 至 {df['collection_time'].max()}")
return df
# sources_config.json 配置文件示例
CONFIG_EXAMPLE = """
[
{
"name": "PubMed Cancer Genomics",
"type": "pubmed",
"query": "cancer genomics",
"max_pages": 10,
"field_mapping": {
"title": "title",
"authors": "authors",
"abstract": "abstract",
"journal": "journal",
"publication_date": "publication_date",
"doi": "doi",
"pmid": "pmid"
}
},
{
"name": "NASA GISTEMP",
"type": "nasa_climate",
"urls": ["https://data.giss.nasa.gov/gistemp/"],
"field_mapping": {
"year": "Year",
"temperature_anomaly": "J-D",
"location": "location"
}
},
{
"name": "arXiv Machine Learning",
"type": "arxiv",
"query": "machine learning",
"max_results": 200,
"field_mapping": {
"title": "title",
"authors": "authors",
"abstract": "abstract",
"publication_date": "published",
"doi": "doi",
"url": "url",
"categories": "categories"
}
}
]
"""
# 运行聚合系统
async def main():
# 创建聚合器
aggregator = ScienceDataAggregator('sources_config.json')
# 采集所有数据
raw_data = await aggregator.aggregate_all()
# 去重
unique_data = aggregator.deduplicate(raw_data)
# 创建数据集
dataset = aggregator.create_merged_dataset(unique_data)
# 保存数据集
dataset.to_csv('merged_science_data.csv', index=False)
dataset.to_parquet('merged_science_data.parquet', compression='snappy')
print(f"\\n数据集已保存,共 {len(dataset)} 条记录")
if __name__ == "__main__":
asyncio.run(main())
9.2 数据集质量评估
python
# quality_assessment.py
import pandas as pd
import numpy as np
from sklearn.metrics import cohen_kappa_score
import matplotlib.pyplot as plt
import seaborn as sns
class DatasetQualityAssessor:
"""数据集质量评估器"""
def __init__(self, df):
self.df = df
self.quality_report = {}
def assess_completeness(self):
"""评估完整性"""
completeness = {}
for col in self.df.columns:
non_null = self.df[col].notna().sum()
completeness[col] = {
'non_null_count': non_null,
'completeness_rate': non_null / len(self.df) * 100,
'missing_count': len(self.df) – non_null
}
self.quality_report['completeness'] = completeness
return completeness
def assess_accuracy(self, reference_columns=None):
"""评估准确性"""
accuracy = {}
# 数值列的基本统计
numeric_cols = self.df.select_dtypes(include=[np.number]).columns
for col in numeric_cols:
accuracy[col] = {
'min': self.df[col].min(),
'max': self.df[col].max(),
'mean': self.df[col].mean(),
'std': self.df[col].std(),
'outliers_count': self.detect_outliers(self.df[col])
}
self.quality_report['accuracy'] = accuracy
return accuracy
def detect_outliers(self, series, method='iqr'):
"""检测异常值"""
if method == 'iqr':
Q1 = series.quantile(0.25)
Q3 = series.quantile(0.75)
IQR = Q3 – Q1
lower_bound = Q1 – 1.5 * IQR
upper_bound = Q3 + 1.5 * IQR
outliers = series[(series < lower_bound) | (series > upper_bound)]
return len(outliers)
return 0
def assess_consistency(self):
"""评估一致性"""
consistency = {}
# 检查日期格式一致性
date_cols = [col for col in self.df.columns if 'date' in col.lower()]
for col in date_cols:
try:
pd.to_datetime(self.df[col])
consistency[col] = {'date_format_consistent': True}
except:
consistency[col] = {'date_format_consistent': False}
# 检查分类变量的一致性
categorical_cols = self.df.select_dtypes(include=['object']).columns
for col in categorical_cols:
unique_values = self.df[col].nunique()
consistency[col] = {
'unique_values': unique_values,
'most_common': self.df[col].mode().iloc[0] if unique_values > 0 else None
}
self.quality_report['consistency'] = consistency
return consistency
def assess_timeliness(self, date_column='collection_time'):
"""评估时效性"""
if date_column in self.df.columns:
self.df[date_column] = pd.to_datetime(self.df[date_column])
latest = self.df[date_column].max()
earliest = self.df[date_column].min()
timeliness = {
'earliest_record': earliest,
'latest_record': latest,
'time_span_days': (latest – earliest).days,
'recency_days': (pd.Timestamp.now() – latest).days
}
self.quality_report['timeliness'] = timeliness
return timeliness
return {}
def assess_uniqueness(self, key_fields=None):
"""评估唯一性"""
if key_fields is None:
key_fields = ['data_id', 'doi', 'pmid']
uniqueness = {}
for field in key_fields:
if field in self.df.columns:
duplicate_count = self.df[field].duplicated().sum()
uniqueness[field] = {
'total_count': len(self.df),
'unique_count': self.df[field].nunique(),
'duplicate_count': duplicate_count,
'duplicate_rate': duplicate_count / len(self.df) * 100
}
self.quality_report['uniqueness'] = uniqueness
return uniqueness
def generate_quality_report(self):
"""生成质量报告"""
self.assess_completeness()
self.assess_accuracy()
self.assess_consistency()
self.assess_timeliness()
self.assess_uniqueness()
# 计算综合得分
scores = []
# 完整性得分
completeness_scores = [v['completeness_rate'] for v in self.quality_report['completeness'].values()]
scores.append(np.mean(completeness_scores))
# 唯一性得分
uniqueness_scores = [100 – v['duplicate_rate'] for v in self.quality_report['uniqueness'].values()]
scores.append(np.mean(uniqueness_scores))
# 时效性得分(最近30天内为100分,每超过30天减10分)
if 'timeliness' in self.quality_report:
recency_days = self.quality_report['timeliness']['recency_days']
timeliness_score = max(0, 100 – (recency_days // 30) * 10)
scores.append(timeliness_score)
overall_score = np.mean(scores) if scores else 0
self.quality_report['overall_score'] = overall_score
self.quality_report['quality_grade'] = self.get_grade(overall_score)
return self.quality_report
def get_grade(self, score):
"""获取质量等级"""
if score >= 90:
return 'A'
elif score >= 80:
return 'B'
elif score >= 70:
return 'C'
elif score >= 60:
return 'D'
else:
return 'F'
def visualize_quality(self):
"""可视化质量评估结果"""
fig, axes = plt.subplots(2, 2, figsize=(15, 10))
# 完整性热图
completeness_df = pd.DataFrame(self.quality_report['completeness']).T
sns.heatmap(completeness_df[['completeness_rate']].T,
annot=True, fmt='.1f', cmap='YlOrRd', ax=axes[0,0])
axes[0,0].set_title('数据完整性 (%)')
# 唯一性条形图
uniqueness_df = pd.DataFrame(self.quality_report['uniqueness']).T
uniqueness_df[['unique_count', 'duplicate_count']].plot(
kind='bar', ax=axes[0,1]
)
axes[0,1].set_title('数据唯一性')
axes[0,1].set_ylabel('记录数')
# 数值分布
if self.quality_report.get('accuracy'):
accuracy_df = pd.DataFrame(self.quality_report['accuracy']).T
if not accuracy_df.empty:
accuracy_df[['mean', 'std']].plot(kind='bar', ax=axes[1,0])
axes[1,0].set_title('数值列统计')
# 质量综合得分
axes[1,1].text(0.5, 0.5,
f"总体质量得分: {self.quality_report['overall_score']:.1f}\\n"
f"质量等级: {self.quality_report['quality_grade']}",
horizontalalignment='center',
verticalalignment='center',
transform=axes[1,1].transAxes,
fontsize=16)
axes[1,1].set_title('质量综合评估')
axes[1,1].axis('off')
plt.tight_layout()
plt.savefig('data_quality_report.png', dpi=300, bbox_inches='tight')
plt.show()
# 使用示例
def assess_dataset_quality(filepath):
# 读取数据集
df = pd.read_csv(filepath)
# 创建评估器
assessor = DatasetQualityAssessor(df)
# 生成质量报告
report = assessor.generate_quality_report()
# 打印报告
print("=== 数据集质量评估报告 ===")
print(f"总体得分: {report['overall_score']:.1f}")
print(f"质量等级: {report['quality_grade']}")
print(f"\\n评估时间: {pd.Timestamp.now()}")
# 可视化
assessor.visualize_quality()
return report
# 运行评估
quality_report = assess_dataset_quality('merged_science_data.csv')
第十章:最佳实践与未来展望
10.1 爬虫开发最佳实践
尊重robots.txt:在开始爬取前检查网站的robots.txt文件
设置合理延迟:避免对目标服务器造成压力
数据验证:始终验证爬取数据的完整性和准确性
错误处理:实现完善的错误处理和重试机制
日志记录:记录详细的运行日志以便排查问题
版本控制:使用Git管理爬虫代码
文档编写:为爬虫编写清晰的文档和使用说明
10.2 伦理与法律考量
-
遵守网站的服务条款和使用政策
-
仅用于合法研究和教育目的
-
尊重知识产权和数据版权
-
不在短时间内发送大量请求
-
如需商业使用,获取相应授权
10.3 未来技术趋势
AI辅助爬虫:使用机器学习识别页面结构和数据模式
无头浏览器进化:更轻量、更快速的浏览器自动化
分布式爬虫:利用云计算进行大规模数据采集
实时数据流:支持实时数据采集和处理
智能反反爬虫:更智能的请求模式模拟
10.4 扩展阅读与资源
-
官方文档:Scrapy, Selenium, Playwright
-
社区资源:Stack Overflow, GitHub开源项目
-
相关书籍:《Web Scraping with Python》《Python网络数据采集》
-
在线课程:Coursera, Udemy上的爬虫课程

