【问题标题】:How to change the size and distribution of a PySpark Dataframe according to the values of its rows & columns?如何根据行和列的值更改 PySpark Dataframe 的大小和分布?
【发布时间】:2021-01-02 20:13:18
【问题描述】:

我有一个大的 PySpark DataFrame,我想按照下面的示例进行操作。我认为形象化它比描述它更容易。因此,为了便于说明,让我们采用一个简单的 DataFrame df

df.show()
+----------+-----------+-----------+
|  series  | timestamp |   value   |
+----------+-----------+-----------+
|    ID1   |    t1     | value1_1  |
|    ID1   |    t2     | value2_1  |
|    ID1   |    t3     | value3_1  |
|    ID2   |    t1     | value1_2  |
|    ID2   |    t2     | value2_2  |
|    ID2   |    t3     | value3_2  |
|    ID3   |    t1     | value1_3  |
|    ID3   |    t2     | value2_3  |
|    ID3   |    t3     | value3_3  |
+----------+-----------+-----------+

在上述DataFrame中,series列中包含的三个唯一值中的每一个(即ID1ID2ID3)都有相应的值(在values列下)同时发生(即timestamp 列中的相同条目)。

从这个 DataFrame 中,我想进行一个转换,最终得到以下 DataFrame,例如,命名为 results。可以看出,DataFrame的大小发生了变化,甚至列都根据原始DataFrame的条目进行了重命名。

result.show()
+-----------+-----------+-----------+-----------+
| timestamp |    ID1    |    ID2    |    ID3    |
+-----------+-----------+-----------+-----------+
|    t1     |  value1_1 |  value1_2 |  value1_3 |
|    t2     |  value2_1 |  value2_2 |  value2_3 |
|    t3     |  value3_1 |  value3_2 |  value3_3 |
+-----------+-----------+-----------+-----------+

result 中的列顺序是任意的,不应影响最终答案。此说明性示例仅在 series 中包含三个唯一值(即 ID1ID2ID3)。理想情况下,我想编写一段代码,自动检测series 中的唯一值,从而生成一个新的对应列。有谁知道我可以从哪里开始?我尝试按timestamp 分组,然后使用聚合函数collect_set 收集一组不同的seriesvalue,但没有运气:(

提前非常感谢!

马里奥安萨斯

【问题讨论】:

    标签: python dataframe apache-spark pyspark apache-spark-sql


    【解决方案1】:

    只是一个简单的支点:

    import pyspark.sql.functions as F
    
    result = df.groupBy('timestamp').pivot('series').agg(F.first('value'))
    

    确保df 中的每一行都是不同的;否则重复的条目可能会被静默去重。

    【讨论】:

      【解决方案2】:

      扩展 mck 的回答,我找到了提高 pivot 性能的方法。 pivot 是一个非常昂贵的操作,因此,对于 Spark 2.0 及更高版本,建议提供列数据(如果已知)作为函数的参数,如下面的代码所示。这将提高 DataFrames 代码的性能,比这个问题中提出的说明性代码要大得多。鉴于 series 的值是事先已知的,我们可以使用:

      import pyspark.sql.functions as F
      
      series_list = ('ID1', 'ID2', 'ID3')
      result = df.groupBy('timestamp').pivot('series', series_list).agg(F.first('value'))
      result.show()
      +---------+--------+--------+--------+
      |timestamp|     ID1|     ID2|     ID3|
      +---------+--------+--------+--------+
      |       t1|value1_1|value1_2|value1_3|
      |       t2|value2_1|value2_2|value2_3|
      |       t3|value3_1|value3_2|value3_3|
      +---------+--------+--------+--------+
      

      【讨论】:

        猜你喜欢
        • 2018-07-03
        • 2022-07-04
        • 1970-01-01
        • 1970-01-01
        • 2020-05-28
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2017-11-23
        相关资源
        最近更新 更多