【问题标题】:Pyskark Dataframe: Transforming unique elements in rows to columnsPyspark Dataframe:将行中的唯一元素转换为列
【发布时间】:2017-11-08 23:37:28
【问题描述】:

我有以下格式的 Pyspark 数据框:

+------------+---------+
|    date    |  query  |
+------------+---------+
| 2011-08-11 | Query 1 |
| 2011-08-11 | Query 1 |
| 2011-08-11 | Query 2 |
| 2011-08-12 | Query 3 |
| 2011-08-12 | Query 3 |
| 2011-08-13 | Query 1 |
+------------+---------+

我需要将其转换为将每个唯一查询转换为按日期分组的列,并将每个查询的计数插入数据框的行中。我希望输出是这样的:

+------------+---------+---------+---------+
|    date    | Query 1 | Query 2 | Query 3 |
+------------+---------+---------+---------+
| 2011-08-11 |       2 |       1 |       0 |
| 2011-08-12 |       0 |       0 |       2 |
| 2011-08-13 |       1 |       0 |       0 |
+------------+---------+---------+---------+

我正在尝试以this answer 为例,但我不太了解代码,尤其是make_row 函数中的return 语句。

有没有办法在转换 DataFrame 时计算查询? 也许像

import pyspark.sql.functions as func

grouped = (df
    .map(lambda row: (row.date, (row.query, func.count(row.query)))) # Just an example. Not sure how to do this.
    .groupByKey())

这是一个可能包含数十万行和查询的数据框,因此我更喜欢 RDD 版本而不是使用 .collect() 的选项

谢谢!

【问题讨论】:

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


    【解决方案1】:

    您可以使用groupBy.pivotcount 作为聚合函数:

    from pyspark.sql.functions import count
    df.groupBy('date').pivot('query').agg(count('query')).na.fill(0).orderBy('date').show()
    
    +--------------------+-------+-------+-------+
    |                date|Query 1|Query 2|Query 3|
    +--------------------+-------+-------+-------+
    |2011-08-11 00:00:...|      2|      1|      0|
    |2011-08-12 00:00:...|      0|      0|      2|
    |2011-08-13 00:00:...|      1|      0|      0|
    +--------------------+-------+-------+-------+
    

    【讨论】:

    • 我正在努力使用该格式的 DataFrame 执行一些操作。我将如何使用这个命令来创建这个 DataFrame,但是以一种转置的方式? (即标题作为日期(时间戳)并将每个查询作为新行?我尝试过queries_df = df.groupBy('query').pivot('query_time').agg(count('query')).na.fill(0),但我无法将“x 轴”上的日期作为标题。
    • 为我工作;如果您的意思是日期格式不正确,您可能需要先将日期转换为字符串。 df.withColumn("date", date_format("date", "YYYY-MM-dd")).groupBy('query').pivot('date').agg(count('query')).na.fill(0) 并导入 from pyspark.sql.functions import date_format
    • 哦,对不起!你说得对!再次感谢您:) 最后一个后续问题......现在我正在尝试对除第一列之外的所有列求和,并将结果添加到 DataFrame 中的新列中。 import pyspark.sql.functions as F 然后是newDF = queries_df.withColumn('my_sum', F.sum(queries_df[i] for i in queries_df.columns[1:])).show() 但这给了我TypeError: Column is not iterable。有什么想法吗?
    • 我在post问了一个更详细的问题
    猜你喜欢
    • 1970-01-01
    • 2022-09-24
    • 2017-11-02
    • 2018-05-27
    • 1970-01-01
    • 1970-01-01
    • 2022-06-11
    • 2021-09-25
    • 2020-02-20
    相关资源
    最近更新 更多