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

如何加速Pandas合并.bin.txt文件并导出的过程?

加速大规模.bin.txt文件合并的优化方案

我手上有大约1500个.bin.txt文件,用Excel PowerQuery加载分析时速度极慢(耗时15分钟以上),因此编写了Python脚本先合并所有文件再导入PowerQuery。但当前脚本合并这些文件耗时超20分钟,而全量数据规模将是当前的10倍甚至更多,现寻求加速该合并过程的方法。

初始脚本

def combineBinFiles():
    root = tk.Tk()
    root.withdraw()
    # 获取bin文件所在文件夹
    folder_selected = filedialog.askdirectory()
    print(folder_selected)
    os.chdir(folder_selected)
    files = os.listdir(folder_selected)
    print(files)
    df = pd.DataFrame()
    temp_df = pd.DataFrame(columns= ['Timestamp', 'Wind speed', 'Own consumption'])
    for file in files:
        if file.endswith('.bin.txt'):
            print("正在读取文件: " + file)
            # 按空格分隔读取文件
            temp_df = pd.read_csv(file, delimiter=" ", header=0)
            # 从列名提取日期
            date = temp_df.columns[0] 
            # 给每个时间戳拼接日期,再加入数据框
            temp_df[date] = temp_df[date].apply(lambda x: date + ' ' + x)
            df = pd.concat([df, temp_df], axis=0, ignore_index=True)
            
    df.to_csv('combinedBinFile.csv', index=False)
    print(df)

combineBinFiles()

文件格式示例

14_07_2023  .WindSpeed  .Power  
17 50 00 006    10,53   0   
17 50 00 016    10,53   0   
17 50 00 026    10,53   0   
17 50 00 036    10,53   0   
17 50 00 046    10,53   0   
17 50 00 056    10,53   0   

尝试的多线程分块脚本

def worker(q, df_list):
    while not q.empty():
        file = q.get()
        if file.endswith('.bin.txt'):
            print("正在读取文件: " + file)
            temp_df = pd.read_csv(file, delimiter="\t", header=0, engine='python')
            date = temp_df.columns[0] 
            temp_df[date] = temp_df[date].apply(lambda x: date + ' ' + x)
            # 删除最后一列
            temp_df = temp_df.iloc[:, :-1]
            temp_df.columns = ['Timestamp', 'Wind speed', 'Own consumption']
            df_list.append(temp_df)
        q.task_done()

def combineBinFilesThreaded():
    root = tk.Tk()
    root.withdraw()
    folder_selected = filedialog.askdirectory()
    print(folder_selected)

    os.chdir(folder_selected)
    files = os.listdir(folder_selected)
    print(files)

    df_list = []
    q = queue.Queue()

    # 创建4个工作线程
    for i in range(4):
        t = threading.Thread(target=worker, args=(q, df_list))
        t.daemon = True
        t.start()

    # 将文件加入队列
    for file in files:
        q.put(file)

    # 等待队列中所有任务完成
    q.join()
    print("所有文件读取完成")
    # 合并所有数据框
    df = pd.concat(df_list, axis=0, ignore_index=True)
    print("数据框合并完成")
    chunksize = 100000
    # 检查目标文件是否存在,避免重复追加
    if os.path.exists('combinedBinFile.csv'):
        os.remove('combinedBinFile.csv')
        
    for i in range(0, len(df), chunksize):
        print("正在写入块: " + str(i))
        df.iloc[i:i+chunksize].to_csv('combinedBinFile.csv', index=False, mode='a')

    print("写入完成")

核心优化方案

1. 用多进程替代多线程

Python的全局解释器锁(GIL)会限制多线程在CPU密集型任务中的效率,改用multiprocessing可以充分利用多核CPU,大幅提升文件处理速度。

2. 优化pd.read_csv参数

  • 使用sep='\s+'适配文件中不规则的空格分隔,避免因分隔符问题读取错误
  • 指定decimal=','让数值列直接识别为浮点数,省去后续格式转换开销
  • 用usecols指定仅读取需要的列,减少内存占用
  • 保留默认engine='c',比python引擎速度快数倍

3. 矢量化操作替代apply

将逐行循环的apply替换为矢量化字符串拼接df[date] = date + ' ' + df[date],避免逐行处理的低效。

4. 边处理边写入,避免内存过载

无需在内存中合并所有DataFrame,每个文件处理完成后直接追加到输出CSV,降低内存占用,避免因内存不足触发磁盘交换拖慢速度。

优化后的多进程脚本

import os
import tkinter as tk
from tkinter import filedialog
import multiprocessing as mp
import pandas as pd

def process_file(file_path):
    if not file_path.endswith('.bin.txt'):
        return None
    print(f"正在处理文件: {os.path.basename(file_path)}")
    # 优化read_csv参数,仅读取必要列、适配分隔符和小数格式
    df = pd.read_csv(
        file_path,
        sep='\s+',
        header=0,
        decimal=',',
        usecols=[0, 1, 2]
    )
    date = df.columns[0]
    # 矢量化拼接日期与时间戳
    df[date] = date + ' ' + df[date]
    df.columns = ['Timestamp', 'Wind speed', 'Own consumption']
    return df

def save_chunk(output_path, df_chunk, is_first):
    # 首次写入带表头,后续追加不带表头
    df_chunk.to_csv(
        output_path,
        index=False,
        mode='w' if is_first else 'a',
        header=is_first
    )

def combineBinFilesMultiprocess():
    root = tk.Tk()
    root.withdraw()
    folder_selected = filedialog.askdirectory()
    print(folder_selected)

    output_path = os.path.join(folder_selected, 'combinedBinFile.csv')
    # 清理已存在的输出文件
    if os.path.exists(output_path):
        os.remove(output_path)

    # 生成所有文件的完整路径
    file_paths = [os.path.join(folder_selected, f) for f in os.listdir(folder_selected)]
    first_write = True

    # 用CPU核心数创建进程池,最大化利用硬件资源
    with mp.Pool(mp.cpu_count()) as pool:
        # 无序映射处理文件,提升效率
        for result_df in pool.imap_unordered(process_file, file_paths):
            if result_df is not None:
                save_chunk(output_path, result_df, first_write)
                first_write = False

    print("所有文件合并完成")

if __name__ == '__main__':
    combineBinFilesMultiprocess()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 02:12:06