Python爬虫 -- 美股吧多线程(二)

目标

在大规模爬取数据前,先定一个能达到的小目标,比方说先爬个10万条数据。

爬虫爬数据太慢了,怎么爬快点?
程序中途中断了怎么办,好不容易爬了这么多数据,又要重头开始爬吗/(ㄒoㄒ)/
数据有重复的,占用多余的空间,影响统计怎么办?

这些都是刚开始爬取大规模数据都会遇到的问题,这次就来说说解决这些问题的思路。
涉及的知识点如下:

  • 多线程的生产者和消费者模型
  • 断点数据的记录和恢复
  • 数据入库前的去重

多线程的生产者和消费者模型

1. 单线程

  • 其实可以看成生产者生产一个任务(比如构造出一个url),然后消费者执行这个任务(爬取url对应的网站),消费者任务还没执行完时,生产者就不会生产任务,所以他们相对任务来说是一对一的,同步执行的。当生产速率远远大于消费速率,这时生产者也会被拖累。


    生产者和消费者单线程版本
  • 还有一个问题就是,我们的爬虫任务,在做请求网络时,实际上cpu大多数时候都在等待网络返回的包,没有完全发挥出cpu的计算能力,所以要想办法让cpu在等待网络响应时,也能动起来做些计算任务。而多线程就是为了解决这种情况出现的,操作系统会自动安排cpu在不同的线程中切换,提高cpu利用率。所以这就是为什么要开启多线程的原因。

2. 多线程

  • 主要考虑的是线程的管理,任务的调度,还有线程间共享数据会出现的冲突问题。

  • 先来看下面这张图。整个流程就是,我们给每个生产者和消费者都开启不同的线程,生产者生产任务,放入队列中存储,消费者从队列中取出任务,并执行任务,cpu会在不同的线程中切换,来执行任务的。


    生产者和消费者多线程版本
  • 生产者不需要管任务是否被消费掉,只需要不停生产就行,而消费者也不需要去等待生产者,只要队列中有任务,就取出来执行就行了,这样就解决了任务的调度问题。

  • 因为队列本身是线程安全的,换句话说就是队列中的数据,在多线程下不会出现数据冲突的问题,这样也就解决了队列中共享数据的冲突问题。对于其他共享的数据 ,可以使用线程锁,来给共享数据加锁,这样在释放锁之前,其他线程就不会访问这些数据,防止冲突。

  • 因为生产者和消费者中间隔了一个队列,使得他们互不干扰,解耦程度高。这样也带来一个好处 ,当生产速率远大于消费速率,这时可以添加多个消费者(实际就是多开几个线程),提高消费速率,与生产速率达到相对的平衡,提高资源利用率。

关于多线程知识的补充

如果对于线程的基础概念不了解的,可以看下参考链接,讲解十分到位。
廖雪峰 进程和线程
菜鸟教程 python多线程

多线程代码思路

前面的都是介绍偏概念的知识,估计耐心看的人不多,这里就讲讲代码中的思路

  1. 导入threading线程库,队列库Queue,python3的是queue。
  2. 创建任务队列,必须要设置队列的大小,否则队列中的任务数量会无限增长,占用大量内存。
task_queue = Queue.Queue(maxsize=thread_max_count*10)
  1. 创建生产者线程,threading.Thread(target=producer)需要传入线程要执行函数的名字,注意函数名不要加括号。调用start后,线程才会真正的跑起来。
"""负责对生产者线程的创建"""
def producer_manager():
    thread = threading.Thread(target=producer)
    thread.start()
  1. 生产者执行的生产任务
