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

Python异步调用改造求助:批量处理股票代码性能优化问题

问题:6000+股票批量爬取期权数据速度过慢,异步改造卡壳求助

嗨各位大佬,我是Python新手,这是我写的第一个Python程序——用来从Yahoo Finance爬取股票期权数据,然后插入SQL Server数据库。处理少量股票的时候一切正常,但跑包含6000+股票代码的文件时,速度慢到离谱。我知道有其他性能优化路子,但想优先通过异步调用来提升效率。

之前我试着把read_ticker_file改成async函数,还在yfin_options()调用前加了await,但完全没效果。我猜得重构整个调用逻辑,但现在卡壳了,恳请各位大佬给我指条异步改造的明路!

我当前的代码:

import logging
import pyodbc
import config
import yahoo_fin as yfin
from yahoo_fin import options
from datetime import datetime, date
from selenium import webdriver

def main():
    read_ticker_file()

def init_selenium():
    driver = webdriver.Chrome(config.CHROME_DRIVER)
    return driver

def yfin_options(symbol):
    logging.basicConfig(filename='yfin.log', level=logging.INFO)
    logging.basicConfig(filename='no_options.log', level=logging.ERROR)
    try:
        # get all options dates (in epoch) from dropdown on yahoo finance options page
        dates = get_exp_dates(symbol)
        # iterate each date to get all calls and insert into sql db
        for date in dates:
            arr = yfin.options.get_calls(symbol, date)
            arr_length = len(arr.values)
            i = 0
            for x in range(0, arr_length):
                strike = str(arr.values[i][2])
                volume = str(arr.values[i][8])
                open_interest = str(arr.values[i][9])
                convert_epoch = datetime.fromtimestamp(int(date))
                try:
                    sql_insert(symbol, strike, volume, open_interest, convert_epoch)
                    i += 1
                except Exception as insert_fail:
                    print("I failed at sqlinsert {0}".format(insert_fail))
            file_name_dir = "C:\\temp\\rh\\options{0}{1}.xlsx".format(symbol, date)
            logging.info(arr.to_excel(file_name_dir))
    except Exception as e:
        bad_tickers_file_dir = config.BAD_TICKERS
        f = open(bad_tickers_file_dir, "a")
        f.write(symbol)
        f.write('\n')

def sql_insert(symbol, strike, volume, open_interest, exp_date):
    conn_string = ('Driver={SQL Server};'
                   'Server={0};'
                   'Database={1};'
                   'Trusted_Connection=yes;').format(config.SERVER, config.DATABASE)
    conn = pyodbc.connect(conn_string)
    cursor = conn.cursor()
    insert_string = """INSERT INTO dbo.options (Ticker, Strike, Volume, OpenInterest, expDate) VALUES (?, ?, ?, ?, ?)"""
    cursor.execute(insert_string, symbol, strike, volume, open_interest, str(exp_date))
    conn.commit()

def get_exp_dates(symbol):
    url = "https://finance.yahoo.com/quote/" + symbol + "/options?p=" + symbol
    chromedriver = init_selenium()
    chromedriver.get(url)
    # Yahoo Finance options dropdown class name (find better way to do this)
    select_dropdown = chromedriver.find_element_by_css_selector("div[class='Fl(start) Pend(18px)'] > select")
    options_list = [x for x in select_dropdown.find_elements_by_tag_name("option")]
    dates = []
    for element in options_list:
        dates.append(element.get_attribute("value"))
    return dates

def read_ticker_file():
    file1 = open(config.TICKER_FILE, 'r')
    lines = file1.readlines()
    count = 0
    # loop to read each ticker in file
    for line in lines:
        count += 1
        line = line.strip('\n')
        line = line.strip()
        yfin_options(line)

if __name__ == "__main__":
    main()

异步改造方案与思路

先给你说清楚:你之前的改造没生效,是因为只改了外层函数的async标记,但底层所有IO操作(网络爬取、数据库插入、文件读写)都是同步阻塞的,相当于还是串行执行,根本没用到异步的并行能力。要真正提升效率,得把整个IO链路都改成异步,再用异步任务组来并行处理多个股票。

核心改造步骤

1. 替换同步依赖为异步版本

原来的工具都是同步的,得换成异步替代:

  • 网络请求:把selenium(慢到离谱的同步浏览器工具)换成aiohttp直接请求Yahoo的公开API,速度快N倍,还适合异步并行;yahoo_fin是同步库,直接用aiohttp封装Yahoo的API就行,不用依赖它。
  • 数据库操作:把pyodbc换成asyncio-pyodbc,还要用连接池复用数据库连接(你原来每次插入都新建连接,这是极大的性能浪费)。
  • 文件读写:用aiofiles替代原生的open,实现异步文件操作。

先装依赖:

pip install aiohttp asyncio-pyodbc aiofiles python-dotenv

2. 重构核心逻辑为异步函数

下面是关键部分的改造示例,你可以照着这个思路调整你的代码:

异步数据库连接池

import asyncio
import asyncio_pyodbc
import config

