【问题标题】:Pandas iterations to Pyspark FunctionPandas 对 Pyspark 函数的迭代
【发布时间】:2020-04-21 12:31:26
【问题描述】:

给定一个数据集,我尝试创建一个逻辑,在该逻辑中,我需要在两列中强制执行连续性,即最后一个目的地(在列中)是每个 id 的确切下一个起点(从列中)。比如这张表

+----+-------+-------+
| id | from  | to    |
+----+-------+-------+
|  1 | A     | B     |
|  1 | C     | A     |
|  2 | D     | D     |
|  2 | F     | G     |
|  2 | F     | F     |
+----+-------+-------+

理想情况下应该是这样的:

+----+-------+-------+
| id | from  | to    |
+----+-------+-------+
|  1 | A     | B     |
|  1 | B     | C     |
|  1 | C     | A     |
|  2 | D     | D     |
|  2 | D     | F     |
|  2 | F     | G     |
|  2 | G     | F     |
|  2 | F     | F     |
+----+-------+-------+

使用 Pandas,我通过逐行循环并检查 previous_row['to'] == current_row['from'] 来做到这一点,这也是一个使用 groupby 可以避免的 id 检查,如下所示

for i in range(len(df)):
    if (i < (len(df)-1)):
        if (new.ix[i,"to"] != new.ix[i+1,"from"]) & (new.ix[i,"id"] == new.ix[i+1,"id"]): 
            new_index = i + 0.5
            line = pd.DataFrame({"id":new.ix[i,"id"],
                             "from":new.ix[i,"to"],"to":new.ix[i+1,"from"],}, index = [new_index])
            appendings = pd.concat([appendings,line])
        else:
            pass
    else:
        pass

是否可以将其“翻译”为 pyspark rdds?

我知道在 Pyspark 中循环远非最佳,无法复制循环和 if-else 逻辑。

我考虑过对列进行分组和压缩,然后处理单个列。这样做的主要问题在于,我可以在“错误”的行上生成一个标志,但是如果不使用索引操作,就无法插入新行。

【问题讨论】:

  • 在spark中,shuffle后行的顺序是不确定的(即groupby,window),需要指定至少一列进行排序,这样行的顺序才能确定。

标签: python pandas pyspark


【解决方案1】:

这不是pyspark 的答案,而是向您展示如何在pandas 中实现无循环任务的部分答案。

你可以试试:

def f(sub_df):
    return sub_df.assign(to_=np.roll(sub_df.To, 1)) \
                .apply(lambda x: [[x.From, x.To]] if x.to_ == x.From else [[x.to_, x.From], [x.From, x.To]], axis=1) \
                .explode() \
                .apply(pd.Series)


out = df.groupby('id').apply(f) \
        .reset_index(level=1, drop=True) \
        .rename(columns={0: "from", 1: "to"})

工作流程:

  • 通过id 使用groupby 对数据帧进行分组
  • 对于每个组:
    • 创建一个新列(此处名称为to_)以保留上一行。 np.roll 执行循环移位以保留最后一个值。
    • 根据如果current from == previous to:返回当前行或添加新行以进行转换。
    • 使用explodelistlist 分解为每行一个列表。
    • 使用apply(pd.Series)将该列转换为两列
  • 然后,对于输出数据帧,使用 reset_index 删除 level 1 索引
  • 并使用rename 重命名列

完整代码

# Import module
import pandas as pd
import numpy as np

# create dataset
df = pd.DataFrame({"id": [1,1,2,2,2], "From": ["A", "C", "D", "F", "F"], "To": ["B", "D", "D", "G", "F"]})
# print(df)


def f(sub_df):
    return sub_df.assign(to_=np.roll(sub_df.To, 1)) \
                .apply(lambda x: [[x.From, x.To]] if x.to_ == x.From else [[x.to_, x.From], [x.From, x.To]], axis=1) \
                .explode() \
                .apply(pd.Series)


out = df.groupby('id').apply(f) \
        .reset_index(level=1, drop=True) \
        .rename(columns={0: "from", 1: "to"})
print(out)
#    from to
# id
# 1     D  A
# 1     A  B
# 1     B  C
# 1     C  D
# 2     F  D
# 2     D  D
# 2     D  F
# 2     F  G
# 2     G  F
# 2     F  F

下一步是将其翻译成PySpark。尝试一下,并随时根据您的尝试提出一个新问题。

【讨论】:

    猜你喜欢
    • 2018-10-28
    • 1970-01-01
    • 2016-08-11
    • 1970-01-01
    • 2021-02-07
    • 1970-01-01
    • 2022-06-29
    • 1970-01-01
    • 2018-07-05
    相关资源
    最近更新 更多