使用ProcessPoolExecutor多进程生成Pandas DataFrame时遇提取错误
问题:多进程创建Pandas DataFrame无法提取,报错TypeError
我尝试用多进程同时创建两个Pandas DataFrame,两个函数的返回结果均为DataFrame,但使用concurrent.futures.ProcessPoolExecutor时无法提取这些DataFrame,代码如下:
from datetime import date import concurrent.futures import sgs import numpy as np import pandas as pd data_inicial_cdi = '02/01/1998' data_inicial_selic = '04/01/2021' data_final = date.today().strftime('%d/%m/%Y') def taxa_CDI(*knargs): dt1 = data_inicial_cdi dt2 = data_final taxa_CDI = sgs.dataframe([4389], start=data_inicial_cdi, end=data_final) taxa_CDI.rename(columns = {4389: 'CDI_taxa'}, inplace=True) taxa_CDI.dropna(inplace=True) taxa_CDI['CDI_fator'] = (1+taxa_CDI['CDI_taxa']/100) ** (1/252) taxa_CDI['CDI_acum'] = np.cumprod(taxa_CDI['CDI_fator'] ) taxa_CDI['Data'] = taxa_CDI.index total_registro_cdi = len(taxa_CDI) index_cdi = np.array(range(total_registro_cdi)) taxa_CDI.set_index(index_cdi, inplace=True) return taxa_CDI def taxa_Selic(*knargs): dt1 = data_inicial_selic dt2 = data_final taxa_Selic = sgs.dataframe([1178], start=data_inicial_selic, end=data_final) taxa_Selic.rename(columns={1178: 'Selic_taxa'}, inplace=True) taxa_Selic.dropna(inplace=True) taxa_Selic['Selic_fator'] = (1 + taxa_Selic['Selic_taxa'] / 100) ** (1 / 252) taxa_Selic['Selic_acum'] = np.cumprod(taxa_Selic['Selic_fator']) taxa_Selic['Data'] = taxa_Selic.index total_registro_selic = len(taxa_Selic) index_selic = np.array(range(total_registro_selic)) taxa_Selic.set_index(index_selic, inplace=True) return taxa_Selic if __name__ == '__main__': with concurrent.futures.ProcessPoolExecutor(max_workers=2) as executor: cod1 = executor.submit(taxa_CDI(data_inicial_cdi, data_final)) cod2 = executor.submit(taxa_Selic(data_inicial_selic, data_final)) df_cdi = cod1.result() df_selic = cod2.result() print(cod1, cod2)
报错提示:
TypeError: 'DataFrame' object is not callable
解决方案
核心问题
你调用executor.submit()时犯了关键错误:直接执行了函数(比如taxa_CDI(data_inicial_cdi, data_final)),这会让函数立刻在主进程中运行并返回DataFrame,随后你把这个DataFrame对象传给了submit(),但submit()需要的是函数本身和它的参数,不是函数执行后的结果。
修正步骤
修改
submit()的调用方式:
把错误的调用:cod1 = executor.submit(taxa_CDI(data_inicial_cdi, data_final)) cod2 = executor.submit(taxa_Selic(data_inicial_selic, data_final))改成:
cod1 = executor.submit(taxa_CDI, data_inicial_cdi, data_final) cod2 = executor.submit(taxa_Selic, data_inicial_selic, data_final)这样
submit()会把函数和参数传递给子进程执行,返回Future对象,调用result()就能获取对应的DataFrame。优化函数参数(可选但推荐):
你的两个函数定义了*knargs但并未使用,而是依赖全局变量,建议改成明确接收参数的形式,避免全局变量依赖,让代码更健壮:def taxa_CDI(dt1, dt2): taxa_CDI = sgs.dataframe([4389], start=dt1, end=dt2) taxa_CDI.rename(columns = {4389: 'CDI_taxa'}, inplace=True) taxa_CDI.dropna(inplace=True) taxa_CDI['CDI_fator'] = (1+taxa_CDI['CDI_taxa']/100) ** (1/252) taxa_CDI['CDI_acum'] = np.cumprod(taxa_CDI['CDI_fator'] ) taxa_CDI['Data'] = taxa_CDI.index total_registro_cdi = len(taxa_CDI) index_cdi = np.array(range(total_registro_cdi)) taxa_CDI.set_index(index_cdi, inplace=True) return taxa_CDI def taxa_Selic(dt1, dt2): taxa_Selic = sgs.dataframe([1178], start=dt1, end=dt2) taxa_Selic.rename(columns={1178: 'Selic_taxa'}, inplace=True) taxa_Selic.dropna(inplace=True) taxa_Selic['Selic_fator'] = (1 + taxa_Selic['Selic_taxa'] / 100) ** (1 / 252) taxa_Selic['Selic_acum'] = np.cumprod(taxa_Selic['Selic_fator']) taxa_Selic['Data'] = taxa_Selic.index total_registro_selic = len(taxa_Selic) index_selic = np.array(range(total_registro_selic)) taxa_Selic.set_index(index_selic, inplace=True) return taxa_Selic
修正后效果
修改完成后,df_cdi和df_selic就能正常获取两个函数生成的DataFrame,多进程也能正常工作。
内容的提问来源于stack exchange,提问作者Carlos
相关产品推荐
相关产品推荐

