目标
在大规模爬取数据前,先定一个能达到的小目标,比方说先爬个10万条数据。
爬虫爬数据太慢了,怎么爬快点?
程序中途中断了怎么办,好不容易爬了这么多数据,又要重头开始爬吗/(ㄒoㄒ)/
数据有重复的,占用多余的空间,影响统计怎么办?
这些都是刚开始爬取大规模数据都会遇到的问题,这次就来说说解决这些问题的思路。
涉及的知识点如下:
- 多线程的生产者和消费者模型
- 断点数据的记录和恢复
- 数据入库前的去重
多线程的生产者和消费者模型
1. 单线程
-
其实可以看成生产者生产一个任务(比如构造出一个url),然后消费者执行这个任务(爬取url对应的网站),消费者任务还没执行完时,生产者就不会生产任务,所以他们相对任务来说是一对一的,同步执行的。当生产速率远远大于消费速率,这时生产者也会被拖累。
- 还有一个问题就是,我们的爬虫任务,在做请求网络时,实际上cpu大多数时候都在等待网络返回的包,没有完全发挥出cpu的计算能力,所以要想办法让cpu在等待网络响应时,也能动起来做些计算任务。而多线程就是为了解决这种情况出现的,操作系统会自动安排cpu在不同的线程中切换,提高cpu利用率。所以这就是为什么要开启多线程的原因。
2. 多线程
主要考虑的是线程的管理,任务的调度,还有线程间共享数据会出现的冲突问题。
-
先来看下面这张图。整个流程就是,我们给每个生产者和消费者都开启不同的线程,生产者生产任务,放入队列中存储,消费者从队列中取出任务,并执行任务,cpu会在不同的线程中切换,来执行任务的。
生产者不需要管任务是否被消费掉,只需要不停生产就行,而消费者也不需要去等待生产者,只要队列中有任务,就取出来执行就行了,这样就解决了任务的调度问题。
因为队列本身是线程安全的,换句话说就是队列中的数据,在多线程下不会出现数据冲突的问题,这样也就解决了队列中共享数据的冲突问题。对于其他共享的数据 ,可以使用线程锁,来给共享数据加锁,这样在释放锁之前,其他线程就不会访问这些数据,防止冲突。
因为生产者和消费者中间隔了一个队列,使得他们互不干扰,解耦程度高。这样也带来一个好处 ,当生产速率远大于消费速率,这时可以添加多个消费者(实际就是多开几个线程),提高消费速率,与生产速率达到相对的平衡,提高资源利用率。
关于多线程知识的补充
如果对于线程的基础概念不了解的,可以看下参考链接,讲解十分到位。
廖雪峰 进程和线程
菜鸟教程 python多线程
多线程代码思路
前面的都是介绍偏概念的知识,估计耐心看的人不多,这里就讲讲代码中的思路
- 导入threading线程库,队列库Queue,python3的是queue。
- 创建任务队列,必须要设置队列的大小,否则队列中的任务数量会无限增长,占用大量内存。
task_queue = Queue.Queue(maxsize=thread_max_count*10)
- 创建生产者线程,threading.Thread(target=producer)需要传入线程要执行函数的名字,注意函数名不要加括号。调用start后,线程才会真正的跑起来。
"""负责对生产者线程的创建"""
def producer_manager():
thread = threading.Thread(target=producer)
thread.start()
- 生产者执行的生产任务
"""生产者负责请求网站首页,解析出每个帖子的url,和创建出帖子的请求任务并放入任务队列中"""
def producer():
for title_page in range(start_page, total_page):
text = request_title(title_page)
if not text:
continue
tie_list = parse_title(text)
for tie in tie_list:
if tieba.find_one({'link':tie['link']}): #数据去重的判断
print('Data already exists: '+tie['link'])
else:
task_queue.put(tie) #将任务放入队列
log.update_one(run_log,{'$set':{'run_page':title_page}}) #记录断点数据
- 创建消费者线程。将线程放入list中放入方便管理。 调用task_queue.join(),表示若队列中还存在任务,那么主线程就阻塞住,不会往下面执行。
"""负责创建并管理消费者的线程"""
def consumer_manager():
threads = []
while len(threads) < thread_max_count:
thread = threading.Thread(target=consumer)
thread.setName("thread%d" % len(threads))
threads.append(thread)
thread.start()
task_queue.join()
6.消费者执行的任务。这里用了一个while True做死循环,这样线程就不会结束,避免了创建和销毁线程带来的开销,能提高一点运行效率。任务执行完后需要调用queue.task_done()函数,告诉队列已经完成一个任务,只有所有任务都调用过queue.task_done()以后,队列才会解除阻塞,主线程继续往下执行。
"""消费者负责从任务队列中取出任务(即任务的消费),并执行爬取每篇帖子和里面评论的任务"""
def consumer():
while True:
if task_queue.qsize() == 0:
time.sleep(3)
else:
task_count[0] = task_count[0] + 1
#print("run time second %s, ready task counts %d, finish task counts %d, db counts %d"%(str(time.time() - start_time)[:6], task_queue.qsize(),task_count[0],tieba.count()))
print("运行时间:%s秒,队列中剩余任务数%d,已完成任务数%d,数据已保存条数%d"%(str(time.time() - start_time)[:6], task_queue.qsize(),task_count[0],tieba.count()))
tie = task_queue.get()
request_comment(tie)
insert_db(tie)
task_queue.task_done()
断点记录和恢复
为了每次开始爬取数据不必重头开始,必须要记录下上次断点的位置
- 记录的手段,可以使用csv或者是数据库,这里用的是mongodb。
- 记录的位置。这个需要考虑一下。如果放在每个消费者线程中的话,记录的位置会比较多,到时候恢复起来比较麻烦。所以还是放在生产者线程中记录会比较好,每次分页请求导航页时,使用数据更新的方式将页数记录下来,恢复时读出页数,从这里开始继续爬取。
数据更新的函数update_one,接受的第一个参数,是表示查询的位置的,第二个参数里 '$set‘ 是固定用法,后面是更新的数据。
log.update_one(run_log,{'$set':{'run_page':title_page}})
- 因为恢复的时候可能会存在重复的数据,所以还需要做去重处理。
去重
- 去重最简单的方法,就是在写入数据库前,查找有没有这条记录,有的话就不写入,比较适合数据量不多时采用的方法。但在海量数据时,会受到空间和时间效率的限制,这时可以采用性能更加优秀的Bloom-Filter,即布隆过滤器算法,不过本人没研究过,这里不详细讨论了。
- 查询的条件要有唯一性,本文中就是直接对 URL 进行查询。
- 如果需要对数据库中某个字段频繁查询的话,会涉及到查询的效率问题,那就需要对这个字段做索引。索引就像书的目录,如果查找某内容在没有目录的帮助下,只能全篇从头到尾查找翻阅,这导致效率非常的低下;如果在借助目录情况下,就能很快的定位内容所在位置,效率会直线提高。
- 对字段做索引时需要的条件,该字段最好是能满足唯一性的,比如 ID,URL 这些数据,这样查找返回的值只有一个。还有这个字段的内容不能频繁变化,因为数据库引擎会对索引维护,其实就是对索引进行排序,索引值经常变化就会加大排序的负担,影响性能。
运行结果
可以放在服务器上运行,抓取了10万条帖子的数据。我用个人电脑运行时,可能是因为网络或者路由器的问题,最多只能开5个线程,多了容易出现请求超时的情况,所以最后没办法只能放在服务器上去跑,速度挺快的,可以开20个线程,每秒能抓取十几个帖子,只不过cpu是瓶颈,一运行cpu就满载,想以后再考虑优化下吧。
完整代码
以下是python2.7版本的代码。
使用python3.6运行的话,import Queue要变成小写的import queue,还有用queue.Queue()来创建队列。
# -*- coding: utf-8 -*-
import requests
from bs4 import BeautifulSoup
import pymongo
import re, math
import time,sys
import threading
import Queue
thread_max_count = 20
total_page = 1501
db_name = 'tieba3'
"""若请求超时,则重试请求,重试次数在5次以内"""
def request(method, url, **kwargs):
retry_count = 5
while retry_count > 0:
try:
res = requests.get(url, **kwargs) if method == 'get' else requests.post(url, **kwargs)
return res.text
except:
print('retry...', url)
retry_count -= 1
"""请求网站的导航页,获取帖子数据"""
def request_title(title_page=1):
title_url = "http://guba.eastmoney.com/list,cjpl_" + str(title_page) + ".html"
return request('get', title_url, timeout=5)
"""解析导航页帖子的标题数据,包括阅读数,评论数,标题,作者,发布时间,评论的总页数"""
def parse_title(text):
article_list = []
soup = BeautifulSoup(text, 'lxml')
host_url = 'http://guba.eastmoney.com'
elem_article = soup.find_all(name='div', class_='articleh')
for item in elem_article:
article_dict = {'read_count': '', 'comment_count': '', 'page': '', 'title': '', 'tie': '', 'author': '',
'time': '', 'link': '', 'comment': ''}
article_dict['read_count'] = item.select_one("span.l1").text
article_dict['comment_count'] = item.select_one("span.l2").text
article_dict['page'] = int(math.ceil(int(article_dict['comment_count']) / 30.0))
article_dict['title'] = item.select_one("span.l3 > a").text
article_dict['author'] = item.select_one("span.l4 > a").text if item.select_one("span.l4 > a") else u'匿名作者'
article_dict['time'] = item.select_one("span.l5").text
href = item.select_one("span.l3 > a").get("href")
article_dict['link'] = host_url + href if href[:1] == '/' else host_url + '/' + href
article_dict['comment'] = []
article_list.append(article_dict)
return article_list
"""根据评论的总页数,拼接出每一个评论页的url"""
def get_comment_urls(tie):
comment_urls = []
for cur_page in range(1, tie['page'] + 1 if tie['page'] > 0 else tie['page'] + 2):
comment_url = tie['link'][:-5] + '_' + str(cur_page) + ".html"
comment_urls.append(comment_url)
return comment_urls
"""请求评论页的数据"""
def request_comment(tie):
"""跳过一些不是帖子的链接"""
if re.compile(r'news,cjpl').search(tie['link']) == None:
return
print(tie['link']+' '+threading.currentThread().name)
for comment_url in get_comment_urls(tie):
text = request('get', comment_url, timeout=5)
parse_comment(text, tie)
"""解析出评论页的数据,包括作者,时间,评论内容和计算评论楼层"""
def parse_comment(text, tie):
soup = BeautifulSoup(text, 'lxml')
if (soup.find(name='div', id='zw_body')):
tie['tie'] = soup.find(name='div', id='zw_body').text.replace(u'\u3000', u'')
div_list = soup.find(id="mainbody").find_all(name='div', class_="zwlitxt")
for item in div_list:
comment_info = {"author": '', "time": '', "content": '', "lou": 0}
comment_info['author'] = item.find(name='span', class_="zwnick").text
comment_info['lou'] = len(tie['comment']) + 1
comment_info['time'] = item.find(name='div', class_="zwlitime").text[3:]
if (item.find(name='div', class_="zwlitext stockcodec")):
comment_info['content'] = item.find(name='div', class_="zwlitext stockcodec").text
comment_info['content'] = u"没有评论内容" if comment_info['content'] == '' else comment_info['content']
else:
comment_info['content'] = u"没有评论内容"
tie['comment'].append(comment_info)
"""消费者负责从任务队列中取出任务(即任务的消费),并执行爬取每篇帖子和里面评论的任务"""
def consumer():
while True:
if task_queue.qsize() == 0:
time.sleep(3)
else:
task_count[0] = task_count[0] + 1
#print("run time second %s, ready task counts %d, finish task counts %d, db counts %d"%(str(time.time() - start_time)[:6], task_queue.qsize(),task_count[0],tieba.count()))
print("运行时间:%s秒,队列中剩余任务数%d,已完成任务数%d,数据已保存条数%d"%(str(time.time() - start_time)[:6], task_queue.qsize(),task_count[0],tieba.count()))
tie = task_queue.get()
request_comment(tie)
insert_db(tie)
task_queue.task_done()
"""负责创建并管理消费者的线程"""
def consumer_manager():
threads = []
while len(threads) < thread_max_count:
thread = threading.Thread(target=consumer)
thread.setName("thread%d" % len(threads))
threads.append(thread)
thread.start()
task_queue.join()
"""数据保存"""
def insert_db(tie):
tieba.insert_one(tie)
"""生产者负责请求网站首页,解析出每个帖子的url,和创建出帖子的请求任务并放入任务队列中"""
def producer():
for title_page in range(start_page, total_page):
text = request_title(title_page)
if not text:
continue
tie_list = parse_title(text)
for tie in tie_list:
if tieba.find_one({'link':tie['link']}): #数据去重的判断
print('Data already exists: '+tie['link'])
else:
task_queue.put(tie)
log.update_one(run_log,{'$set':{'run_page':title_page}}) #记录断点数据
"""负责对生产者线程的创建"""
def producer_manager():
thread = threading.Thread(target=producer)
thread.start()
if __name__ == '__main__':
start_time = time.time()
task_count = [0]
client = pymongo.MongoClient('localhost', 27017)
test = client['test']
tieba = test[db_name]
"""
创建一个log数据库,记录断点的位置,每次重新运行就从断点为位置重爬,
这里记录的断点数据是帖子在首页的页数
"""
log = test['log']
run_log = {'db_name':db_name}
if not log.find_one(run_log):
log.insert_one(run_log)
start_page = 1
else:
start_page = log.find_one(run_log)['run_page']
print('start_page',start_page)
"""使用帖子的链接作为索引,可以提高去重时的查询效率"""
tieba.create_index('link')
"""必须要设置队列的大小,否则队列中的任务数量会无限增长,占用大量内存"""
task_queue = Queue.Queue(maxsize=thread_max_count*10)
producer_manager() #创建生产者线程
consumer_manager() #创建消费者线程