告别for循环:批量行情架构如何把全市场5000只股票扫描耗时压到1秒
2026/9/15 4:34:55 网站建设 项目流程

先把结论放在前面:如果你还在用for循环,把5000只股票挨个请求一遍行情接口,那5分钟只是最理想状态下的耗时;行情一波动、连接一超时、接口一限流,10分钟扫不完也正常。我最初做QuantDash这个全市场行情扫描项目时,第一版就是这样被“慢性崩溃”折磨的。后来把架构改成批量行情拉取,配合本地缓存增量计算,全市场扫描的常规耗时直接压到了1秒上下。这篇文章专门讲清楚:为什么循环请求会慢到崩溃,批量行情架构到底改了什么,以及你在复现这套方案时会踩到哪些真实存在的坑。

1. 串行轮询5000只股票的病灶诊断

1.1 一次请求的往返成本,比你想象的高一个量级

很多人对“请求一次行情接口”的耗时没有概念。单独调一只股票的实时报价,在本地网络良好的情况下,单次HTTP往返大概是50到150毫秒。如果你用的是免费公开接口,DNS解析、TCP握手、TLS握手都会算进去,部分接口甚至还会在服务端做一次行情快照的组装,耗时再往上走。

我们来算一笔简单的账:

  • 单次请求平均耗时:按80ms算;
  • 5000只股票串行请求:5000 × 80ms = 400秒;
  • 折算成分钟:大约6.7分钟。

如果单次请求因为超时被重试,每重试一次就多出至少100ms,而且超时往往发生在行情剧烈波动时,此时网络拥塞、服务端压力大,重试失败的概率更高。串行模式下的耗时不是线性增长,而是指数级恶化。

这还只是“请求”本身。拿到JSON以后你还要解析、清洗、落库,这些CPU和I/O时间在实际运行中同样计入总耗时。

1.2 循环请求放大延迟的方式:串行、超时重试与队头阻塞

循环请求真正的病灶有三个。

第一个是串行等待。前一只股票没有返回结果,后一只股票的请求就发不出去。你可能会想,用多线程不就行了吗?但多线程解决的是“等待不阻塞”,解决不了“接口整体吞吐有限”的问题。

第二个是超时重试。默认情况下,requests.get如果不加超时,会一直挂在那里。某个接口偶尔抖动一下,你的循环就会卡在那一只股票上,后面的4999只全部堵死。就算你加了超时,一旦触发重试,这部分时间照样会被放大。

第三个是队头阻塞。如果你用一个共享的解析队列,前面一条脏数据(比如某个股票返回了空字符串、字段缺了price)解析失败,会拖累整个消费链路,后面攒了一堆待处理数据,最终导致内存堆积、程序越来越慢,看起来就像“崩溃”了一样。

1.3 为什么简单加线程池也不可靠

有人会说,那我开20个线程,每个线程处理250只股票,总耗时不是能降到20到30秒吗?

在理想环境下确实可以。但真实行情环境下,问题往往出在接口限流上。免费公开接口对单个IP的并发和频率都有隐性限制:

  • 你瞬间发起20个并发,每个线程每秒请求数十次;
  • 服务端大概率会返回429 Too Many Requests或直接断开连接;
  • 线程池里的线程并不会“友好”地退避,它们会反复重试,结果整个线程池被无效请求塞满;
  • 连接池也被占满,新的请求只能排队等待,甚至出现超过urllib3默认连接池上限的报错。

更隐蔽的一个问题是连接死锁:线程池有8个线程,但连接池只有5个可用连接,8个线程同时去抢连接,谁都拿不到,程序既不报错也不退出,CPU却被打满。这就像经典的哲学家就餐问题里“循环等待”那一步——每个资源持有人在等别人释放资源,整个系统卡死。我在初版QuantDash里就遇到过一模一样的情况,后来才明白,必须把“并发”和“限流”都交给一个统一调度的模块去管,而不是放任线程自己去抢。

