Python multiprocessing正确使用及批量数据提取效率优化咨询
多进程提取用户记录方案优化咨询
我首次尝试使用multiprocessing提取约50万条用户记录,当前测试阶段配置变量为500条。原单循环执行耗时过长,因此改用多进程方案,目前启动10个Process可正常运行,但预估全量执行仍需约4小时。我计划提升至20个左右的进程运行,但担忧计算机出现性能问题导致程序崩溃,想咨询当前多进程用法是否正确,是否存在更优的实现方案?
初始版本代码
from pyETT import ett_parser import pandas as pd import time from datetime import datetime from multiprocessing import Process import sys c = 10 x1,y1 = 1,50 x2,y2 = 51,100 x3,y3 = 101,150 x4,y4 = 151,200 x5,y5 = 201,250 x6,y6 = 251,300 x7,y7 = 301,350 x8,y8 = 351,400 x9,y9 = 401,450 x10,y10 = 451,500 m_cols = ('user-name','elo','rank','wins','losses','last-online') def run1(): print('Running query 1...') df_master = pd.DataFrame(columns = m_cols) for i in range(x1, y1): try: if int(i) % int(c) == 0: print('Loop1 is at:', i) user_id = i line = ett.ett_parser.get_user(user_id) temp_df = pd.DataFrame(line, index=[i]) df_master = df_master.append(temp_df, ignore_index = True) except Exception: print("Error_1:",i ) #Export to excel file_name = 'export_file1_' + datetime.now().strftime("%H_%M_%S") + '.xlsx' df_master.to_excel(file_name, index = False) print('DataFrame(' + file_name + ') is written to Excel File successfully.') def run2(): print('Running query2...') m_cols = ('user-name','elo','rank','wins','losses','last-online') df_master = pd.DataFrame(columns = m_cols) for i in range(x2, y2): try: if int(i) % int(c) == 0: print('Loop2 is at:', i) user_id = i line = ett.ett_parser.get_user(user_id) temp_df = pd.DataFrame(line, index=[i]) df_master = df_master.append(temp_df, ignore_index = True) except Exception: print("Error_2:",i ) #Export to excel file_name = 'export_file2_' + datetime.now().strftime("%H_%M_%S") + '.xlsx' df_master.to_excel(file_name, index = False) print('DataFrame(' + file_name + ') is written to Excel File successfully.') def run3(): print('Running query3...') df_master = pd.DataFrame(columns = m_cols) for i in range(x3, y3): try: if int(i) % int(c) == 0: print('Loop3 is at:', i) user_id = i line = ett.ett_parser.get_user(user_id) temp_df = pd.DataFrame(line, index=[i]) df_master = df_master.append(temp_df, ignore_index = True) except Exception: print("Error_3:",i ) #Export to excel file_name = 'export_file3_' + datetime.now().strftime("%H_%M_%S") + '.xlsx' df_master.to_excel(file_name, index = False) print('DataFrame(' + file_name + ') is written to Excel File successfully.') def run4(): print('Running query4...') df_master = pd.DataFrame(columns = m_cols) for i in range(x4, y4): try: if int(i) % int(c) == 0: print('Loop4 is at:', i) user_id = i line = ett.ett_parser.get_user(user_id) temp_df = pd.DataFrame(line, index=[i]) df_master = df_master.append(temp_df, ignore_index = True) except Exception: print("Error_4:",i ) #Export to excel file_name = 'export_file4_' + datetime.now().strftime("%H_%M_%S") + '.xlsx' df_master.to_excel(file_name, index = False) print('DataFrame(' + file_name + ') is written to Excel File successfully.') def run5(): print('Running query5...') df_master = pd.DataFrame(columns = m_cols) for i in range(x5, y5): try: if int(i) % int(c) == 0: print('Loop5 is at:', i) user_id = i line = ett.ett_parser.get_user(user_id) temp_df = pd.DataFrame(line, index=[i]) df_master = df_master.append(temp_df, ignore_index = True) except Exception: print("Error_5:",i ) #Export to excel file_name = 'export_file5_' + datetime.now().strftime("%H_%M_%S") + '.xlsx' df_master.to_excel(file_name, index = False) print('DataFrame(' + file_name + ') is written to Excel File successfully.') def run6(): print('Running query6...') df_master = pd.DataFrame(columns = m_cols) for i in range(x6, y6): try: if int(i) % int(c) == 0: print('Loop6 is at:', i) user_id = i line = ett.ett_parser.get_user(user_id) temp_df = pd.DataFrame(line, index=[i]) df_master = df_master.append(temp_df, ignore_index = True) except Exception: print("Error_6:",i ) #Export to excel file_name = 'export_file6_' + datetime.now().strftime("%H_%M_%S") + '.xlsx' df_master.to_excel(file_name, index = False) print('DataFrame(' + file_name + ') is written to Excel File successfully.') def run7(): print('Running query7...') df_master = pd.DataFrame(columns = m_cols) for i in range(x7, y7): try: if int(i) % int(c) == 0: print('Loop7 is at:', i) user_id = i line = ett.ett_parser.get_user(user_id) temp_df = pd.DataFrame(line, index=[i]) df_master = df_master.append(temp_df, ignore_index = True) except Exception: print("Error_7:",i ) #Export to excel file_name = 'export_file7_' + datetime.now().strftime("%H_%M_%S") + '.xlsx' df_master.to_excel(file_name, index = False) print('DataFrame(' + file_name + ') is written to Excel File successfully.') def run8(): print('Running query8...') df_master = pd.DataFrame(columns = m_cols) for i in range(x8, y8): try: if int(i) % int(c) == 0: print('Loop8 is at:', i) user_id = i line = ett.ett_parser.get_user(user_id) temp_df = pd.DataFrame(line, index=[i]) df_master = df_master.append(temp_df, ignore_index = True) except Exception: print("Error_8:",i ) #Export to excel file_name = 'export_file8_' + datetime.now().strftime("%H_%M_%S") + '.xlsx' df_master.to_excel(file_name, index = False) print('DataFrame(' + file_name + ') is written to Excel File successfully.') def run9(): print('Running query9...') df_master = pd.DataFrame(columns = m_cols) for i in range(x9, y9): try: if int(i) % int(c) == 0: print('Loop9 is at:', i) user_id = i line = ett.ett_parser.get_user(user_id) temp_df = pd.DataFrame(line, index=[i]) df_master = df_master.append(temp_df, ignore_index = True) except Exception: print("Error_9:",i ) #Export to excel file_name = 'export_file9_' + datetime.now().strftime("%H_%M_%S") + '.xlsx' df_master.to_excel(file_name, index = False) print('DataFrame(' + file_name + ') is written to Excel File successfully.') def run10(): print('Running query10...') df_master = pd.DataFrame(columns = m_cols) for i in range(x10, y10): try: if int(i) % int(c) == 0: print('Loop10 is at:', i) user_id = i line = ett.ett_parser.get_user(user_id) temp_df = pd.DataFrame(line, index=[i]) df_master = df_master.append(temp_df, ignore_index = True) except Exception: print("Error_10:",i ) #Export to excel file_name = 'export_file10_' + datetime.now().strftime("%H_%M_%S") + '.xlsx' df_master.to_excel(file_name, index = False) print('DataFrame(' + file_name + ') is written to Excel File successfully.') def main(): p = Process(target=run1) p.start() #p.join() p2 = Process(target=run2) p2.start() p3 = Process(target=run3) p3.start() p4 = Process(target=run4) p4.start() p5 = Process(target=run5) p5.start() p6 = Process(target=run6) p6.start() p7 = Process(target=run7) p7.start() p8 = Process(target=run8) p8.start() p9 = Process(target=run9) p9.start() p10 = Process(target=run10) p10.start() p10.join() if __name__ == '__main__': start = time.time() print('starting main') main() print('finishing main',time.time()-start)
优化后代码
参考社区回答实现的代码,可满足需求且代码大幅精简:
from concurrent.futures import ThreadPoolExecutor from multiprocessing import cpu_count from pyETT import ett_parser import pandas as pd import time def main(): USER_ID_COUNT = 50 MAX_WORKERS = 2 * cpu_count() + 1 dataframe_list = [] #user_array = [] user_ids = list(range(1, USER_ID_COUNT)) def obtain_user_record(user_id): return ett_parser.get_user(user_id) with ThreadPoolExecutor(MAX_WORKERS) as executor: for user_id, user_record in zip(user_ids, executor.map(obtain_user_record, user_ids)): if user_record: dataframe_list.append(user_record) df_master = pd.DataFrame.from_dict(dataframe_list,orient='columns') print(df_master) if __name__ == '__main__': start = time.time() print('starting main') main() print('finishing main', time.time() - start)
内容的提问来源于stack exchange,提问作者cmccall95
相关产品推荐
相关产品推荐

