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

Python多线程无提速问题求助:核心代码段并行增益失效

问题描述

代码功能正常,但启用5个线程后未实现预期提速,耗时与单线程一致。性能瓶颈位于代码中TRY 10至TRY 20的区间,初步怀疑是temp_df_time = series.to_frame().T.rename(columns=column_mapping)这一DataFrame转换环节导致,而非IO操作问题。附上完整代码:

final_dfs_time = []

# Initialize a lock for synchronizing access to df_list_time
lock = threading.Lock()
json_list_lock = threading.Lock()

def append_to_dataframe(file_path):
    df_temp_time = None
    if os.path.getsize(file_path) > 0:
        try:
            print(f"try 1 {file_path}")
            # Read the JSON file using buffered I/O
            with open(file_path, 'rb') as file:
                buffered_file = io.BufferedReader(file)
                json_list = json.load(buffered_file)
                print(f"try 2 {file_path}")

            # Initialize an empty list to store temporary DataFrames
            df_list_time = []

            print(f"try 10")


            for entry in json_list:  # Iterate over json_list, not json_data
                try:
                    with json_list_lock:
                        for epoch_time, values in entry.items():
                            # Create a Series for each key-value pair
                            
                            series = pd.Series(values, name=epoch_time)
                            
                            # Convert the series into a DataFrame
                            temp_df_time = series.to_frame().T.rename(columns=column_mapping)


                            # Add 'time' column
                            # Convert epoch_time to a pandas Series with the same length as temp_df_time
                            epoch_time_series = pd.Series(epoch_time, index=temp_df_time.index)
                            # Assign the Series to the 'time' column directly
                            
                            temp_df_time['time'] = epoch_time_series
                            
                            # Acquire the lock before appending to df_list_time
                            with lock:
                                # print(f"LOCK {file_path}")
                                # Append the DataFrame to the list
                                
                                df_list_time.append(temp_df_time)
                                

                except ValueError as e:
                  print(f"Error reading {file_path}: {e}")



            print(f"TRY 20")


            # Concatenate the DataFrames into one DataFrame
            df_temp_time = pd.concat((df for df in df_list_time), ignore_index=True) #stuck here for 1
            print(f"try4 {file_path}")
            print(f"Appended {file_path} for user {user_id_dir}")

        except ValueError as e:
            print(f"Error reading {file_path}: {e}")

    else:
        print(f"Skipping empty file: {file_path}")

    print(f"try5 {file_path}")
    return df_temp_time


num_threads = 5
# Create a ThreadPoolExecutor with the desired number of threads
with ThreadPoolExecutor(max_workers=num_threads) as executor:
    # Map the function to each file path and process concurrently
    results = executor.map(append_to_dataframe, file_paths)

# Concatenate all resulting DataFrames
final_dfs_time = list(results)

total_df_time = pd.concat(final_dfs_time, ignore_index=True)
print("all added together")

瓶颈原因分析

  1. 全局锁导致线程串行化:代码中使用的json_list_lock和lock是全局锁,且核心循环全程持有json_list_lock,每次添加临时DataFrame还要获取lock。这直接让所有线程的核心处理逻辑变成串行执行,完全丧失多线程并行优势,耗时自然和单线程一致。
  2. 小DataFrame频繁创建+拼接的低效:循环中反复创建单一行的DataFrame,再逐个加入列表后拼接,这种小对象的频繁创建会带来极大的性能开销,是Pandas处理数据的典型低效场景。
  3. 冗余的Series转换:为添加time列创建epoch_time_series属于多余操作,直接赋值标量即可实现相同效果,无需额外生成Series。

优化建议

1. 移除不必要的全局锁

每个线程处理独立文件,df_list_time是线程内部局部变量,json_list是每个线程从独立文件读取的局部数据,完全不需要锁保护。直接删除两个锁的定义和使用,让线程真正并行执行。

2. 批量构建数据,避免逐行创建DataFrame

改用列表收集字典格式的原始数据,最后一次性转换为DataFrame,大幅降低Pandas对象创建开销:

def append_to_dataframe(file_path):
    df_temp_time = None
    if os.path.getsize(file_path) > 0:
        try:
            print(f"try 1 {file_path}")
            with open(file_path, 'rb') as file:
                buffered_file = io.BufferedReader(file)
                json_list = json.load(buffered_file)
            print(f"try 2 {file_path}")

            # 用列表收集所有行数据,避免逐行创建小DataFrame
            data_rows = []

            print(f"try 10")
            for entry in json_list:
                try:
                    for epoch_time, values in entry.items():
                        # 直接映射列名并添加time字段
                        mapped_row = {column_mapping[k]: v for k, v in values.items()}
                        mapped_row['time'] = epoch_time
                        data_rows.append(mapped_row)
                except ValueError as e:
                    print(f"Error reading {file_path}: {e}")

            print(f"TRY 20")
            # 一次性转换为DataFrame
            df_temp_time = pd.DataFrame(data_rows) if data_rows else None
            print(f"try4 {file_path}")
            print(f"Appended {file_path} for user {user_id_dir}")

        except ValueError as e:
            print(f"Error reading {file_path}: {e}")
    else:
        print(f"Skipping empty file: {file_path}")
    print(f"try5 {file_path}")
    return df_temp_time

3. 改用ProcessPoolExecutor提升CPU密集场景性能

Pandas底层操作多受GIL限制,多线程在CPU密集场景下提速有限。如果核心处理是CPU密集型,建议替换为多进程绕过GIL:

from concurrent.futures import ProcessPoolExecutor

# 替换ThreadPoolExecutor为ProcessPoolExecutor
with ProcessPoolExecutor(max_workers=num_threads) as executor:
    results = executor.map(append_to_dataframe, file_paths)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 19:30:24