并发请求下Pandas写入CSV仅存单条数据的问题排查求助
问题描述
我有一个调用Flask应用API的脚本,希望将请求的statuscode和elapsed time存入Pandas DataFrame并写入CSV文件,用于生成响应时间可视化图表。但目前CSV文件中仅保留一条记录,而打印时能看到所有请求的状态码和耗时;我尝试编写了write_df函数,但最终在send_api_request函数中直接处理请求变量,使用ThreadPoolExecutor并发执行请求,代码如下:
import requests import datetime import concurrent.futures import csv import pandas as pd HOST = 'http://127.0.0.1:5000' API_PATH = '/' ENDPOINT = HOST + API_PATH MAX_THREADS = 8 CONCURRENT_THREADS = 10 csv_path = "flasktests.csv" try: file = open(csv_path, 'w', newline='') writer = csv.writer(file) except: print("error opening or writing to the CSV file!") def send_api_request(): try: #print ('Sending API request: ', ENDPOINT) r = requests.get(ENDPOINT) if r.status_code == 200: #print('Received: ', r.status_code, r.elapsed) responses = {"statuscode":[r.status_code], "elapsed time": [r.elapsed]} statuscode = r.status_code elapsedtime = r.elapsed print(statuscode, elapsedtime) df = pd.DataFrame([statuscode,elapsedtime], columns=["statuscode","elapsed time"]) df.to_csv(csv_path, index=False) elif r.status_code == 417: print('Received error code:', r.status_code, r.json()) except Exception as e: print("error",str(e)) def write_df(statuscode, elapsedtime): print(statuscode,elapsedtime) df = pd.DataFrame({"statuscode":[statuscode], "elapsed time": [elapsedtime]}) df.to_csv(csv_path, index=False) print(df) with concurrent.futures.ThreadPoolExecutor(max_workers=MAX_THREADS) as executor: futures = [executor.submit(send_api_request) for x in range (CONCURRENT_THREADS)] executor.shutdown(wait=True)
想请教问题出在哪里?
问题原因
- CSV写入覆盖:每次调用
df.to_csv()默认用mode='w'(覆盖模式),每个线程写入时都会清空文件只留当前行,最终只剩最后一条记录。 - 并发文件竞争:多线程同时写同一个文件,会导致数据丢失、覆盖甚至文件损坏。
- DataFrame结构错误:
pd.DataFrame([statuscode,elapsedtime], columns=["statuscode","elapsed time"])是把两个值拆成两行,而不是同一行的两列,数据结构完全错误。 - 无效文件操作:开头打开的
file和writer变量完全没用到,属于冗余代码。
修复方案
核心逻辑改成先收集所有线程的请求结果,再统一写入CSV,彻底解决覆盖和竞争问题,同时修正DataFrame创建逻辑:
- 让
send_api_request返回每条请求的结果(包括异常和错误状态码); - 等待所有线程执行完毕后,收集所有结果并转换为DataFrame;
- 一次性将完整的DataFrame写入CSV。
修复后的代码
import requests import concurrent.futures import pandas as pd HOST = 'http://127.0.0.1:5000' API_PATH = '/' ENDPOINT = HOST + API_PATH MAX_THREADS = 8 CONCURRENT_THREADS = 10 csv_path = "flasktests.csv" def send_api_request(): try: r = requests.get(ENDPOINT) # 把耗时转成秒数,方便后续可视化处理 elapsed_sec = r.elapsed.total_seconds() if r.status_code == 200: print(r.status_code, elapsed_sec) return {"statuscode": r.status_code, "elapsed time": elapsed_sec} elif r.status_code == 417: err_msg = r.json() print('Received error code:', r.status_code, err_msg) return {"statuscode": r.status_code, "elapsed time": elapsed_sec, "error": err_msg} # 处理其他状态码 print('Received status code:', r.status_code) return {"statuscode": r.status_code, "elapsed time": elapsed_sec} except Exception as e: err_msg = str(e) print("error:", err_msg) return {"statuscode": None, "elapsed time": None, "error": err_msg} # 收集所有线程的结果 results = [] with concurrent.futures.ThreadPoolExecutor(max_workers=MAX_THREADS) as executor: futures = [executor.submit(send_api_request) for _ in range(CONCURRENT_THREADS)] for future in concurrent.futures.as_completed(futures): res = future.result() if res: results.append(res) # 统一写入CSV df = pd.DataFrame(results) df.to_csv(csv_path, index=False) print(f"成功写入{len(results)}条记录到{csv_path}")
关键说明
- 将
r.elapsed转为秒数,后续做可视化时更易处理; - 新增异常和非200/417状态码的记录,让CSV数据更完整;
- 移除冗余的文件打开代码,简化逻辑;
- 单线程统一写入CSV,彻底避免多线程文件操作的冲突。
内容的提问来源于stack exchange,提问作者YannickMetz
相关产品推荐
相关产品推荐

