【问题标题】:PySpark - Split array in all columns and merge as rowsPySpark - 在所有列中拆分数组并合并为行
【发布时间】:2018-02-26 21:52:29
【问题描述】:

PySpark 中有没有办法同时在所有列中展开数组/列表并将展开的数据分别合并/压缩成行?

列数可能是动态的,具体取决于其他因素。

来自数据框

|col1   |col2   |col3   |
|[a,b,c]|[d,e,f]|[g,h,i]|
|[j,k,l]|[m,n,o]|[p,q,r]|

到数据框

|col1|col2|col3|
|a   |d   |g   |
|b   |e   |h   |
|c   |f   |i   |
|j   |m   |p   |
|k   |n   |q   |
|l   |o   |r   |

【问题讨论】:

  • 你能假设数组的长度吗?例如,它们总是相同的吗?和列数一样吗?
  • 数组长度和列数变化且是动态的

标签: apache-spark pyspark


【解决方案1】:

这是使用rddflatMap() 的一种方法:

cols = df.columns
df.rdd.flatMap(lambda x: zip(*[x[c] for c in cols])).toDF(cols).show(truncate=False)
#+----+----+----+
#|col1|col2|col3|
#+----+----+----+
#|a   |d   |g   |
#|b   |e   |h   |
#|c   |f   |i   |
#|j   |m   |p   |
#|k   |n   |q   |
#|l   |o   |r   |
#+----+----+----+

【讨论】:

    【解决方案2】:

    试试这个,

    import pyspark.sql.functions as F
    from pyspark.sql.types import *
    
    a = [(['a','b','c'],['d','e','f'],['g','h','i']),(['j','k','l'],['m','n','o'],['p','q','r'])]
    a = sql.createDataFrame(a,['a','b','c'])
    
    cols = ['col1','col2','col3']
    splits = [F.udf(lambda val:val[0],StringType()),F.udf(lambda val:val[1],StringType()),F.udf(lambda val:val[2],StringType())]
    
    def exploding(cols):
        return F.explode(F.array([F.struct([F.col(c).getItem(i).alias(c)\
                                      for c in colnames]) for i in range(3)]))
    
    a = a.withColumn("new_col", exploding(["a", "b", "c"]))\
                    .select([s('new_col').alias(c) for s,c in zip(splits,cols)])
    a.show()
    

    【讨论】:

    • 在玩具问题中,它按预期工作,但在没有时很难管理。列表中的项目和编号。的列是动态的。
    猜你喜欢
    • 2016-07-21
    • 1970-01-01
    • 1970-01-01
    • 2017-04-22
    • 2012-09-25
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多