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

并行Python API数据采集脚本性能优化及磁盘写入改进问询

问题描述与性能优化方案

问题背景

我运行着一个从外部API获取数据的Python脚本,接收任意ID编号和产品名称作为输入,将API返回的数据打印并保存至文件。数据吞吐量因产品名称而异,大致为每秒1行(一天中不同时段速度会更慢)。

可以针对多个产品并行运行该脚本,但发现程序持续运行时CPU负载逐渐升高(其他软件占用负载有限),在8核i7-8650 ThinkPad(Ubuntu 22.04系统)上,CPU负载曾高达4。我猜测磁盘写入(如df.to_csv(...))是主要瓶颈,想了解如何优化代码性能,实现数据的流式持续写入。

原始代码核心问题分析

原始代码的性能瓶颈主要来自两点:

  • 内存与CPU开销:每次收到API数据时,都会创建新的DataFrame并通过pd.concat合并到全局DataFrame。随着数据量累积,内存占用持续攀升,且concat本身是高开销操作,会导致CPU负载逐步升高。
  • 磁盘IO浪费:每次合并后调用df.to_csv重写整个文件,意味着所有历史数据都要被重复写入,磁盘IO开销极大,进一步推高CPU负载。

优化后的代码实现

优化思路是直接流式写入CSV文件,彻底抛弃用Pandas维护大内存DataFrame的方案,每次仅写入新的一行数据:

from client import *
from wrapper import *
from contract import *
import time
import csv
import datetime
import threading
import sys
import os

# 接收命令行参数
cl = sys.argv[1]
sym = sys.argv[2]


class TestApp(EClient, EWrapper):

    def __init__(self):
        EClient.__init__(self, wrapper=self)

    def reqIds(self, numIds: int):
        return super().reqIds(numIds)

    def contractDetails(self, reqId: int, contractDetails: ContractDetails):
        super().contractDetails(reqId, contractDetails)
        print(f"合约详情: {contractDetails}")

    def contractDetailsEnd(self, reqId: int):
        print("合约详情获取完毕")
        self.disconnect()
        
    def updateMktDepth(self, reqId: TickerId, position: int, operation: int, side: int, price: float, size: int):
        super().updateMktDepth(reqId, position, operation, side, price, size)
                
        current_time = datetime.datetime.now()
        data = [current_time, sym, position, operation, side, price, size]        
        print(f'{current_time}: 请求ID: {reqId} 标的: {sym} 位置: {position} 操作: {operation} 方向: {side}, 价格: {price} 数量 {size}')
               
        # 以追加模式写入单行数据
        with open(f'/home/chris/data/{sym}_l2.csv', mode='a', newline='') as file:
            wr = csv.writer(file)
            wr.writerow(data)


def main():
    try:            
        app = TestApp()
        app.connect("127.0.0.1", 7496, cl)
        print('连接成功')

        t = threading.Thread(name=f'API_worker_{sym}', target=app.run)
        t.start()
        print("API线程已启动")
    
        c = Contract()
        c.localSymbol = sym
        c.secType = 'FUT'
        c.exchange = 'CME'
        c.currency = 'USD'     
        time.sleep(1)

        # 预处理文件:重命名旧文件并创建带表头的新文件
        file_path = f'/home/chris/data/{sym}_l2'
        timestamp = datetime.datetime.now().strftime('%Y%m%d%H%M')
        headers = ['时间','标的','位置','操作','方向','价格','数量']

        if os.path.isfile(f'{file_path}.csv'):
            os.rename(f'{file_path}.csv', f'{file_path}_{timestamp}.csv')
            
            with open(f'{file_path}.csv', 'w', newline='') as file:
                wr = csv.writer(file)
                wr.writerow(headers)
        else:
            with open(f'{file_path}.csv', 'w', newline='') as file:
                wr = csv.writer(file)
                wr.writerow(headers)

        app.reqMktDepth(cl, c, 20, 0, [])  
    except KeyboardInterrupt:
        print('用户中断,程序结束')
        app.disconnect()

if __name__ == "__main__":
    main()

优化效果说明

  • 内存占用稳定:不再维护全局DataFrame,每次仅处理单行数据,内存占用始终处于低水平。
  • 磁盘IO开销骤降:每次仅追加写入一行数据,避免了重写整个文件的高开销操作,CPU负载会明显降低并保持稳定。
  • 功能完整性保留:原有的文件预处理(重命名旧文件、写入表头)和数据打印功能完全保留,同时实现了高效的流式写入。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 21:37:53