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

如何为大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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 03:15:01