使用for循环处理10K标的OHLC数据时程序崩溃的解决方法
解决10K标的OHLC数据全量写入数据库不崩溃的方案
核心问题分析
原程序崩溃的主要原因是内存过载:把10K标的的所有数据先缓存到temp_data列表中,最后一次性合并写入,导致内存被大量占用;其次是逐bar创建单行DataFrame的低效写法,进一步加剧内存消耗;另外缺少异常处理,单个标的的API请求失败会直接中断全流程。
具体优化方案
1. 分批写入数据库,释放内存
不再把所有数据存到内存,处理完一批标的就写入数据库并清空临时缓存,避免内存持续累积。
2. 优化单标的数据转换逻辑
避免循环每个bar创建单行DataFrame,改为一次性将单个标的的所有bar数据转换为完整DataFrame,大幅减少内存碎片和开销。
3. 添加异常处理,容错单个标的失败
对API请求、数据转换过程加异常捕获,避免单个标的出错导致整个程序崩溃,同时记录失败的标的方便后续补拉。
4. 控制API请求速率,避免触发限流
Alpaca API有调用频率限制,连续请求10K标的可能触发限流,添加短延迟或重试机制降低风险。
修改后的完整代码
import config from alpaca_trade_api.rest import REST, TimeFrame import sqlite3 import pandas as pd import datetime from dateutil.relativedelta import relativedelta import time import traceback # 初始化参数 BATCH_SIZE = 500 # 每批写入的标的数量,可根据内存调整 RETRY_DELAY = 2 # API请求失败后的重试间隔(秒) MAX_RETRIES = 3 # 单个标的最大重试次数 start_date = (datetime.datetime.now() - relativedelta(years=2)).date() start_date = pd.Timestamp(start_date, tz='America/New_York').isoformat() end_date = pd.Timestamp(datetime.datetime.now(), tz='America/New_York').isoformat() conn = sqlite3.connect('allStockData.db') api = REST(config.api_key_id, config.api_secret, base_url=config.base_url) origin_symbols = pd.read_sql_query("SELECT symbol, name from stock", conn) df_dict = origin_symbols.to_dict('records') # 预先创建表结构(避免后续写入出错) sample_df = pd.DataFrame(columns=['date', 'symbol', 'open', 'high', 'low', 'close', 'volume', 'vwap']) sample_df.to_sql('daily_ohlc_init', if_exists='replace', con=conn, index=False) startTime = datetime.datetime.now() batch_data = [] failed_symbols = [] for idx, key in enumerate(df_dict, 1): symbol = key['symbol'] print(f"Processing {idx}/{len(df_dict)}: {symbol}") success = False retries = 0 while not success and retries < MAX_RETRIES: try: # 获取单个标的的所有K线数据 barsets = list(api.get_bars_iter(symbol, TimeFrame.Day, start_date, end_date)) if not barsets: print(f"No data found for {symbol}") success = True continue # 一次性转换为DataFrame,避免逐行创建的低效写法 bar_data = { 'date': [bar.t.date() for bar in barsets], 'symbol': [symbol] * len(barsets), 'open': [bar.o for bar in barsets], 'high': [bar.h for bar in barsets], 'low': [bar.l for bar in barsets], 'close': [bar.c for bar in barsets], 'volume': [bar.v for bar in barsets], 'vwap': [bar.vw for bar in barsets] } df_single = pd.DataFrame(bar_data) batch_data.append(df_single) # 达到批次大小则写入数据库 if len(batch_data) >= BATCH_SIZE: pd.concat(batch_data).to_sql('daily_ohlc_init', if_exists='append', con=conn, index=False) batch_data = [] # 清空批次缓存 print(f"Written batch up to {idx} symbols") success = True time.sleep(0.1) # 控制请求速率,避免触发API限流 except Exception as e: retries += 1 print(f"Failed to fetch {symbol} (retry {retries}/{MAX_RETRIES}): {str(e)}") traceback.print_exc() time.sleep(RETRY_DELAY) if not success: failed_symbols.append(symbol) # 写入剩余的批次数据 if batch_data: pd.concat(batch_data).to_sql('daily_ohlc_init', if_exists='append', con=conn, index=False) # 记录失败的标的,方便后续补拉 if failed_symbols: pd.DataFrame({'symbol': failed_symbols}).to_sql('failed_symbols', if_exists='replace', con=conn, index=False) print(f"Failed to process {len(failed_symbols)} symbols, saved to failed_symbols table") endTime = datetime.datetime.now() print(f'Time elapsed to pull data was {endTime - startTime}') conn.close()
额外说明
BATCH_SIZE可根据你的内存情况调整,内存较小就设为200左右,内存充裕可设为1000- 若频繁遇到限流,可适当增大
time.sleep(0.1)的延迟时间 failed_symbols表会记录拉取失败的标的,后续可以单独编写脚本补拉这些数据
内容的提问来源于stack exchange,提问作者a7dc
相关产品推荐
相关产品推荐

