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

如何在Python中高效处理大量独立文件及LLM请求优化

大规模数据集处理优化方案

问题背景

需处理540,000条样本,分为100个独立目录,每个目录约含5400个CSV文件。当前采用Python单文件处理+bash多目录进程的方式,但整体耗时过长,现寻求文件I/O与处理流程的优化方案。

目录结构如下:

├── dir1
│   ├── file11.csv
│   ├── file12.csv
│   └── ...
├── dir2
│   ├── file21.csv
│   ├── file22.csv
│   └── ...
├── ...
└── dir99

当前处理代码:

def split_and_truncate(input_string, max_tokens=512):
    # Split the input string into tokens using space as a delimiter
    #print(F"Input String {input_string}")
    tokens = input_string.split()
    #print(F"Tokens {tokens}")

    # Check if the token count exceeds the specified limit
    if len(tokens) > max_tokens:
        # Truncate the string by joining the first max_tokens tokens with spaces
        truncated_string = ' '.join(tokens[:max_tokens])
    else:
        truncated_string = input_string

    return truncated_string

def execute_prompt(prompt):
    print(prompt)
    args = {"batch": [prompt]}
    request = requests.post(F"http://{node}:5000", json=args)
    results = request.json()
    print(results["data"][0])
    response = results["data"][0]
    return response

for dir in sorted_list:
    for filename in os.listdir(os.path.join(subdir, dir)):
        print(filename)
        file_name_vic = f'context_{mode_vic}_{filename}.json'
        context_file_vic = F'{subdir}/{dir}/{filename}/{file_name_vic}'

        summary_path = f"{subdir}/{dir}/{filename}/summary_{filename}.csv"
        
        if not os.path.exists(context_file_vic):
            with open(output_filename, "w") as output_file:
                output_file.write(f"{filename}\n")
            continue

        if os.path.exists(summary_path):
            continue
        df_vic = pd.read_json(context_file_vic, orient="records", dtype=object)

        df_vic = df_vic[['X', 'Y', 'text']]
        df_vic["summary"] = np.nan


        for index, row in df_vic.iterrows():
            print(F"length of input {len(row['text'].split())}")

            input_string = split_and_truncate(row['text'])


            prompt= f"""
            ### Instruction:
            Different instructions

            ### Input: 
            Reads the input:  "{input_string}"

            ### Output: """
            
            if pd.isna(df_vic.at[index, "summary"]):
                response2 = execute_prompt(prompt)
                df_vic.at[index, "summary"] = response2


            df_vic.to_csv(summary_path, index=False)

核心优化方案

1. LLM请求批量处理,削减网络开销

当前单条请求的往返延迟是核心瓶颈之一,改为批量提交请求:

  • 修改请求函数支持批量prompt提交,根据LLM服务器承载能力设置批量大小(比如20-50条)
  • 收集所有未处理样本的prompt,一次性发送请求,减少HTTP交互次数
  • 示例修改:
    def execute_batch_prompts(prompts):
        args = {"batch": prompts}
        # 添加超时避免进程挂起
        request = requests.post(f"http://{node}:5000", json=args, timeout=30)
        results = request.json()
        return results["data"]
    
    # 批量处理逻辑
    pending_prompts = []
    pending_indices = []
    batch_size = 20
    
    for index, row in df_vic.iterrows():
        if pd.notna(df_vic.at[index, "summary"]):
            continue
        input_string = split_and_truncate(row['text'])
        prompt = f"""
        ### Instruction:
        Different instructions
    
        ### Input: 
        Reads the input:  "{input_string}"
    
        ### Output: """
        pending_prompts.append(prompt)
        pending_indices.append(index)
    
        if len(pending_prompts) >= batch_size:
            responses = execute_batch_prompts(pending_prompts)
            for idx, resp in zip(pending_indices, responses):
                df_vic.at[idx, "summary"] = resp
            pending_prompts = []
            pending_indices = []
    
    # 处理剩余未提交的批量
    if pending_prompts:
        responses = execute_batch_prompts(pending_prompts)
        for idx, resp in zip(pending_indices, responses):
            df_vic.at[idx, "summary"] = resp
    

2. 优化文件I/O操作

消除循环内的频繁写入,减少磁盘开销:

  • 移除循环中的df_vic.to_csv,仅在整个文件处理完成后统一写入一次
  • 用pathlib替代字符串拼接路径,提升代码可读性与路径处理效率
  • 示例修改:
    from pathlib import Path
    
    # 路径处理
    context_file_vic = Path(subdir) / dir / filename / file_name_vic
    summary_path = Path(subdir) / dir / filename / f"summary_{filename}.csv"
    
    # 移除循环内的写入,处理完所有行后统一保存
    df_vic.to_csv(summary_path, index=False)
    

3. 提升DataFrame处理效率

避免低效的行迭代方式,改用向量化操作:

  • 替换iterrows()为apply或批量筛选逻辑,iterrows()是pandas中最慢的行遍历方式
  • 示例修改:
    def generate_prompt(text):
        input_string = split_and_truncate(text)
        return f"""
        ### Instruction:
        Different instructions
    
        ### Input: 
        Reads the input:  "{input_string}"
    
        ### Output: """
    
    # 筛选未处理的行并批量生成prompt
    mask = df_vic["summary"].isna()
    df_pending = df_vic[mask].copy()
    df_pending["prompt"] = df_pending["text"].apply(generate_prompt)
    
    # 批量请求并更新结果
    if len(df_pending) > 0:
        responses = execute_batch_prompts(df_pending["prompt"].tolist())
        df_vic.loc[mask, "summary"] = responses
    

4. 细化并行处理粒度

在目录级并行基础上,进一步实现文件级并行:

  • 用concurrent.futures.ProcessPoolExecutor实现多进程并行处理文件,进程数根据CPU核心与LLM服务器负载调整
  • 示例:
    from concurrent.futures import ProcessPoolExecutor
    
    def process_single_file(context_file_path):
        # 封装单个文件的完整处理逻辑
        dir_name = context_file_path.parent.parent.name
        filename = context_file_path.parent.name
        summary_path = context_file_path.parent / f"summary_{filename}.csv"
        
        if summary_path.exists():
            return
        
        df_vic = pd.read_json(context_file_path, orient="records", dtype=object)
        df_vic = df_vic[['X', 'Y', 'text']]
        df_vic["summary"] = np.nan
    
        # 批量处理逻辑(同上)
        # ...
    
        df_vic.to_csv(summary_path, index=False)
    
    # 收集所有待处理的文件路径
    context_files = []
    for dir in sorted_list:
        dir_path = Path(subdir) / dir
        for filename in os.listdir(dir_path):
            file_path = dir_path / filename / f'context_{mode_vic}_{filename}.json'
            if file_path.exists():
                summary_path = dir_path / filename / f"summary_{filename}.csv"
                if not summary_path.exists():
                    context_files.append(file_path)
    
    # 并行处理
    with ProcessPoolExecutor(max_workers=8) as executor:
        executor.map(process_single_file, context_files)
    

5. 细节优化

  • 移除所有不必要的print语句,减少控制台I/O开销
  • 优化split_and_truncate函数,避免全量拆分字符串:
    def split_and_truncate(input_string, max_tokens=512):
        if max_tokens <= 0:
            return ""
        count = 0
        truncate_pos = len(input_string)
        for i, char in enumerate(input_string):
            if char == ' ':
                count += 1
                if count == max_tokens:
                    truncate_pos = i
                    break
        return input_string[:truncate_pos].strip()
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 05:07:35