如何通过Python asyncio Queue加速股票数据爬虫进程?
问题
我想要构建一个简易爬虫,从维基百科链接获取标普500公司股票代码列表,再到雅虎财经查询这些代码对应的股价。为提升健壮性并深入学习Python asyncio,我采用了asyncio.Queue实现。目前爬虫可正常运行,但爬取仅15个代码就耗时62秒。以下是我的三段代码,请问有哪些优化手段可以提速?
维基百科爬虫代码
from selectolax.lexbor import LexborHTMLParser import httpx class WikiWorker: def __init__(self): self._url = "https://en.wikipedia.org/wiki/List_of_S%26P_500_companies" @staticmethod def _parse_company_symbols(page_html): html = LexborHTMLParser(page_html) table = html.css_first("table#constituents") for td in table.css("tr")[1:15]: symbol = td.css_first("td").text(strip=True) yield symbol def get_sp_500_companies(self): response = httpx.get(self._url) html = response.text yield from self._parse_company_symbols(html) if __name__ == "__main__": wikiWorker = WikiWorker() asyncio.run(wikiWorker.get_sp_500_companies())
雅虎财经爬虫代码
import asyncio from selectolax.lexbor import LexborHTMLParser import httpx class YahooFinancePriceScheduler: def __init__(self, input_queue: asyncio.Queue, **kwargs): self._input_queue = input_queue async def run(self): while True: val = await self._input_queue.get() print(val) yahooFinanceWorker = YahooFinanceWorker(symbol=val) price = yahooFinanceWorker.get_price() print(price) self._input_queue.task_done() class YahooFinanceWorker: def __init__(self, symbol, **kwargs): self._symbol = symbol base_url = "https://finance.yahoo.com/quote/" self._url = f"{base_url}{self._symbol}" headers = { "User-Agent": "Mozilla/5.0 (X11; Linux x86_64; rv:120.0) Gecko/20100101 Firefox/120.0" } def get_price(self): print(self._url) response = httpx.get(self._url, headers=self.headers, follow_redirects=True) html = LexborHTMLParser(response.text) try: price = float( html.css_first('[data-test="qsp-price"]') .text(strip=True) .replace(",", "") ) return price except AttributeError: return None if __name__ == "__main__": yfps = YahooFinancePriceScheduler("MMM") yfps.run()
爬虫入口代码
import asyncio import time from workers.wikiWorker import WikiWorker from workers.yahooFinanceWorker import YahooFinancePriceScheduler NUM_WORKERS = 5 SYMBOL_QUEUE_MAX_SIZE = 100 async def main(): symbol_queue = asyncio.Queue() tasks = [] wikiWorker = WikiWorker() scraper_start_time = time.time() for _ in range(NUM_WORKERS): yahooFinancePriceScheduler = YahooFinancePriceScheduler( input_queue=symbol_queue ) tasks.append(asyncio.create_task(yahooFinancePriceScheduler.run())) for symbol in wikiWorker.get_sp_500_companies(): await symbol_queue.put(symbol) await symbol_queue.join() for task in tasks: task.cancel() print("Extracting time took:", round(time.time() - scraper_start_time, 1)) if __name__ == "__main__": asyncio.run(main())
优化方案
1. 替换同步HTTP请求为异步请求(核心优化)
当前雅虎财经爬虫用的是httpx.get()同步请求,这会直接阻塞asyncio事件循环,完全浪费了异步框架的并发优势——这是耗时62秒的核心原因。
修改方案:
- 将
YahooFinanceWorker.get_price()改为异步方法,使用httpx.AsyncClient发送请求 - 复用全局AsyncClient实例,避免每次请求重建连接,提升连接复用效率
修改后的YahooFinanceWorker:
class YahooFinanceWorker: def __init__(self, symbol, **kwargs): self._symbol = symbol base_url = "https://finance.yahoo.com/quote/" self._url = f"{base_url}{self._symbol}" headers = { "User-Agent": "Mozilla/5.0 (X11; Linux x86_64; rv:120.0) Gecko/20100101 Firefox/120.0" } async def get_price(self, client: httpx.AsyncClient): response = await client.get(self._url, headers=self.headers, follow_redirects=True) html = LexborHTMLParser(response.text) try: price = float( html.css_first('[data-test="qsp-price"]') .text(strip=True) .replace(",", "") ) return price except AttributeError: return None
修改Scheduler的run方法,传入全局AsyncClient:
class YahooFinancePriceScheduler: def __init__(self, input_queue: asyncio.Queue, client: httpx.AsyncClient, **kwargs): self._input_queue = input_queue self._client = client async def run(self): while True: val = await self._input_queue.get() worker = YahooFinanceWorker(symbol=val) price = await worker.get_price(self._client) self._input_queue.task_done()
2. 调高并发数
当前设置的NUM_WORKERS = 5并发数过低,无法充分利用异步优势。可以根据网站反爬限制调整到10-20之间(注意不要过度触发反爬):
NUM_WORKERS = 15 # 可根据实际测试调整
3. 优化维基百科爬虫的异步适配
维基爬虫仅请求一次,对整体速度影响不大,但可以改成异步方法保持代码风格统一:
async def get_sp_500_companies(self): async with httpx.AsyncClient() as client: response = await client.get(self._url) html = response.text yield from self._parse_company_symbols(html)
同时入口代码中的入队逻辑也要改成异步遍历:
async for symbol in wikiWorker.get_sp_500_companies(): await symbol_queue.put(symbol)
4. 减少冗余IO操作
代码中大量的print()语句属于同步IO操作,会拖慢异步任务执行速度,测试完成后可以注释或删除,或者改用异步日志库记录。
5. 添加请求延迟与重试机制
雅虎财经会对高频请求限流,添加随机延迟和重试逻辑可以避免被封禁,同时保证请求成功率:
- 添加随机延迟:在每个请求完成后加入
await asyncio.sleep(random.uniform(0.3, 0.8)) - 实现重试:可以用
tenacity库或自定义重试逻辑,处理请求超时、4xx/5xx错误
最终优化后的入口代码示例
import asyncio import time import random import httpx from workers.wikiWorker import WikiWorker from workers.yahooFinanceWorker import YahooFinancePriceScheduler NUM_WORKERS = 15 SYMBOL_QUEUE_MAX_SIZE = 200 async def main(): symbol_queue = asyncio.Queue(maxsize=SYMBOL_QUEUE_MAX_SIZE) tasks = [] # 创建全局异步HTTP客户端 async with httpx.AsyncClient() as client: wiki_worker = WikiWorker() start_time = time.time() # 启动异步工作任务 for _ in range(NUM_WORKERS): scheduler = YahooFinancePriceScheduler( input_queue=symbol_queue, client=client ) tasks.append(asyncio.create_task(scheduler.run())) # 异步获取符号并存入队列 async for symbol in wiki_worker.get_sp_500_companies(): await symbol_queue.put(symbol) await symbol_queue.join() # 取消并回收任务 for task in tasks: task.cancel() await asyncio.gather(*tasks, return_exceptions=True) print(f"Extracting time took: {round(time.time() - start_time, 1)}s") if __name__ == "__main__": asyncio.run(main())
内容的提问来源于stack exchange,提问作者dougj
相关产品推荐
相关产品推荐