2. QuantDash批量行情架构的核心:一次拉一批,控制权留在调度器手里

2.1 把“一只一只问”拆成“一片一片拿”:批量行情接口的正确用法

现在市面上大部分行情接口,无论是免费还是收费的,都支持一次请求多个股票代码,常见的批量上限在50到100个代码左右。这才是全市场扫描的正确打开方式:不是5000次请求,而是把代码池切成批次,一次拉取一片。

QuantDash的调度模块核心逻辑非常简单:

class QuoteDispatcher: def __init__(self, codes: list[str], batch_size: int = 80, max_workers: int = 8): self.codes = codes self.batch_size = batch_size self.max_workers = max_workers self.pending = deque(codes) def _fetch_batch(self, batch: list[str]) -> dict: # 批量行情接口示例,一次性传入多个代码 resp = requests.get( QUOTE_HTTP_ENDPOINT, params={"codes": ",".join(batch)}, timeout=(3, 5), ) resp.raise_for_status() return normalize_quote(resp.json()) def run(self): with ThreadPoolExecutor(max_workers=self.max_workers) as executor: futures = [] while self.pending: batch = [] for _ in range(self.batch_size): if self.pending: batch.append(self.pending.popleft()) futures.append(executor.submit(self._fetch_batch, batch)) for future in as_completed(futures): yield future.result()

这里有几个关键点:

  • timeout=(3, 5)表示连接超时3秒,读取超时5秒,避免单次请求无限挂起;
  • 批次大小80,意味着5000只股票只需要约63个请求;
  • max_workers=8,是可控的并发上限,不是无脑开100个线程;
  • 每个批次内部是一次大的HTTP往返,而不是5000次小往返。

很多人把优化重点放在“异步”上,其实真正的收益大头在“批量”。异步能让你在等待时不阻塞CPU,但该发出的请求数一个都不会少;批量则是直接把请求数量砍掉一个量级。

2.2 并发数、批次大小与重试窗口:调度参数的平衡点

调度参数不是什么神秘的东西,但调不好就会反复被限流。我试过一组比较稳的参数组合:

参数建议值理由
批量大小60-100低于30,请求数太多;高于接口限制,直接报错
最大并发数8-16GIL和连接池压力可控,服务端不容易触发限流
连接超时3秒给足DNS和TCP握手时间,但不无限等待
读取超时5-10秒批量响应数据量大,太短容易误杀
批次间退避0-0.2秒给接口一点喘息空间,大幅降低429概率
单批次失败重试2次超过2次直接跳过该批次,记录日志后续补单

这些参数不是死板的。如果你用的是高可用付费数据源,并发可以翻倍;如果你用的是免费公共接口,建议把退避调到0.3秒左右,跑全市场也不会慢到哪里去。

从实测来看,80只一批、8个并发、每批平均响应200ms,扫描完5000只股票大约需要:

63批 × 200ms ÷ 8并发 ≈ 1.6秒

再加上解析和入库的时间,整体控制在2秒内没有问题。如果响应压到100ms以下,并发提到16,那么1秒达成全市场首轮扫描是完全现实的。

2.3 请求层只做请求,行情层只做聚合:职责划分的重要性

初版QuantDash写得很随意,请求股票、解析数据、算指标、推送结果全堆在一个函数里,后续改参数特别痛苦,排查限流问题时更是无从下手。后来我把架构拆成了三层:

  • 接入层只负责“拿到一批原始行情JSON”;
  • 聚合层负责把批次数据拆成单股票快照,按symbol + datetime为维度写入内存;
  • 服务层负责指标计算、过滤条件筛选和结果推送。

每一层之间用队列连接,接入层不知道指标怎么算,服务层也不知道HTTP连接池是怎么管理的。这样做的直接好处是:当行情接口变了,比如从A接口切换到B接口,我只需改写接入层;当扫描策略变了,比如新增一个均线交叉条件,我只需在服务层加规则。

