伪并发
异步和多路复用IO
线程存在空闲 from multiprocessing.dummy import Pool
1from multiprocessing.dummy import Pool as ThreadPool
2 pool = ThreadPool(20)
3 pool.map(job_worker, result_cursor)
4 pool.close()
5 pool.join()
1"""
2可以实现并发
3但是,请求发送出去后和返回之前,中间时期线程空闲
4编写方式:
5 - 直接返回处理
6 - 通过回调函数处理
7"""
8########### 编写方式一 ###########
9"""
10from concurrent.futures import ThreadPoolExecutor
11import requests
12import time
13
14def task(url):
15 response = requests.get(url)
16 print(url,response)
17 # 写正则表达式
18
19
20pool = ThreadPoolExecutor(7)
21url_list = [
22 'http://www.cnblogs.com/wupeiqi',
23 'http://huaban.com/favorite/beauty/',
24 'http://www.bing.com',
25 'http://www.zhihu.com',
26 'http://www.sina.com',
27 'http://www.baidu.com',
28 'http://www.autohome.com.cn',
29]
30for url in url_list:
31 pool.submit(task,url)
32
33pool.shutdown(wait=True)
34"""
效果相同,增加会掉函数
1########### 编写方式二 ###########
2from concurrent.futures import ThreadPoolExecutor
3import requests
4import time
5
6def task(url):
7 """
8 下载页面
9 :param url:
10 :return:
11 """
12 response = requests.get(url)
13 return response
14
15def done(future,*args,**kwargs):
16 response = future.result()
17 print(response.status_code,response.content)
18
19pool = ThreadPoolExecutor(7)
20url_list = [
21 'http://www.cnblogs.com/wupeiqi',
22 'http://huaban.com/favorite/beauty/',
23 'http://www.bing.com',
24 'http://www.zhihu.com',
25 'http://www.sina.com',
26 'http://www.baidu.com',
27 'http://www.autohome.com.cn',
28]
29for url in url_list:
30 v = pool.submit(task,url)
31 v.add_done_callback(done)
32
33pool.shutdown(wait=True)
多进程实现并发
1"""
2可以实现并发
3但是,请求发送出去后和返回之前,中间时期进程空闲
4编写方式:
5 - 直接返回处理
6 - 通过回调函数处理
7"""
8
9########### 编写方式一 ###########
10"""
11from concurrent.futures import ProcessPoolExecutor
12import requests
13import time
14
15def task(url):
16 response = requests.get(url)
17 print(url,response)
18 # 写正则表达式
19
20
21pool = ProcessPoolExecutor(7)
22url_list = [
23 'http://www.cnblogs.com/wupeiqi',
24 'http://huaban.com/favorite/beauty/',
25 'http://www.bing.com',
26 'http://www.zhihu.com',
27 'http://www.sina.com',
28 'http://www.baidu.com',
29 'http://www.autohome.com.cn',
30]
31for url in url_list:
32 pool.submit(task,url)
33
34pool.shutdown(wait=True)
35"""
36
37########### 编写方式二 ###########
38from concurrent.futures import ProcessPoolExecutor
39import requests
40import time
41
42def task(url):
43 response = requests.get(url)
44 return response
45
46def done(future,*args,**kwargs):
47 response = future.result()
48 print(response.status_code,response.content)
49
50pool = ProcessPoolExecutor(7)
51url_list = [
52 'http://www.cnblogs.com/wupeiqi',
53 'http://huaban.com/favorite/beauty/',
54 'http://www.bing.com',
55 'http://www.zhihu.com',
56 'http://www.sina.com',
57 'http://www.baidu.com',
58 'http://www.autohome.com.cn',
59]
60for url in url_list:
61 v = pool.submit(task,url)
62 v.add_done_callback(done)
63
64pool.shutdown(wait=True)
65
异步IO(多线程+协程)
1# 协程只是切换,不能控制什么时候切回来,异步IO实现回调
2角色:使用者
3 - 多线程
4 - 多线程
5 - 协程(微线程) + 异步IO =》 1个线程发送N个Http请求
6 - asyncio
7 - 示例1:asyncio.sleep(5)
8 - 示例2:自己封装Http数据包
9 - 示例3:asyncio+aiohttp
10 aiohttp模块:封装Http数据包 pip3 install aiohttp
11 - 示例4:asyncio+requests
12 requests模块:封装Http数据包 pip3 install requests
13 - gevent(内部异步IO+切换),greenlet+异步IO
14 pip3 install greenlet
15 pip3 install gevent
16 - 示例1:gevent+requests
17 - 示例2:gevent(协程池,最多发多少个请求)+requests
18 - 示例3:gevent+requests => grequests
19 pip3 install grequests
20
21 - Twisted
22 pip3 install twisted
23 - Tornado
24 pip3 install tornado
25
26 =====> gevent > Twisted > Tornado > asyncio
异步IO
1import asyncio
2
3
4@asyncio.coroutine
5def func1():
6 print('before...func1......')
7 yield from asyncio.sleep(5)
8 print('end...func1......')
9
10
11tasks = [func1(), func1()]
12
13loop = asyncio.get_event_loop()
14loop.run_until_complete(asyncio.gather(*tasks))
15loop.close()
异步IO实现tcp发http
1import asyncio
2
3
4@asyncio.coroutine
5def wget(host):
6 print('wget %s...' % host)
7 reader, writer = yield from asyncio.open_connection(host, 80)
8 header = 'GET / HTTP/1.0\r\nHost: %s\r\n\r\n' % host
9 writer.write(header.encode('utf-8'))
10 yield from writer.drain()
11 while True:
12 line = yield from reader.readline()
13 if line == b'\r\n':
14 break
15 print('%s header > %s' % (host, line.decode('utf-8').rstrip()))
16 # Ignore the body, close the socket
17 writer.close()
18
19
20loop = asyncio.get_event_loop()
21tasks = [wget(host) for host in ['www.sina.com.cn', 'www.sohu.com', 'www.163.com']]
22loop.run_until_complete(asyncio.wait(tasks))
23loop.close()
24
异步IO实现发http
1import aiohttp
2import asyncio
3
4
5@asyncio.coroutine
6def fetch_async(url):
7 print(url)
8 response = yield from aiohttp.request('GET', url)
9 print(url, response)
10 response.close()
11
12
13tasks = [fetch_async('http://www.baidu.com/'), fetch_async('http://www.chouti.com/')]
14
15event_loop = asyncio.get_event_loop()
16results = event_loop.run_until_complete(asyncio.gather(*tasks))
17event_loop.close()
异步IO + requests
1import asyncio
2import requests
3
4
5@asyncio.coroutine
6def fetch_async(func, *args):
7 loop = asyncio.get_event_loop()
8 future = loop.run_in_executor(None, func, *args)
9 response = yield from future
10 print(response.url, response.content)
11
12
13tasks = [
14 fetch_async(requests.get, 'http://www.cnblogs.com/wupeiqi/'),
15 fetch_async(requests.get, 'http://dig.chouti.com/pic/show?nid=4073644713430508&lid=10273091')
16]
17
18loop = asyncio.get_event_loop()
19results = loop.run_until_complete(asyncio.gather(*tasks))
20loop.close()
gevent + requests
1import gevent
2
3import requests
4from gevent import monkey
5
6monkey.patch_all()
7
8
9def fetch_async(method, url, req_kwargs):
10 print(method, url, req_kwargs)
11 response = requests.request(method=method, url=url, **req_kwargs)
12 print(response.url, response.content)
13
14# ##### 发送请求 #####
15gevent.joinall([
16 gevent.spawn(fetch_async, method='get', url='https://www.python.org/', req_kwargs={}),
17 gevent.spawn(fetch_async, method='get', url='https://www.yahoo.com/', req_kwargs={}),
18 gevent.spawn(fetch_async, method='get', url='https://github.com/', req_kwargs={}),
19])
20
21# ##### 发送请求(协程池控制最大协程数量) #####
22# from gevent.pool import Pool
23# pool = Pool(None)
24# gevent.joinall([
25# pool.spawn(fetch_async, method='get', url='https://www.python.org/', req_kwargs={}),
26# pool.spawn(fetch_async, method='get', url='https://www.yahoo.com/', req_kwargs={}),
27# pool.spawn(fetch_async, method='get', url='https://www.github.com/', req_kwargs={}),
28# ])
封装gevent + requests
1import grequests
2
3
4request_list = [
5 grequests.get('http://httpbin.org/delay/1', timeout=0.001),
6 grequests.get('http://fakedomain/'),
7 grequests.get('http://httpbin.org/status/500')
8]
9
10
11# ##### 执行并获取响应列表 #####
12# response_list = grequests.map(request_list)
13# print(response_list)
14
15
16# ##### 执行并获取响应列表(处理异常) #####
17# def exception_handler(request, exception):
18# print(request,exception)
19# print("Request failed")
20
21# response_list = grequests.map(request_list, exception_handler=exception_handler)
22# print(response_list)
twisted
1from twisted.internet import defer
2from twisted.web.client import getPage
3from twisted.internet import reactor
4
5
6def one_done(arg):
7 print('----------------------------------------------- %s' % arg)
8
9
10def all_done(arg):
11 print('done===========================================')
12 reactor.stop()
13
14
15@defer.inlineCallbacks # 发送Http请求,立即返回
16def task(url):
17 res = getPage(bytes(url, encoding='utf8')) # 发送Http请求
18 res.addCallback(one_done)
19 yield res
20
21
22url_list = [
23 'http://www.cnblogs.com',
24 'http://www.cnblogs.com',
25 'http://www.cnblogs.com',
26 'http://www.cnblogs.com',
27]
28
29defer_list = [] # [特殊,特殊,特殊(已经向url发送请求)]
30for url in url_list:
31 v = task(url)
32 defer_list.append(v)
33
34d = defer.DeferredList(defer_list)
35d.addBoth(all_done) # d特殊对象里有特殊url发送列表
36
37reactor.run() # 死循环 DeferredList查询,检测完成对象执行one_done,有计数器,所有执行完,执行all_done
tornado
1from tornado.httpclient import AsyncHTTPClient
2from tornado.httpclient import HTTPRequest
3from tornado import ioloop
4
5COUNT = 0
6def handle_response(response):
7 global COUNT
8 COUNT -= 1
9 if response.error:
10 print("Error:", response.error)
11 else:
12 print(response.body)
13 # 方法同twisted
14 # ioloop.IOLoop.current().stop()
15 if COUNT == 0:
16 ioloop.IOLoop.current().stop()
17
18def func():
19 url_list = [
20 'http://www.baidu.com',
21 'http://www.bing.com',
22 ]
23 global COUNT
24 COUNT = len(url_list)
25 for url in url_list:
26 print(url)
27 http_client = AsyncHTTPClient()
28 http_client.fetch(HTTPRequest(url), handle_response)
29
30
31ioloop.IOLoop.current().add_callback(func)
32ioloop.IOLoop.current().start() # 死循环
自己实现IO
1角色:NB开发者
2
3 1. socket客户端、服务端
4 连接阻塞
5 setblocking(0): 无数据(连接无响应;数据未返回)就报错 传0或者false所有socket都不会阻塞(包括连接和接收)
6
7 2. IO多路复用 就是while循环监听多个socket对象
8 客户端:
9 try:
10 socket对象1.connet()
11 socket对象2.connet()
12 socket对象3.connet()
13 except Ex..
14 pass
15
16 while True:
17 r(接收端),w(发送端),e(异常) = select.select([socket对象1,socket对象2,socket对象3,],[socket对象1,socket对象2,socket对象3,],[],0.05)
18 r = [socket对象1,] # 表示有人给我发送数据
19 xx = socket对象1.recv()
20 w = [socket对象1,] # 表示我已经和别人创建连接成功:
21 socket对象1.send('"""GET /index HTTP/1.0\r\nHost: baidu.com\r\n\r\n"""')
22
23
24 3.
25 class Foo:
26
27 def fileno(self):
28 obj = socket()
29 return obj.fileno()
30
31 r,w,e = select.select([socket对象?,对象?,对象?,Foo()],[],[])
32 # 对象必须有: fileno方法,并返回一个文件描述符
33
34 ========
35 a. select内部:对象.fileno()
36 b. Foo()内部封装socket文件描述符
37
38 IO多路复用: 就是while循环监听多个socket对象
39 异步IO: 非阻塞的socket+IO多路复用
自己实现异步IO
1class HttpRequest:
2 def __init__(self, sk, host, callback):
3 self.socket = sk
4 self.host = host
5 self.callback = callback
6
7 def fileno(self):
8 return self.socket.fileno()
9
10
11class HttpResponse:
12 def __init__(self, recv_data):
13 self.recv_data = recv_data
14 self.header_dict = {}
15 self.body = None
16
17 self.initialize()
18
19 def initialize(self):
20 headers, body = self.recv_data.split(b'\r\n\r\n', 1)
21 self.body = body
22 header_list = headers.split(b'\r\n')
23 for h in header_list:
24 h_str = str(h, encoding='utf-8')
25 v = h_str.split(':', 1)
26 if len(v) == 2:
27 self.header_dict[v[0]] = v[1]
28
29
30class AsyncRequest:
31 def __init__(self):
32 self.conn = []
33 self.connection = [] # 用于检测是否已经连接成功
34
35 def add_request(self, host, callback):
36 try:
37 sk = socket.socket()
38 sk.setblocking(0)
39 sk.connect((host, 80,))
40 except BlockingIOError as e:
41 pass
42 request = HttpRequest(sk, host, callback)
43 self.conn.append(request)
44 self.connection.append(request)
45
46 def run(self):
47
48 while True:
49 rlist, wlist, elist = select.select(self.conn, self.connection, self.conn, 0.05)
50 for w in wlist:
51 print(w.host, '连接成功...')
52 # 只要能循环到,表示socket和服务器端已经连接成功
53 tpl = "GET / HTTP/1.0\r\nHost:%s\r\n\r\n" % (w.host,)
54 w.socket.send(bytes(tpl, encoding='utf-8'))
55 self.connection.remove(w)
56 for r in rlist:
57 # r,是HttpRequest
58 recv_data = bytes()
59 while True:
60 try:
61 chunck = r.socket.recv(8096)
62 recv_data += chunck
63 except Exception as e:
64 break
65 response = HttpResponse(recv_data)
66 r.callback(response)
67 r.socket.close()
68 self.conn.remove(r)
69 if len(self.conn) == 0:
70 break
71
72
73def f1(response):
74 print('保存到文件', response.header_dict)
75
76
77def f2(response):
78 print('保存到数据库', response.header_dict)
79
80
81url_list = [
82 {'host': 'www.baidu.com', 'callback': f1},
83 {'host': 'cn.bing.com', 'callback': f2},
84 {'host': 'www.cnblogs.com', 'callback': f2},
85]
86
87req = AsyncRequest()
88for item in url_list:
89 req.add_request(item['host'], item['callback'])
90
91req.run()