Queue
成都創(chuàng)新互聯公司主營宿豫網站建設的網絡公司,主營網站建設方案,重慶APP開發(fā),宿豫h5小程序制作搭建,宿豫網站營銷推廣歡迎宿豫等地區(qū)企業(yè)咨詢Tornado的tornado.queue模塊為基于協程的應用程序實現了一個異步生產者/消費者模式的隊列。這與python標準庫為多線程環(huán)境實現的queue模塊類似。
一個協程執(zhí)行到yieldqueue.get會暫停,直到隊列中有條目。如果queue有上限,一個協程執(zhí)行yieldqueue.put將會暫停,直到隊列中有空閑的位置。
在一個queue內部維護了一個未完成任務的引用計數,每調用一次put操作便會增加引用計數,而調用task_done操作將會減少引用計數。
下面是一個簡單的web爬蟲的例子:
最開始,queue只包含一個基準url。當一個worker從中取出一個url后,它會從對應的頁面中解析中所包含的url并將其放入隊列,然后調用task_done減少引用計數一次。
最后,worker會取出一個url,而這個url頁面中的所有url都已經被處理過了,這時隊列中也沒有url了。這時調用task_done會將引用計數減少至0.
這樣,在main協程里,join操作將會解除掛起并結束主協程。
這個爬蟲使用了HTMLParse來解析html頁面。
import time from datetime import timedelta try: from HTMLParser import HTMLParser from urlparse import urljoin, urldefrag except ImportError: from html.parser import HTMLParser from urllib.parse import urljoin, urldefrag from tornado import httpclient, gen, ioloop, queues base_url = 'http://www.tornadoweb.org/en/stable/' concurrency = 10 @gen.coroutine def get_links_from_url(url): """Download the page at `url` and parse it for links. Returned links have had the fragment after `#` removed, and have been made absolute so, e.g. the URL 'gen.html#tornado.gen.coroutine' becomes 'http://www.tornadoweb.org/en/stable/gen.html'. """ try: response = yield httpclient.AsyncHTTPClient().fetch(url) print('fetched %s' % url) html = response.body if isinstance(response.body, str) \ else response.body.decode() urls = [urljoin(url, remove_fragment(new_url)) for new_url in get_links(html)] except Exception as e: print('Exception: %s %s' % (e, url)) raise gen.Return([]) raise gen.Return(urls) #用于從一個包含片段的url中提取中真正的url. def remove_fragment(url): pure_url, frag = urldefrag(url) return pure_url def get_links(html): class URLSeeker(HTMLParser): def __init__(self): HTMLParser.__init__(self) self.urls = [] #從所有a標簽中提取中href屬性。 def handle_starttag(self, tag, attrs): href = dict(attrs).get('href') if href and tag == 'a': self.urls.append(href) url_seeker = URLSeeker() url_seeker.feed(html) return url_seeker.urls @gen.coroutine def main(): q = queues.Queue() start = time.time() fetching, fetched = set(), set() @gen.coroutine def fetch_url(): current_url = yield q.get() try: if current_url in fetching: return print('fetching %s' % current_url) fetching.add(current_url) urls = yield get_links_from_url(current_url) fetched.add(current_url) for new_url in urls: # Only follow links beneath the base URL if new_url.startswith(base_url): yield q.put(new_url) finally: q.task_done() @gen.coroutine def worker(): while True: yield fetch_url() q.put(base_url) # Start workers, then wait for the work queue to be empty. for _ in range(concurrency): worker() yield q.join(timeout=timedelta(seconds=300)) assert fetching == fetched print('Done in %d seconds, fetched %s URLs.' % ( time.time() - start, len(fetched))) if __name__ == '__main__': import logging logging.basicConfig() io_loop = ioloop.IOLoop.current() io_loop.run_sync(main)
另外有需要云服務器可以了解下創(chuàng)新互聯scvps.cn,海內外云服務器15元起步,三天無理由+7*72小時售后在線,公司持有idc許可證,提供“云服務器、裸金屬服務器、高防服務器、香港服務器、美國服務器、虛擬主機、免備案服務器”等云主機租用服務以及企業(yè)上云的綜合解決方案,具有“安全穩(wěn)定、簡單易用、服務可用性高、性價比高”等特點與優(yōu)勢,專為企業(yè)上云打造定制,能夠滿足用戶豐富、多元化的應用場景需求。