async def init_db_pool():
    # 创建数据库连接池,复用连接提升插入效率
    pool = await asyncio_pyodbc.create_pool(
        dsn=f"Driver={{SQL Server}};Server={config.SERVER};Database={config.DATABASE};Trusted_Connection=yes;",
        min_size=5,  # 最小连接数
        max_size=20  # 最大连接数,根据你的数据库配置调整
    )
    return pool

异步获取期权数据(替代selenium和yahoo_fin)

直接调用Yahoo的官方API,比爬页面高效太多:

import aiohttp
from datetime import datetime

async def fetch_exp_dates(session, symbol):
    """异步获取股票的期权到期日期戳"""
    url = f"https://query1.finance.yahoo.com/v7/finance/options/{symbol}"
    async with session.get(url) as resp:
        if resp.status != 200:
            return []
        data = await resp.json()
        # 提取到期日期的时间戳
        return data['optionChain']['result'][0]['expirationDates']

async def fetch_calls(session, symbol, date):
    """异步获取指定到期日的看涨期权数据"""
    url = f"https://query1.finance.yahoo.com/v7/finance/options/{symbol}?date={date}"
    async with session.get(url) as resp:
        if resp.status != 200:
            return []
        data = await resp.json()
        calls = data['optionChain']['result'][0]['options'][0]['calls']
        # 整理需要插入数据库的字段
        return [
            {
                'symbol': symbol,
                'strike': call['strike'],
                'volume': call.get('volume', 0),  # 处理可能为空的字段
                'open_interest': call.get('openInterest', 0),
                'exp_date': datetime.fromtimestamp(date)
            }
            for call in calls
        ]

异步批量插入数据库

绝对不要单条插入,改成批量插入,这比异步还能提升更多性能:

async def batch_insert_options(pool, options_data):
    """批量插入期权数据到数据库"""
    if not options_data:
        return
    async with pool.acquire() as conn:
        async with conn.cursor() as cur:
            insert_query = """
                INSERT INTO dbo.options (Ticker, Strike, Volume, OpenInterest, expDate)
                VALUES (?, ?, ?, ?, ?)
            """
            # 准备批量插入的参数
            params = [
                (d['symbol'], d['strike'], d['volume'], d['open_interest'], d['exp_date'])
                for d in options_data
            ]
            await cur.executemany(insert_query, params)
            await conn.commit()

异步处理单个股票

import logging
import aiofiles

async def process_ticker(session, pool, symbol):
    """异步处理单个股票的期权数据爬取与插入"""
    logging.basicConfig(filename='yfin.log', level=logging.INFO)
    try:
        exp_dates = await fetch_exp_dates(session, symbol)
        all_calls = []
        for date in exp_dates:
            calls = await fetch_calls(session, symbol, date)
            all_calls.extend(calls)
            # 异步保存到Excel(可选,pandas是同步的,这里可以用线程池包装,或者暂时保留同步写法)
            # df = pd.DataFrame(calls)
            # df.to_excel(f"C:\\temp\\rh\\options{symbol}{date}.xlsx")
        if all_calls:
            await batch_insert_options(pool, all_calls)
        logging.info(f"✅ 成功处理股票: {symbol}")
    except Exception as e:
        logging.error(f"❌ 处理股票{symbol}失败: {str(e)}")
        # 异步记录失败的股票
        async with aiofiles.open(config.BAD_TICKERS, 'a', encoding='utf-8') as f:
            await f.write(f"{symbol}\n")

主函数异步化(核心并行逻辑)

用asyncio.gather并行处理所有股票,同时用Semaphore控制并发数(避免被Yahoo封IP):

async def main():
    # 初始化数据库连接池和aiohttp会话
    pool = await init_db_pool()
    async with aiohttp.ClientSession() as session:
        # 异步读取股票文件
        async with aiofiles.open(config.TICKER_FILE, 'r', encoding='utf-8') as f:
            lines = await f.readlines()
        # 整理股票列表,去掉空行和换行符
        tickers = [line.strip() for line in lines if line.strip()]
        
        # 限制并发数:同时处理10个股票,可根据Yahoo的限制调整(别太高,容易被封)
        semaphore = asyncio.Semaphore(10)

        # 包装函数,实现并发控制
        async def bounded_process(ticker):
            async with semaphore:
                await process_ticker(session, pool, ticker)

        # 并行执行所有股票的处理任务
        await asyncio.gather(*[bounded_process(ticker) for ticker in tickers])
    
    # 关闭数据库连接池
    pool.close()
    await pool.wait_closed()

if __name__ == "__main__":
    asyncio.run(main())

3. 关键注意事项

  • 并发数控制:一定要用Semaphore限制同时处理的股票数量,建议从10开始慢慢调整,太高会被Yahoo的API封禁IP。
  • 批量操作优先:数据库单条插入的效率极低,批量插入能把速度提升几十倍,这个优化比异步还重要。
  • 错误隔离:异步环境下每个任务的异常要单独处理,确保一个股票处理失败不会影响其他股票。
  • 日志优化:你原来的logging.basicConfig调用了两次,后面的会覆盖前面的,建议只在程序开头配置一次日志。

内容的提问来源于stack exchange,提问作者user2880722

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.29 22:38:11