【问题标题】:Paralle apply function on df in pythonpython中df上的并行应用函数
【发布时间】:2020-09-19 14:12:21
【问题描述】:

我有一个函数可以遍历两个列表:项目和日期。 该函数返回更新的项目列表。 目前它使用 apply 运行,这在数百万行上效率不高。 我想通过并行化来提高效率。

项目列表中的项目按时间顺序排列,以及对应的日期列表(item_list 和 date_list 大小相同)。

这是df:

Date        item_list            date_list

12/05/20    [I1,I3,I4]    [10/05/20, 11/05/20, 12/05/20 ]
11/05/20    [I1,I3]       [11/05/20 , 14/05/20]

这就是我想要的df:

Date        item_list     date_list             items_list_per_date  

12/05/20    [I1,I3,I4]    [10/05/20, 11/05/20, 12/05/20]   [I1,I3]
11/05/20    [I1,I3]       [11/05/20 , 14/05/20]               nan

这是我的代码:

def get_item_list_per_date(date, items_list, date_list):

    if str(items_list)=="nan" or str(date_list)=="nan":
        return np.nan

    new_date_list = []
    for d in list(date_list):
        new_date_list.append(pd.to_datetime(d))

    if (date in new_date_list) and (len(new_date_list)>1):
        loc = new_date_list.index(date)
    else:
        return np.nan

    updated_items_list = items_list[:loc]

    if len(updated_items_list )==0:
        return np.nan

    return updated_items_list 

df['items_list_per_date'] = df.progress_apply(lambda x: get_item_list_per_date(date=x['date'], items_list=x['items_list'], date_list=x['date_list']),axis=1)

我很想把它并行化,你能帮忙吗?

【问题讨论】:

  • 你能分享train数据的样本吗?
  • 是的,对不起。我刚刚更新了它
  • Date 列和date_list 列中日期的值是字符串,对吧?
  • 不,它是一个实际日期(带小时、分钟......)。 new_date_list.append(pd.to_datetime(d)) - 因为当我创建日期列表时解析发生了变化。
  • 你能分享df['Date'].dtypetype(df['date_list'].iloc[0][0])的输出吗

标签: python parallel-processing processing-efficiency


【解决方案1】:

用途:

import multiprocessing as mp

def fx(df):
    def __fx(s):
        date = s['Date']
        date_list = s['date_list']
        if date in date_list:
            loc = date_list.index(date)
            return s['item_list'][:loc]
        else:
            return np.nan

    return df.apply(__fx, axis=1)

def parallel_apply(df):
    dfs = filter(lambda d: not d.empty, np.array_split(df, mp.cpu_count()))
    pool = mp.Pool()
    per_date = pd.concat(pool.map(fx, dfs))
    pool.close()
    pool.join()
    return per_date

df['items_list_per_date'] = parallel_apply(df)

结果:

#print(df)

Date        item_list     date_list             items_list_per_date  

12/05/20    [I1,I3,I4]    [10/05/20, 11/05/20, 12/05/20]   [I1,I3]
11/05/20    [I1,I3]       [11/05/20 , 14/05/20]               nan

【讨论】:

  • 在对照常规应用功能检查后,我发现它需要更多时间:/。我正在使用google colab,也许有问题?只有两个 cpu 但它应该比常规应用更快。
猜你喜欢
  • 2018-05-01
  • 1970-01-01
  • 1970-01-01
  • 2020-12-26
  • 2022-01-04
  • 1970-01-01
  • 1970-01-01
  • 2018-04-14
  • 1970-01-01
相关资源
最近更新 更多