如何为大DataFrame应用自定义函数并每N行写入一次CSV
实现方案及代码示例
你提出的批次拆分+增量写入的思路完全符合该场景的需求,既可以避免长时间运行中途故障导致全部数据丢失,也不会占用过多内存资源,以下是可直接落地的实现方案:
1. 核心实现逻辑
- 提前设定单批次处理行数
N,可根据请求耗时、容错需求灵活调整,通常建议设置为100~500行 - 对原始DataFrame按固定行数拆分批次,使用pandas内置的分组逻辑即可实现
- 遍历每个批次,对目标列应用自定义GET请求函数
- 每个批次处理完成后追加写入CSV,仅第一个批次写入表头,后续批次跳过表头避免重复
- 可额外增加断点续跑逻辑,运行中断后无需重新处理已完成的批次
2. 可复用代码示例
import pandas as pd import requests import os from tqdm import tqdm # 可选,用于显示处理进度 # 自定义GET请求函数,根据实际业务逻辑调整 def custom_get_request(target_value): try: # 此处替换为你的实际请求逻辑 resp = requests.get(target_value, timeout=15) resp.raise_for_status() # 可替换为你实际需要返回的内容,比如json解析后的指定字段 return resp.text except Exception as e: # 出错时返回异常信息,避免单个请求失败导致整个批次中断 return f"请求异常:{str(e)}" # 配置项 BATCH_SIZE = 200 # 每批次处理行数,可自行调整 INPUT_FILE = "your_source_data.csv" # 原始数据文件路径 OUTPUT_FILE = "processed_result.csv" # 结果输出路径 TARGET_COLUMN = "request_url" # 要应用自定义函数的列名 # 读取原始数据 df = pd.read_csv(INPUT_FILE) # 生成批次分组键,每N行对应同一个分组 df["batch_id"] = df.index // BATCH_SIZE total_batch = df["batch_id"].nunique() # ---------------------- 断点续跑逻辑(可选,推荐启用) ---------------------- processed_batch = set() if os.path.exists(OUTPUT_FILE): # 读取已处理的批次,避免重复运行 processed_df = pd.read_csv(OUTPUT_FILE, usecols=["batch_id"]) processed_batch = set(processed_df["batch_id"].unique()) # ------------------------------------------------------------------------- # 遍历批次处理 for batch_id, batch_df in tqdm(df.groupby("batch_id"), total=total_batch, desc="处理进度"): # 跳过已处理的批次 if batch_id in processed_batch: continue # 对目标列应用自定义函数 batch_df["request_result"] = batch_df[TARGET_COLUMN].apply(custom_get_request) # 增量写入CSV batch_df.to_csv( OUTPUT_FILE, mode="a", # 追加写入模式 header=(batch_id == 0 and not os.path.exists(OUTPUT_FILE)), # 仅首个批次写表头 index=False, encoding="utf-8-sig" )
3. 可选优化方向
- 单批次内可改用异步请求库
aiohttp替代同步的requests,并发处理请求可大幅提升整体处理速度 - 可给请求逻辑增加重试装饰器(比如
tenacity库的@retry),减少偶发网络波动导致的请求失败 - 处理前可先对目标列做去重处理,重复的请求目标无需多次发送请求,节省资源和时间
内容的提问来源于stack exchange,提问作者Theol
相关产品推荐
相关产品推荐

