如何优化基于Concurrent Futures的Python MySQL数据处理脚本性能
哇,10万+IP跑10小时确实够闹心的,我帮你拆解下核心瓶颈,从数据库到代码逻辑给你梳理几个能快速提效的优化方向:
这部分是大头,你的脚本现在每个IP/设备都要发起多次独立查询,10万量级下数据库的开销直接拉满了。
把计算逻辑从Python移到数据库,用批量聚合代替循环单查
现在你在Python里循环每个IP、每个设备,多次查询数据库再计算,完全可以用一条复杂的SQL直接算出所有IP+device的最终结果。比如用JOIN关联SSD_fail_predictor、SSD_size和lba_written,再用GROUP BY做聚合计算,把analyse_data里的逻辑全部搬到SQL里。这样数据库的并行计算能力能发挥出来,比Python循环快N倍。给查询字段加联合索引
检查下这几个表的索引情况:lba_written表:给machineip, device, attribute, time_stamp加联合索引,避免全表扫描:CREATE INDEX idx_lba_machine_device_attr_time ON lba_written(machineip, device, attribute, time_stamp);SSD_fail_predictor表:给machine_ip, device加联合索引:CREATE INDEX idx_ssd_ip_device ON SSD_fail_predictor(machine_ip, device);
索引能让数据库快速定位数据,直接把单查询时间从秒级压到毫秒级。
替换字符串拼接SQL为参数化查询
你现在用format拼接SQL,不仅有SQL注入风险,还会让数据库无法复用执行计划。改用参数化查询,让MySQL缓存执行计划,提升重复查询的速度:# 原写法 # db_engine.execute("""select ... inet_ntoa(machine_ip)='{0}' and device='{1}'""". format(ip, device)) # 改成参数化 db_engine.execute("""select model,network,region_num,(select size from SSD_size where model=a.model) from SSD_fail_predictor a where inet_ntoa(machine_ip)=%s and device=%s""", (ip, device))
现在的多进程/线程并没有发挥最大作用,反而可能因为资源浪费拖慢速度。
复用数据库连接池,不要每个IP都新建连接
你现在在process_ips里每个IP都创建新的db_engine然后销毁,连接的创建销毁开销极大。可以给每个进程初始化一次连接池,复用连接:# 定义进程初始化函数 def init_process(): global db_engine # 创建带连接池的引擎 db_engine = sqlalchemy.create_engine( "mysql+pymysql://root:hdfs@xxx.xx.xx.xxx/hardware_perf", pool_size=10, # 按需调整 pool_recycle=3600 ) # 在main里使用ProcessPoolExecutor时传入初始化函数 with ProcessPoolExecutor(max_workers=args.threads, initializer=init_process) as executor: results = executor.map(process_ips, ips_batch)这样每个进程复用一个连接池,避免频繁创建销毁连接的开销。
调整并发数和批量大小
现在的batch_size=50和threads=8可能不是最优配置。可以测试更大的batch_size(比如200-500),同时根据数据库的max_connections调整线程/进程数——如果数据库最大连接数是100,那max_workers设为15-20比较合适,太多会导致数据库连接阻塞。合并小任务,减少进程间通信开销
多进程的IPC(进程间通信)有开销,如果每个任务只处理一个IP,开销占比会很高。可以把多个IP打包成一个任务,比如每个任务处理10个IP,减少进程切换和IPC的开销。
减少频繁的文件IO操作
现在每个错误情况都要打开写入文件,频繁的IO会拖慢速度。可以把错误信息先缓存到内存列表里,批量写入——比如每处理1000个IP写一次,或者所有任务完成后统一写入。简化attribute查询逻辑
你现在查所有distinct attribute再取最后一个,不如直接取最新的一条:# 原逻辑 # results = db_engine.execute( "select distinct attribute from lba_written where machineip='{0}' and device='{1}' order by time_stamp".format(str(ip), str(device))) # 改成直接取最新的 result = db_engine.execute("SELECT attribute FROM lba_written WHERE machineip=%s AND device=%s ORDER BY time_stamp DESC LIMIT 1", (ip, device)).fetchone() if result: attribute = result[0] else: # 处理无数据的情况 with open("nullip.txt", 'a') as f: f.write("IP: {0} with device: {1} has no attributes\n".format(ip, device)) return 1
如果你的数据库支持异步,可以用aiomysql结合asyncio来处理查询,异步IO在IO密集型任务(这里大部分时间是数据库查询)上的性能比同步多进程/线程更高,能进一步压缩耗时。
内容的提问来源于stack exchange,提问作者Yogesh Chandra