如果你在做一个稍大一点的个人项目,强烈建议从一开始就按这种职责拆分。不需要引入重框架,Python自带的queue.Queue或者concurrent.futures就足够支撑。

3. 从“拿全市场行情”到“1秒扫完全市场”的增量之道

3.1 历史K线先落地,实时快照只做增量合并

如果每次扫描都重新去请求5000只股票的完整K线历史,那神仙架构也快不了。真正的优化思路是:把“历史数据”和“实时增量”分开。

以日线级别的均线扫描为例:

  • 每个交易日收盘后跑一次批处理,把全市场的日线数据拉下来,存到本地SQLite或Parquet文件;
  • 盘中扫描时,不再重新拉历史,而是拉最近一个时间窗口的实时行情快照;
  • 把快照和本地历史数据按symbol + date拼接成最新的一条数据;
  • 基于拼接后的数据计算指标。

这样做的收益极大地摊薄了扫描成本。实时接口只需返回当前价格、涨跌幅、成交量这几个字段,而不是完整的历史OHLCV序列。

QuantDash里我专门加了一个SnapshotMerger

class SnapshotMerger: def __init__(self, store: LocalStore): self.store = store def merge(self, snapshot: dict) -> pd.DataFrame: code = snapshot["symbol"] local_df = self.store.load_history(code) # 用最新快照替换最后一行,而不是追加 local_df.iloc[-1, local_df.columns.get_loc("close")] = snapshot["price"] local_df.iloc[-1, local_df.columns.get_loc("volume")] = snapshot["volume"] return local_df

注意这里用的是“替换最后一行”,而不是“追加一行”。因为实时快照对应的是当前未收盘的K线,和本地最后一根K线是同一个时间窗口。如果你在盘中扫描时不断追加,K线数量会膨胀,指标计算的耗时也会越来越大。

3.2 把五千次指标计算改成一次向量化扫描

全市场扫描不仅要拿行情,还要算指标、跑过滤条件。如果你用循环遍历5000只股票,每只都算一遍RSI、MACD、均线,那耗时显然不会好看。

量化扫描的正确姿势是向量化。所谓向量化,就是不要对单只股票调用循环计算函数,而是把全市场的数据组织成一个大DataFrame或NumPy数组,对整个数组执行同一种运算。

举个例子,筛选出“5日均线大于10日均线,且今日涨幅超过3%”的股票,在QuantDash里是这样写的:

def screen_market(snapshot_map: dict[str, pd.DataFrame]) -> pd.DataFrame: rows = [] for code, df in snapshot_map.items(): # 每只股票计算一个尾部特征行 df = df.tail(10).copy() df["ma5"] = df["close"].rolling(5).mean() df["ma10"] = df["close"].rolling(10).mean() last = df.iloc[-1] rows.append({ "symbol": code, "ma5": last["ma5"], "ma10": last["ma10"], "pct_chg": last["pct_chg"], }) result = pd.DataFrame(rows) return result[(result["ma5"] > result["ma10"]) & (result["pct_chg"] > 3)]

这里rollling(5).mean()在每只股票的尾部窗口上计算,5000只股票加在一起仍然是毫秒级完成。真正的耗时大头是构造snapshot_map,也就是从本地缓存中读取每只股票的数据,这部分依赖的是本机磁盘I/O和内存速度,而不是网络。

如果连字典循环都想省掉,直接把所有股票拼成一个大DataFrame,用groupby("symbol")transform来算指标,那才是完全的向量化方案。不过实际测试下来,groupby拆组再合并的开销也不算小,对5000只股票来说,普通字典循环已经足够快。

3.3 扫描结果的缓存与推送,才是“1秒”能被肉眼感知的关键

架构优化到最后,用户体感才是最重要的。你的后端扫描确实1秒完成了,但结果经过WebSocket推送、前端图表渲染,如果链路不通畅,用户依然感受不到“秒级”。

QuantDash的处理方式是双缓存:

  • 第一层,内存缓存:扫描结果按scan_id存一份最新实例,供API查询;
  • 第二层,增量推送:每次扫描完成后,只把变化的股票列表diff推给前端,而不是全量推5000条。

