【问题标题】:how to select within groupby using spark sql如何使用 spark sql 在 groupby 中进行选择
【发布时间】:2020-06-19 22:34:28
【问题描述】:

我正在尝试使用 pyspark sql 在列中选择具有所需值的行,但它给了我错误

下面是我的桌子session

session_id  status    session_desc
----------  ------    -----------------
session1    Old       first_description
session1    Active    last_description
session1    Old       next_description
session1    Active    inter_description
session2    Old       next_description
session2    Old       inter_description

下面是我的 spark sql 查询

spark.sql("select session_id, (CASE WHEN status='Active' THEN session_desc END) AS session_description from session group by session_id").show()

但是我遇到了错误

org.apache.spark.sql.AnalysisException: expression 'session.status' is neither present in the group by, nor is it an aggregate function. Add to group by or wrap in first() (or first_value) if you don't care which value you get.;

我需要如下

session_id  session_description
----------  -------------------
session1    last_description      # can be inter_description as well (I don't care)
session2    null

【问题讨论】:

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


    【解决方案1】:

    在外部查询中使用您的 case statement inside subquery 然后 group by

    Example:

    df.show()
    #+----------+------+-----------------+
    #|session_id|status|     session_Desc|
    #+----------+------+-----------------+
    #|  session1|   Old|first_description|
    #|  session1|Active| last_description|
    #|  session1|   Old| next_description|
    #|  session1|Active|inter_description|
    #|  session2|   Old| next_description|
    #|  session2|   Old|inter_description|
    #+----------+------+-----------------+
    
    spark.sql("select session_id,last(session_desc)session_description from (\
    select session_id,case when status='Active' THEN session_desc END as session_desc from tmp)t \
    group by session_id").\
    show()
    
    #+----------+-------------------+
    #|session_id|session_description|
    #+----------+-------------------+
    #|  session1|  inter_description|
    #|  session2|               null|
    #+----------+-------------------+
    

    【讨论】:

    • 赞成,添加一条不相关的评论,但未来的读者如果正在寻找 pyspark:(df.withColumn("session_desc",F.when(F.col("status")=="Active",F.col("session_desc"))) .groupby("session_id").agg(F.first(F.col("session_desc"),ignorenulls=True).alias("session_desc")).show()) :)
    猜你喜欢
    • 2019-10-19
    • 1970-01-01
    • 2022-01-12
    • 2021-11-24
    • 2014-02-14
    • 1970-01-01
    • 2020-03-28
    • 2021-11-22
    • 1970-01-01
    相关资源
    最近更新 更多