【问题标题】:Convert a RDD of Tuples of Varying Sizes to a DataFrame in Spark在 Spark 中将不同大小的元组的 RDD 转换为 DataFrame
【发布时间】:2016-05-31 17:25:33
【问题描述】:

我在使用 python 将以下结构的 RDD 转换为 spark 中的数据帧时遇到了困难。

df1=[['usr1',('itm1',2),('itm3',3)], ['usr2',('itm2',3), ('itm3',5),(itm22,6)]]

转换后,我的数据框应如下所示:

       usr1  usr2
itm1    2.0   NaN
itm2    NaN   3.0
itm22   NaN   6.0
itm3    3.0   5.0

我最初是想将上述 RDD 结构转换为以下内容:

df1={'usr1': {'itm1': 2, 'itm3': 3}, 'usr2': {'itm2': 3, 'itm3': 5, 'itm22':6}}

然后使用python的pandas模块pand=pd.DataFrame(dat2),然后使用spark_df = context.createDataFrame(pand)将pandas数据帧转换回spark数据帧。但是,我相信,通过这样做,我将 RDD 转换为非 RDD 对象,然后再转换回 RDD,这是不正确的。有人可以帮我解决这个问题吗?

【问题讨论】:

  • 这有什么不同,不包括列选择,from your previous question
  • 请注意,在我之前的问题中,我更关心为同一用户处理重复的“itms”(请参阅​​“如果还有一个类似上述元组的计数字段,即 ('itm1' ,3)如何将这个值 3 合并(或添加)到列联表(或实体项矩阵)的最终结果中”。由于给出的答案仍然不清楚(至少从我的角度来看),如果我是能够得到这个问题的解决方案,我可以关闭上一个问题的答案。

标签: python apache-spark tuples pyspark


【解决方案1】:

这样的数据:

rdd = sc.parallelize([
    ['usr1',('itm1',2),('itm3',3)], ['usr2',('itm2',3), ('itm3',5),('itm22',6)]
])

展平记录:

def to_record(kvs):
    user, *vs = kvs  # For Python 2.x use standard indexing / splicing
    for item, value in vs:
        yield user, item, value

records = rdd.flatMap(to_record)

转换为DataFrame:

df = records.toDF(["user", "item", "value"])

枢轴:

result = df.groupBy("item").pivot("user").sum()

result.show()
## +-----+----+----+
## | item|usr1|usr2|
## +-----+----+----+
## | itm1|   2|null|
## | itm2|null|   3|
## | itm3|   3|   5|
## |itm22|null|   6|
## +-----+----+----+

注意:Spark DataFrames 旨在处理较长且相对较细的数据。如果您想生成宽列联表,DataFrames 将没有用处,尤其是在数据密集且您希望为每个特征保留单独列的情况下。

【讨论】:

  • 非常感谢!也感谢您提供的额外信息。
猜你喜欢
  • 1970-01-01
  • 2017-06-02
  • 1970-01-01
  • 2017-11-02
  • 2019-05-14
  • 2018-03-05
  • 1970-01-01
  • 2017-06-13
  • 1970-01-01
相关资源
最近更新 更多