高性能异步协程爬虫
在执行某些IO密集型任务的时候,程序常常会因为等待 IO 而阻塞。为解决这一问题,可以考虑使用python中的协程异步。
从 Python 3.4 开始,Python 中加入了协程的概念,但这个版本的协程还是以生成器对象为基础的,在 Python 3.5 则增加了关键字async/await,使得协程的实现更加方便,本文中通过async/await 来实现协程。
python中使用协程最常用的库是asyncio,可以帮我们检测IO(只能是网络IO【HTTP连接就是网络IO操作】),实现应用程序级别的切换(异步IO)。
1、基本概念:
- event_loop:事件循环,相当于一个无限循环,我们可以把一些函数注册到这个事件循环上,当满足条件发生的时候,就会调用对应的处理方法。可通过asyncio.get_event_loop()方法来生成。
- coroutine:协程对象,我们可以将协程对象注册到事件循环中,它会被事件循环调用。我们可以使用 async 关键字来定义一个方法,这个方法在调用时不会立即被执行,而是返回一个协程对象。这个类似于含有yield关键字的生成器函数,在调用时先返回一个生成器对象。
- task:任务,它是对协程对象的进一步封装,包含了任务的各个状态。可以通过asyncio.ensure_future(coroutine)方法生成,也可以通过event_loop.create_task(coroutine)方法生成。如果有多个任务,我们可以将多个任务添加到一个列表比如tasks中,并将tasks作为参数传入asyncio.wait(tasks)方法,然后将这个整体注册到事件循环中。可以通过task.add_done_callback(callback)方法来给task绑定一个回调函数,在回调函数中,可以通过task.result()方法来获取task中return的返回值
- async/await :它是从 Python 3.5 才出现的,专门用于定义协程。其中,async 定义一个协程,await 用来挂起阻塞方法的执行。await 后面的对象可以是以下3种格式---1)一个原生 coroutine 对象。2)一个由 types.coroutine() 修饰的生成器,这个生成器可以返回 coroutine 对象。3)一个包含__await方法的对象返回的一个迭代器。
2.aiohttp库
要实现真正的异步,必须要使用支持异步操作的请求方式。aiohttp 是一个支持异步请求的库,利用它和 asyncio 配合我们可以非常方便地实现异步请求操作。
安装: pip3 install aiohttp
aiohttp模块的使用:
-
- aiohttp.ClientSession()方法可以实例化一个请求对象sess
- 请求对象sess可以像requests一样发送get、post请求,需要注意的是,代理参数使用proxy='http://ip:port'这种形式,而requests中使用的是proxies={‘http’:'http://ip:port','https':'https://ip:port'},如果代理需要使用身份验证,则使用参数proxy_auth=aiohttp.BasicAuth('用户名','密码'),其它请求参数同requests
- 响应对象response的方法:
- text()--返回字符串类型的响应数据
- read()--返回bytes类型的响应数据
- 通常会使用with关键词进行操作,可以省去close的步骤,在每一个with关键字前需要使用async关键字,在每一个阻塞操作前需要加上await关键字,示例:
async def get(url): async with aiohttp.ClientSession() as sess: async with await sess.get(url) as response: page_data = await response.text() return page_data
3、多任务协程爬虫思路
- 首先使用requests获取待爬取的页面url
- 将url写入列表,使用多任务异步协程爬取列表中的页面数据
4、多任务协程简单实现
-
利用flask实现一个慢速响应的服务器
# coding:utf-8 import time import random from flask import Flask app = Flask(__name__) @app.route('/') def index(): time.sleep(3) return 'hello world--'+str(random.random()) if __name__ == '__main__': app.run(debug=True) -
多任务协程代码
# coding:utf-8 import asyncio import time import aiohttp
# asyncio.set_event_loop_policy(asyncio.WindowsSelectorEventLoopPolicy()) windows系统请求https网站报错时调用此方法 async def get(url): async with aiohttp.ClientSession() as sess: #实例化请求对象 async with await sess.get(url) as response: #发送请求,获取响应对象 page_data = await response.text() #获取响应数据,read()方法获取bytes类型 return page_data #回调函数不需要使用async关键字 def parse(task): print(task.result()) #调用task.result()方法获取返回值 def main(): start_time = time.time() urls = ['http://127.0.0.1:5000/a','http://127.0.0.1:5000/b','http://127.0.0.1:5000/c'] #模拟requests请求获取到的url列表 tasks = [] for url in urls: c = get(url) task = asyncio.ensure_future(c) #创建3个任务对象 task.add_done_callback(parse) #给每个任务对象绑定一个回调函数parse tasks.append(task) loop = asyncio.get_event_loop() #创建事件循环对象 loop.run_until_complete(asyncio.wait(tasks)) #task列表传给asyncio.wait()方法,并注册到事件循环中执行。 end_time = time.time() print('花费时间:{}秒'.format(end_time-start_time)) if __name__ == '__main__': main() -
运行结果
--a-- --c-- --b-- 花费时间:3.014172315597534秒实现了异步请求的效果
4、限制并发数量
使用aiohttp时,python内部会使用select(),操作系统对文件描述符最大数量有限制,linux为1024个,windows为509个。限制并发数量(一般500),若并发的量不大可不作限制。
import asyncio
import aiohttp
async def get_http(url):
async with semaphore:
async with aiohttp.ClientSession() as session:
async with await session.get(url) as res:
global count
count += 1
print(count, res.status)
if __name__ == '__main__':
count = 0
semaphore = asyncio.Semaphore(500)
loop = asyncio.get_event_loop()
url = 'https://www.baidu.com/s?ie=utf-8&f=8&rsv_bp=1&ch=&tn=baiduerr&bar=&wd={0}'
tasks = [get_http(url.format(i)) for i in range(600)]
loop.run_until_complete(asyncio.wait(tasks))
loop.close()
5、多进程+协程异步
-
aiomultiprocess库
- 是一个异步多进程的python库,依赖于aiohttp和asyncio这两个库
- 安装
pip3 install aiomultiprocess - 以访问百度为例,测试多进程和异步协程同时请求
# coding:utf-8 import asyncio,aiohttp import time from multiprocessing.dummy import Pool from aiomultiprocess import Pool headers = { 'User-Agent':'Mozilla/5.0 (Windows NT 6.1; WOW64) AppleWebKit/535.1 (KHTML, like Gecko) Chrome/14.0.835.163 Safari/535.1' } async def get(url): async with aiohttp.ClientSession() as sess: async with await sess.get(url,headers=headers) as response: page_text = await response.text() return page_text async def request(): url = 'https://www.baidu.com/' urls = [url for _ in range(100)] async with Pool(4) as p: result = await p.map(get,urls) return result def multiprocess_crotine(): '''多进程异步协程''' loop = asyncio.get_event_loop() task = asyncio.ensure_future(request()) loop.run_until_complete(task) if __name__ == '__main__': st = time.time() multiprocess_crotine() ed = time.time() print('多进程异步协程耗时:%s秒'%(ed - st))