比如上一轮扫描命中了50只,这一轮命中47只,前端只需要收到“新增2只,移除5只”的增量信息。这个设计把网络传输和前端DOM更新都压到了最小,实际体验上就是“行情K线还在跳,筛选结果已经跟着变了”。

4. 实测对比:耗时、请求数与资源占用

4.1 同数据集下的串行/单并发/批处理对比

我在一台4核8G的Linux服务器上做过一次实测,样本是沪深A股约5000个代码,接口使用同一个公开HTTP行情源,分别用三种方式跑完整扫描:

方式请求次数总耗时备注
纯串行循环50007分52秒期间出现2次超时,触发重试
20线程并发循环500043秒触发大量429,重试增多
QuantDash批量+8并发631.8秒无429,日志无异常

这个对比很能说明问题。20线程并发循环没有把耗时压进10秒以内,原因是请求数太多,触发了服务端限流;而批量方案把请求次数降到63次,天然就远离了限流红线。

另外注意“请求次数”这一列。性能优化不能只看耗时,请求次数直接影响你对接口的依赖强度和被封IP的风险。5000次请求和63次请求,对数据源服务商来说是完全不同的压力等级。

4.2 Docker部署下的性能观察

我在项目初期就决定用Docker部署这个股票实时数据服务,主要是为了让它在不同机器上的行为保持一致。Docker容器里的网络延迟和宿主相比几乎可以忽略,但有几个细节需要留意:

  • 容器内DNS解析有时会变慢,建议在docker-compose.yml里指定dns: 8.8.8.8或本机可用的DNS;
  • 内存限制建议至少开1GB,因为行情快照和DataFrame计算都比较吃内存;
  • 日志不要无脑输出,否则几天下来磁盘会被塞满。

一个可用的compose片段:

services: quantdash: image: quantdash:latest container_name: quantdash restart: unless-stopped environment: - QUOTE_BATCH_SIZE=80 - QUOTE_MAX_WORKERS=8 volumes: - ./data:/app/data ports: - "8080:8080"

./data目录是本地历史K线和缓存数据,挂载出来的目的是容器重建后不用重新拉全市场数据。

4.3 限流下的请求失败率与重试策略

哪怕调度参数再收敛,也拦不住某些接口在交易时段突然抽风。QuantDash的重试策略不是无脑重试,而是带退避和熔断的:

失败次数处理方式说明
第1次失败等待0.5秒后重试大概率是接口抖动
第2次失败等待2秒后重试加重退避
第3次失败放弃当前批次,记录日志避免引爆限流
连续失败超过10个批次熔断开关打开暂停拉取30秒,自动恢复

熔断机制在我实际跑盘中扫描时救了太多次。某个免费接口经常在开盘后第1分钟高负载报错,如果程序傻了似的一直重试,整个扫描任务就会被拖死。熔断后停下来等30秒,等接口缓过来再继续,反而更稳。

5. 我踩过的真实行情数据坑,每个都足以让扫描结果出错

5.1 时间戳不齐导致算错涨跌幅

第一次跑全市场扫描时,我发现某些股票的“涨跌幅”明显不对,排查了半天才发现是时间戳的问题。免费接口返回的字段里有time,但这个时间不一定对应行情快照的时间,有的代码返回的是本地服务器时间,有的返回的是交易所时间,还有的干脆是上个交易日的日期。

后来在normalize_quote函数里强制加了一层时间校准:所有快照统一以交易所的最近交易日和时区为准,如果时间字段无法解析,就放弃该条记录并告警。这个坑不解决,你后续算出来的所有指标都会带上隐性锚定误差。

5.2 停牌、ST与新股池的过滤

“全市场5000只”并不等于“5000只都能正常买入卖出”。停牌股没有实时行情,ST股有交易限制,次新股上市初期数据不足,这些如果不做分层过滤,扫描结果会让人困惑。

