You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何通过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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.04 03:01:06