【发布时间】: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),需要指定至少一列进行排序,这样行的顺序才能确定。