我在QuantDash里维护了一个可配置的股票池过滤器:

  • 剔除当前停牌的股票,基于接口返回的status或成交量为0且涨跌幅为0的特征判断;
  • 默认剔除ST、*ST股票,除非你明确开启“包含ST”开关;
  • 次新股按上市天数过滤,一般要求上市超过60个交易日再纳入指标扫描。

这些规则看着简单,但能帮你过滤掉一堆垃圾信号。尤其是实盘场景,一只停牌股被扫出来并推送到前端,用户点击却发现无法交易,体验非常差。

5.3 接口偶尔返回空串的熔断恢复

免费接口还有一个特色:偶尔返回完全空白的响应,HTTP状态码依然是200。你的代码如果只判断resp.status_code == 200,就会把空响应当成正常数据交给解析器,然后解析器抛出一个莫名其妙的结构错误。

处理方法是:在解析前先判断响应体长度和结构。如果响应体为空、或JSON解析出来是None、或缺少必要的行情字段,都当成当前批次失败处理,走重试逻辑。这样可以把偶发性空响应隔离在单一批次内,不会影响全市场扫描的完成质量。

5.4 试试增量续跑,而不是每次全量重来

另一类问题发生在程序重启后。如果全市场历史K线没存下来,每次启动都要重新拉一遍大数据量接口,不仅慢,还容易再次触发限流。解决思路是增量续跑:

  • 每次落库时记录每只股票的最新K线日期;
  • 重启后扫描本地元数据,只拉取比最近日期新的K线;
  • 如果本地库是空的,才触发全量历史拉取。

这个方案在我的实测中把重启恢复时间从30分钟压缩到了几十秒,尤其是跨周末重启项目时特别有效。

6. 想在小机器上复现:最小落地清单

6.1 选型与安装

一套能跑的QuantDash最小组合不需要大而全的框架,我自己用的组合是:

  • Python 3.11,加上requestspandaspyarrow
  • SQLite存储历史K线,Parquet存储离线全量快照;
  • 一个定时触发脚本,用schedule或系统cron每5秒触发一次增量扫描;
  • 前端仪表盘用WebSocket接后端扫描结果,不轮询全量接口。

初次搭建时,先把“批量请求+本地缓存”跑通,再逐步加指标和推送。不要一上来就把指标库、数据库、消息队列全堆上,那样你会分不清慢到底是慢在哪个环节。

6.2 定时任务和重启策略

交易时段和非交易时段要区别对待:

  • 交易时段运行增量行情拉取,频率可以设为每隔3到5秒;
  • 午间休市和收盘后停止实时拉取,转为历史K线补全任务;
  • 每天收盘后自动执行一次全市场历史K线更新,并清理过期缓存文件。

Docker里配合restart: unless-stopped,容器挂掉会自动拉起来。为了应对极端情况,我在启动脚本里加了一个“首次启动做全量数据校验”的阶段,确保本地数据完整后再开启对外服务。

6.3 可以往哪些方向扩展

架构跑通之后,低成本的延展方向其实很多,我自己就在现有基础上加了三个功能:

  • 预警推送:扫描结果中出现满足条件的股票,通过钉钉或企业微信机器人推送到手机;
  • 自定义指标插件:把均线、RSI、MACD等封装成可插拔函数,新增指标不用改主流程;
  • 收益曲线生成:把买入信号对应的持仓记录和收盘行情回放,生成收益曲线,便于复盘策略参数。

这套批量行情架构本身不挑语言,你换成Go或Java也没问题。关键还是我前面反复强调的那几个点:批量代替串行、增量代替全量、向量化代替循环计算、熔断代替盲目重试。

我今天能跑出1秒全市场扫描,不是因为哪一步操作特别高端,而是把该做对的底层选择都做对了。你如果正打算做类似的项目,不妨先从这个最小清单入手,真的跑起来以后,你会比看任何文章都更直观地理解那5分钟究竟浪费在哪里。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询