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
相关产品推荐
相关产品推荐

