Python 爬虫数据存储
用解析器解析出数据之后,接下来就是存储数据了。保存形式多种多样:
- 文件存储:TXT、JSON、CSV
- 关系型数据库:MySQL、SQLite、PostgreSQL
- 非关系型数据库:MongoDB、Redis
一、文件存储
1.1 TXT 文本存储
TXT 文本几乎兼容任何平台,但不利于检索。适合对检索和数据结构要求不高的场景。
基本写入
import requests
from pyquery import PyQuery as pq
url = 'https://www.example.com'
headers = {'User-Agent': 'Mozilla/5.0 ...'}
html = requests.get(url, headers=headers).text
doc = pq(html)
# 写入文件
file = open('data.txt', 'a', encoding='utf-8')
file.write('内容')
file.write('\n' + '=' * 50 + '\n')
file.close()推荐写法(with as)
# 自动关闭文件
with open('data.txt', 'a', encoding='utf-8') as file:
file.write('内容\n')文件打开模式
| 模式 | 描述 |
|---|---|
r | 只读(默认),文件指针在开头 |
w | 写入,已存在则覆盖,不存在则创建 |
a | 追加,文件指针在末尾 |
rb/wb/ab | 二进制模式 |
r+/w+/a+ | 读写模式 |
1.2 JSON 文件存储
JSON 是轻量级的数据交换格式,结构清晰,是数据交换的极佳方式。
读取 JSON
import json
# 从字符串解析
str_data = '''
[{
"name": "Bob",
"gender": "male",
"birthday": "1992-10-18"
}]
'''
data = json.loads(str_data) # 转为 Python 对象
print(data[0]['name']) # Bob
print(data[0].get('age', 25)) # 25(默认值)
# 从文件读取
with open('data.json', 'r', encoding='utf-8') as file:
data = json.loads(file.read())⚠️ 注意:JSON 字符串必须使用双引号,单引号会导致解析错误。
写入 JSON
import json
data = [{'name': 'Bob', 'age': 20}]
with open('data.json', 'w', encoding='utf-8') as file:
# 格式化输出,支持中文
file.write(json.dumps(data, indent=2, ensure_ascii=False))1.3 CSV 文件存储
CSV 以纯文本形式存储表格数据,比 Excel 更简洁。
写入 CSV
import csv
# 列表写入
with open('data.csv', 'w', newline='', encoding='utf-8') as csvfile:
writer = csv.writer(csvfile)
writer.writerow(['id', 'name', 'age']) # 表头
writer.writerow(['10001', 'Mike', 20])
writer.writerows([['10002', 'Bob', 22], ['10003', 'Jordan', 21]])
# 字典写入
with open('data.csv', 'w', newline='', encoding='utf-8') as csvfile:
fieldnames = ['id', 'name', 'age']
writer = csv.DictWriter(csvfile, fieldnames=fieldnames)
writer.writeheader()
writer.writerow({'id': '10001', 'name': 'Mike', 'age': 20})读取 CSV
import csv
with open('data.csv', 'r', encoding='utf-8') as csvfile:
reader = csv.reader(csvfile)
for row in reader:
print(row) # ['id', 'name', 'age']
# 使用 pandas(推荐)
import pandas as pd
df = pd.read_csv('data.csv')
print(df)二、MySQL 存储
MySQL 是最常用的关系型数据库,Python 通过 PyMySQL 库操作。
2.1 连接数据库
import pymysql
db = pymysql.connect(
host='localhost',
user='root',
password='123456',
port=3306,
db='spiders'
)
cursor = db.cursor()2.2 创建表
sql = '''
CREATE TABLE IF NOT EXISTS students (
id VARCHAR(255) NOT NULL,
name VARCHAR(255) NOT NULL,
age INT NOT NULL,
PRIMARY KEY (id)
)
'''
cursor.execute(sql)
db.close()2.3 插入数据
# 基础插入
sql = 'INSERT INTO students(id, name, age) VALUES (%s, %s, %s)'
try:
cursor.execute(sql, ('20120001', 'Bob', 20))
db.commit() # 必须提交
except:
db.rollback() # 回滚
# 通用插入(字典动态构造)
data = {'id': '20120001', 'name': 'Bob', 'age': 20}
table = 'students'
keys = ', '.join(data.keys())
values = ', '.join(['%s'] * len(data))
sql = f'INSERT INTO {table}({keys}) VALUES ({values})'
try:
cursor.execute(sql, tuple(data.values()))
db.commit()
except:
db.rollback()2.4 更新数据(去重)
# 主键存在则更新,不存在则插入
sql = '''
INSERT INTO {table}({keys}) VALUES ({values})
ON DUPLICATE KEY UPDATE
'''.format(table=table, keys=keys, values=values)
update = ', '.join([f'{key} = %s' for key in data])
sql += update
try:
cursor.execute(sql, tuple(data.values()) * 2)
db.commit()
except:
db.rollback()2.5 查询数据
sql = 'SELECT * FROM students WHERE age >= 20'
try:
cursor.execute(sql)
print('Count:', cursor.rowcount)
# 获取单条
one = cursor.fetchone()
# 获取全部(大数据量慎用)
results = cursor.fetchall()
# 推荐:逐条获取
row = cursor.fetchone()
while row:
print(row)
row = cursor.fetchone()
except:
print('Error')2.6 删除数据
sql = 'DELETE FROM students WHERE age > 20'
try:
cursor.execute(sql)
db.commit()
except:
db.rollback()三、MongoDB 存储
MongoDB 是基于分布式文件存储的非关系型数据库,内容存储形式类似 JSON,非常灵活。
3.1 连接与初始化
import pymongo
# 连接
client = pymongo.MongoClient(host='localhost', port=27017)
# 或使用连接字符串
# client = pymongo.MongoClient('mongodb://localhost:27017/')
# 指定数据库
db = client.test # 或 client['test']
# 指定集合(类似表)
collection = db.students # 或 db['students']3.2 插入数据
student = {
'id': '20170101',
'name': 'Jordan',
'age': 20,
'gender': 'male'
}
# 插入单条(推荐)
result = collection.insert_one(student)
print(result.inserted_id)
# 插入多条
result = collection.insert_many([student1, student2])
print(result.inserted_ids)3.3 查询数据
# 查询单条
result = collection.find_one({'name': 'Mike'})
# 根据 ObjectId 查询
from bson.objectid import ObjectId
result = collection.find_one({'_id': ObjectId('593278c115c2602667ec6bae')})
# 查询多条
results = collection.find({'age': 20})
for result in results:
print(result)
# 条件查询
results = collection.find({'age': {'$gt': 20}}) # 大于 20
results = collection.find({'name': {'$regex': '^M.*'}}) # 正则匹配比较符号:
| 符号 | 含义 | 示例 |
|---|---|---|
$lt | 小于 | {'age': {'$lt': 20}} |
$gt | 大于 | {'age': {'$gt': 20}} |
$lte | 小于等于 | {'age': {'$lte': 20}} |
$gte | 大于等于 | {'age': {'$gte': 20}} |
$ne | 不等于 | {'age': {'$ne': 20}} |
$in | 在范围内 | {'age': {'$in': [20, 23]}} |
$regex | 正则匹配 | {'name': {'$regex': '^M.*'}} |
3.4 计数、排序、分页
# 计数
count = collection.find({'age': 20}).count()
# 排序
results = collection.find().sort('name', pymongo.ASCENDING) # 升序
results = collection.find().sort('name', pymongo.DESCENDING) # 降序
# 分页
results = collection.find().sort('name', pymongo.ASCENDING).skip(2).limit(2)3.5 更新数据
condition = {'name': 'Kevin'}
# 更新单条(推荐)
result = collection.update_one(condition, {'$set': {'age': 26}})
print(result.matched_count, result.modified_count)
# 更新多条
result = collection.update_many({'age': {'$gt': 20}}, {'$inc': {'age': 1}})3.6 删除数据
# 删除单条
result = collection.delete_one({'name': 'Kevin'})
# 删除多条
result = collection.delete_many({'age': {'$lt': 25}})
print(result.deleted_count)四、Redis 存储
Redis 是基于内存的高效键值型数据库,存取效率极高,支持多种数据结构。
4.1 连接 Redis
from redis import StrictRedis, ConnectionPool
# 方式一:直接连接
redis = StrictRedis(host='localhost', port=6379, db=0, password='foobared')
# 方式二:连接池
pool = ConnectionPool(host='localhost', port=6379, db=0, password='foobared')
redis = StrictRedis(connection_pool=pool)
# 方式三:URL 连接
url = 'redis://:foobared@localhost:6379/0'
pool = ConnectionPool.from_url(url)
redis = StrictRedis(connection_pool=pool)
# 测试
redis.set('name', 'Bob')
print(redis.get('name')) # b'Bob'4.2 字符串操作
redis.set('name', 'Bob') # 设置值
redis.get('name') # 获取值
redis.getset('name', 'Mike') # 设置新值并返回旧值
redis.mset({'name1': 'A', 'name2': 'B'}) # 批量设置
redis.mget(['name1', 'name2']) # 批量获取
redis.incr('age', 1) # 自增
redis.decr('age', 1) # 自减
redis.append('name', 'OK') # 追加字符串4.3 列表操作
redis.rpush('list', 1, 2, 3) # 右侧插入
redis.lpush('list', 0) # 左侧插入
redis.llen('list') # 长度
redis.lrange('list', 0, -1) # 获取范围
redis.lindex('list', 1) # 获取指定索引
redis.lpop('list') # 左侧弹出
redis.rpop('list') # 右侧弹出4.4 集合操作
redis.sadd('tags', 'Book', 'Tea', 'Coffee') # 添加元素
redis.srem('tags', 'Book') # 删除元素
redis.smembers('tags') # 所有元素
redis.sismember('tags', 'Book') # 是否存在
redis.sinter(['tags', 'tags2']) # 交集
redis.sunion(['tags', 'tags2']) # 并集
redis.sdiff(['tags', 'tags2']) # 差集4.5 有序集合操作
redis.zadd('grade', {'Bob': 100, 'Mike': 98}) # 添加(带分数)
redis.zrange('grade', 0, -1, withscores=True) # 按分数升序
redis.zrevrange('grade', 0, -1) # 按分数降序
redis.zrank('grade', 'Bob') # 排名(升序)
redis.zcount('grade', 80, 100) # 分数范围计数4.6 散列操作
redis.hset('price', 'cake', 5) # 设置字段
redis.hget('price', 'cake') # 获取字段
redis.hmset('price', {'apple': 3, 'banana': 2}) # 批量设置
redis.hmget('price', ['apple', 'banana']) # 批量获取
redis.hgetall('price') # 获取全部
redis.hkeys('price') # 所有键
redis.hvals('price') # 所有值4.7 键操作
redis.exists('name') # 键是否存在
redis.delete('name') # 删除键
redis.type('name') # 键类型
redis.keys('n*') # 模式匹配
redis.expire('name', 10) # 设置过期时间(秒)
redis.ttl('name') # 获取剩余时间五、Elasticsearch 搜索引擎存储
当需求从「存下来」变成「按内容检索」时,关系型数据库的 LIKE '%关键词%' 会在百万级数据上彻底失效——它无法使用索引,且不支持分词与相关度排序。
5.1 核心概念对应关系
| Elasticsearch | 关系型数据库 | 说明 |
|---|---|---|
| Index(索引) | Table | 文档的集合 |
| Document(文档) | Row | 一条 JSON 记录 |
| Field(字段) | Column | — |
| Mapping(映射) | Schema | 定义字段类型与分词器 |
| 倒排索引 | B+ 树索引 | 检索能力的根本差异来源 |
倒排索引是它的核心:普通索引记录「文档 → 内容」,倒排索引记录「词 → 出现该词的文档列表」。所以搜「爬虫技术」时,ES 不需要扫描任何文档,直接查词表取交集。
5.2 建立索引与映射
from elasticsearch import Elasticsearch
es = Elasticsearch("http://localhost:9200")
# 映射必须建库时定好:字段类型一旦确定不可修改,改类型只能重建索引
mapping = {
"mappings": {
"properties": {
# text:会分词,用于全文检索
"title": {"type": "text", "analyzer": "ik_max_word",
"search_analyzer": "ik_smart"},
"content": {"type": "text", "analyzer": "ik_max_word"},
# keyword:不分词,用于精确匹配、聚合、排序
"category": {"type": "keyword"},
"url": {"type": "keyword"},
"publish_date": {"type": "date"},
"score": {"type": "float"},
}
}
}
es.indices.create(index="articles", body=mapping, ignore=400)⚠️ text 与 keyword 是最容易踩的坑
text会被分词,不能用于精确匹配和聚合。用text字段做term查询几乎永远查不到结果。keyword不分词,整体作为一个词条,不能用于全文检索。
需要两种能力时,用多字段(multi-field)同时定义:"title": {"type": "text", "fields": {"raw": {"type": "keyword"}}},检索用 title,聚合排序用 title.raw。
中文必须安装分词插件(如 IK),否则默认按单字切分,检索质量极差。
5.3 写入数据
# 单条写入:指定 id 可实现幂等(重复写入即更新)
es.index(index="articles", id=doc_hash, document={
"title": "标题", "content": "正文", "category": "技术",
})
# 批量写入:爬虫必用,性能相差一个数量级
from elasticsearch.helpers import bulk
actions = [
{"_index": "articles", "_id": item["url"], "_source": item}
for item in items
]
success, failed = bulk(es, actions, raise_on_error=False)💡 用业务主键做 _id
把 URL 的哈希值作为 _id,重复抓取时会自动覆盖而非新增。这是爬虫场景下最简单的去重方案,省掉单独的去重逻辑。
5.4 查询
# 全文检索 + 高亮 + 过滤 + 分页,爬虫数据检索的典型组合
result = es.search(index="articles", body={
"query": {
"bool": {
"must": [{"match": {"content": "网络爬虫"}}], # 参与相关度打分
"filter": [ # 不打分,可缓存,更快
{"term": {"category": "技术"}},
{"range": {"publish_date": {"gte": "2026-01-01"}}},
],
}
},
"highlight": {"fields": {"content": {}}},
"sort": [{"_score": "desc"}],
"from": 0, "size": 20,
})
for hit in result["hits"]["hits"]:
print(hit["_score"], hit["_source"]["title"])must 与 filter 的区别:must 参与相关度评分(影响排序),filter 只做布尔过滤、结果可缓存因而更快。凡是不需要影响排序的条件(时间范围、分类、状态),一律放 filter。
5.5 何时该用 ES
| 场景 | 建议 |
|---|---|
| 需要全文检索、相关度排序 | ✅ ES 的核心价值 |
| 需要聚合分析(词频、趋势、分面统计) | ✅ 聚合能力强 |
| 只是存下来按主键读取 | ❌ 用 MongoDB,运维成本低得多 |
| 需要事务、强一致 | ❌ ES 是近实时(默认 1 秒刷新),非事务型 |
常见架构:MongoDB 存全量原始数据(数据源真相),ES 存需要检索的字段(检索层)。
六、RabbitMQ 消息队列
严格说 RabbitMQ 不是「存储」,而是任务的中转与解耦。在爬虫里它承担的是 分布式调度 的角色。
6.1 解决什么问题
爬虫规模化后,「抓取」和「解析入库」的速度往往不匹配:抓取快、解析慢,或反之。把两者放在同一进程里,慢的一方会拖死快的一方。
消息队列把它们解耦成生产者与消费者:
收益:两端可独立扩缩容、消费者宕机时任务不丢失、天然支持重试。
6.2 基本使用
import pika, json
conn = pika.BlockingConnection(pika.ConnectionParameters("localhost"))
channel = conn.channel()
# durable=True:队列元数据持久化,broker 重启后队列仍在
channel.queue_declare(queue="scrape_tasks", durable=True)
# 生产者
channel.basic_publish(
exchange="",
routing_key="scrape_tasks",
body=json.dumps({"url": "https://example.com/1"}),
properties=pika.BasicProperties(delivery_mode=2), # 消息本身也持久化
)# 消费者
def callback(ch, method, properties, body):
task = json.loads(body)
try:
process(task)
ch.basic_ack(delivery_tag=method.delivery_tag) # 处理成功才确认
except Exception:
# requeue=False 时进入死信队列,避免坏消息无限循环
ch.basic_nack(delivery_tag=method.delivery_tag, requeue=False)
channel.basic_qos(prefetch_count=1) # 每次只取一条,实现负载均衡
channel.basic_consume(queue="scrape_tasks", on_message_callback=callback)
channel.start_consuming()6.3 三个关键机制
| 机制 | 作用 | 不设置的后果 |
|---|---|---|
持久化(durable + delivery_mode=2) | 队列和消息写磁盘 | broker 重启后任务全部丢失 |
手动 ACK(basic_ack) | 处理成功才确认 | 消费者崩溃时任务丢失 |
prefetch_count | 限制单消费者未确认消息数 | 消息被快消费者一次性抢光,慢消费者空闲 |
⚠️ 必须配死信队列
消息处理失败若用 requeue=True 重新入队,一条永远失败的「毒消息」会无限循环并占满队列。生产环境必须配置死信交换机(DLX),失败消息转入死信队列人工排查。
6.4 与 Redis 队列的选择
| 维度 | RabbitMQ | Redis List/ZSET |
|---|---|---|
| 消息可靠性 | 高(持久化 + ACK + 死信) | 低(LPOP 后即丢失,崩溃即丢任务) |
| 路由能力 | 强(Exchange 支持广播、按键路由) | 无 |
| 吞吐 | 中 | 高 |
| 运维复杂度 | 较高 | 低 |
| 建议 | 任务不能丢、需要复杂路由 | 简单去重队列,Scrapy-Redis 即用此方案 |
多数爬虫项目用 Redis 就够了。只有当「任务丢失会造成实际业务损失」时,RabbitMQ 的额外复杂度才值得。
七、存储方式选择建议
| 存储方式 | 特点 | 适用场景 |
|---|---|---|
| TXT | 简单、兼容性好 | 日志、临时数据 |
| JSON | 结构化、易读 | API 数据、配置文件 |
| CSV | 表格形式、Excel 兼容 | 数据分析、导入导出 |
| MySQL | ACID 事务、复杂查询 | 结构化数据、关系型数据 |
| MongoDB | 灵活 Schema、高性能 | 非结构化数据、大规模存储 |
| Redis | 内存存储、极速读写 | 缓存、队列、去重 |
| Elasticsearch | 倒排索引、全文检索、聚合 | 内容检索、数据分析 |
| RabbitMQ | 可靠投递、解耦生产消费 | 分布式任务调度 |
爬虫常用组合:
- 小规模:CSV / JSON 文件
- 中等规模:MongoDB(灵活,无需建表)
- 大规模分布式:MongoDB + Redis(Redis 用于队列和去重)
- 需要检索:MongoDB(全量)+ Elasticsearch(检索层)
- 任务不可丢失:RabbitMQ 承担调度,MongoDB 承担落库
💡 一条通用原则
存储层要幂等。 爬虫必然会重跑、断点续爬、并发写入,所有写操作都应基于业务主键做 upsert(MongoDB 的 update_one(upsert=True)、MySQL 的 ON DUPLICATE KEY UPDATE、ES 的指定 _id)。做到这一点,就不需要任何额外的去重脚本。
📚 模块与官方文档
文件存储
| 模块 | 安装 | 官方文档 |
|---|---|---|
json | 标准库 | docs.python.org |
csv | 标准库 | docs.python.org |
pandas | pip install pandas | pandas.pydata.org · read_csv / to_csv |
数据库
| 模块 | 安装 | 官方文档 |
|---|---|---|
PyMySQL | pip install pymysql | pymysql.readthedocs.io |
SQLAlchemy | pip install sqlalchemy | docs.sqlalchemy.org |
pymongo | pip install pymongo | pymongo.readthedocs.io · CRUD 教程 |
motor(异步 MongoDB) | pip install motor | motor.readthedocs.io |
redis-py | pip install redis | redis.readthedocs.io · 命令参考 |
搜索与消息队列
| 模块 | 安装 | 官方文档 |
|---|---|---|
elasticsearch-py | pip install elasticsearch | elastic.co/guide · helpers.bulk |
| Elasticsearch 服务端 | — | 官方文档 · Mapping 类型 · Query DSL |
| IK 中文分词插件 | ES 插件 | GitHub |
pika | pip install pika | pika.readthedocs.io |
| RabbitMQ 服务端 | — | rabbitmq.com/docs · 死信交换机 |
celery(任务队列封装) | pip install celery | docs.celeryq.dev |