"""生产者负责请求网站首页,解析出每个帖子的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}})  #记录断点数据
  1. 创建消费者线程。将线程放入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()  #创建消费者线程
最后编辑于
©著作权归作者所有,转载或内容合作请联系作者
  • 序言:七十年代末,一起剥皮案震惊了整个滨河市,随后出现的几起案子,更是在滨河造成了极大的恐慌,老刑警刘岩,带你破解...
    沈念sama阅读 216,039评论 6 498
  • 序言:滨河连续发生了三起死亡事件,死亡现场离奇诡异,居然都是意外死亡,警方通过查阅死者的电脑和手机,发现死者居然都...
    沈念sama阅读 92,223评论 3 392
  • 文/潘晓璐 我一进店门,熙熙楼的掌柜王于贵愁眉苦脸地迎上来,“玉大人,你说我怎么就摊上这事。” “怎么了?”我有些...
    开封第一讲书人阅读 161,916评论 0 351
  • 文/不坏的土叔 我叫张陵,是天一观的道长。 经常有香客问我,道长,这世上最难降的妖魔是什么? 我笑而不...
    开封第一讲书人阅读 58,009评论 1 291
  • 正文 为了忘掉前任,我火速办了婚礼,结果婚礼上,老公的妹妹穿的比我还像新娘。我一直安慰自己,他们只是感情好,可当我...
    茶点故事阅读 67,030评论 6 388
  • 文/花漫 我一把揭开白布。 她就那样静静地躺着,像睡着了一般。 火红的嫁衣衬着肌肤如雪。 梳的纹丝不乱的头发上,一...
    开封第一讲书人阅读 51,011评论 1 295
  • 那天,我揣着相机与录音,去河边找鬼。 笑死,一个胖子当着我的面吹牛,可吹牛的内容都是我干的。 我是一名探鬼主播,决...
    沈念sama阅读 39,934评论 3 416
  • 文/苍兰香墨 我猛地睁开眼,长吁一口气:“原来是场噩梦啊……” “哼!你这毒妇竟也来了?” 一声冷哼从身侧响起,我...
    开封第一讲书人阅读 38,754评论 0 271
  • 序言:老挝万荣一对情侣失踪,失踪者是张志新(化名)和其女友刘颖,没想到半个月后,有当地人在树林里发现了一具尸体,经...
    沈念sama阅读 45,202评论 1 309
  • 正文 独居荒郊野岭守林人离奇死亡,尸身上长有42处带血的脓包…… 初始之章·张勋 以下内容为张勋视角 年9月15日...
    茶点故事阅读 37,433评论 2 331
  • 正文 我和宋清朗相恋三年,在试婚纱的时候发现自己被绿了。 大学时的朋友给我发了我未婚夫和他白月光在一起吃饭的照片。...
    茶点故事阅读 39,590评论 1 346
  • 序言:一个原本活蹦乱跳的男人离奇死亡,死状恐怖,灵堂内的尸体忽然破棺而出,到底是诈尸还是另有隐情,我是刑警宁泽,带...
    沈念sama阅读 35,321评论 5 342
  • 正文 年R本政府宣布,位于F岛的核电站,受9级特大地震影响,放射性物质发生泄漏。R本人自食恶果不足惜,却给世界环境...
    茶点故事阅读 40,917评论 3 325
  • 文/蒙蒙 一、第九天 我趴在偏房一处隐蔽的房顶上张望。 院中可真热闹,春花似锦、人声如沸。这庄子的主人今日做“春日...
    开封第一讲书人阅读 31,568评论 0 21
  • 文/苍兰香墨 我抬头看了看天上的太阳。三九已至,却和暖如春,着一层夹袄步出监牢的瞬间,已是汗流浃背。 一阵脚步声响...
    开封第一讲书人阅读 32,738评论 1 268
  • 我被黑心中介骗来泰国打工, 没想到刚下飞机就差点儿被人妖公主榨干…… 1. 我叫王不留,地道东北人。 一个月前我还...
    沈念sama阅读 47,583评论 2 368
  • 正文 我出身青楼,却偏偏与公主长得像,于是被迫代替她去往敌国和亲。 传闻我的和亲对象是个残疾皇子,可洞房花烛夜当晚...
    茶点故事阅读 44,482评论 2 352

推荐阅读更多精彩内容

  • 1.进程和线程 队列:1、进程之间的通信: q = multiprocessing.Queue()2、...
    一只写程序的猿阅读 1,109评论 0 17
  • Spring Cloud为开发人员提供了快速构建分布式系统中一些常见模式的工具(例如配置管理,服务发现,断路器,智...
    卡卡罗2017阅读 134,650评论 18 139
  • 从哪说起呢? 单纯讲多线程编程真的不知道从哪下嘴。。 不如我直接引用一个最简单的问题,以这个作为切入点好了 在ma...
    Mr_Baymax阅读 2,752评论 1 17
  • NSThread 第一种:通过NSThread的对象方法 NSThread *thread = [[NSThrea...
    攻城狮GG阅读 797评论 0 3
  • 背景 担心了两周的我终于轮到去医院做胃镜检查了!去的时候我都想好了最坏的可能(胃癌),之前在网上查的症状都很相似。...
    Dely阅读 9,237评论 21 42