【问题标题】:Using Spark filter a data frame with conditions使用 Spark 过滤带有条件的数据框
【发布时间】:2017-09-06 06:51:26
【问题描述】:

我有一个看起来像

的数据框
scala> val df = sc.parallelize(Seq(("User 1","X"), ("User 2", "Y"), ("User 3", "X"), ("User 2", "E"), ("User 3", "E"))).toDF("user", "event")

scala> df.show
+------+-----+
|  user|event|
+------+-----+
|User 1|    X|
|User 2|    Y|
|User 3|    X|
|User 2|    E|
|User 3|    E|
+------+-----+

我想找到所有有事件“X”但没有事件“E”的用户

在这种情况下,只有“用户 1”符合条件,因为它没有事件“E”条目。我如何使用 Spark API 来做到这一点?

【问题讨论】:

  • 创建 X 和 E 的 2 个 df 并使用不等于条件加入它们

标签: scala apache-spark apache-spark-sql


【解决方案1】:

可以使用左连接:

val xDF = df.filter(col("event") === "X")
val eDF = df.filter(col("event") === "E")
val result = xDF.as("x").join(eDF.as("e"), List("user"), "left_outer").where(col("e.event").isNull).select(col("x.user"))

结果是:

+------+
|user  |
+------+
|User 1|
+------+

【讨论】:

  • 加入很贵。
【解决方案2】:

您可以通过事件集合对用户进行分组,然后根据特定条件过滤出适合用户的事件。

val result = df.groupBy("user")
    .agg(collect_list("event")
    .as("events"))
    .filter( p => p.getList(1).contains("X") && !p.getList(1).contains("E"))

【讨论】:

    【解决方案3】:
    val tmp = df.groupBy("user").pivot("event").count
    tmp.show
    +------+----+----+----+
    |  user|   E|   X|   Y|
    +------+----+----+----+
    |User 2|   1|null|   1|
    |User 3|   1|   1|null|
    |User 1|null|   1|null|
    +------+----+----+----+
    tmp.filter(  ($"X" isNotNull) and ($"E" isNull) ).show
    +------+----+---+----+
    |  user|   E|  X|   Y|
    +------+----+---+----+
    |User 1|null|  1|null|
    +------+----+---+----+
    tmp.filter(  ($"X" isNotNull) and ($"E" isNull) ).select("user","X").show 
    +------+---+
    |  user|  X|
    +------+---+
    |User 1|  1|
    +------+---+
    

    希望这会有所帮助

    【讨论】:

    • tmp.filter( ($"X" isNotNull) and ($"E" isNull) ).select("user","X").show 已添加!谢谢你的建议。
    【解决方案4】:

    您可以统计每个用户的行数,统计用户和事件的每一行,并过滤两个计数相等且事件列具有 X 值的行。

    import org.apache.spark.sql.functions._
    import org.apache.spark.sql.expressions.Window
    df.withColumn("count", count($"user").over(Window.partitionBy("user")))
        .withColumn("distinctCount", count($"user").over(Window.partitionBy("user", "event")))
        .filter($"count" === $"distinctCount" && $"event" === "X")
        .drop("count", "distinctCount")
    

    你应该得到你想要的结果

    希望回答对你有帮助

    【讨论】:

      猜你喜欢
      • 2021-07-26
      • 2018-09-18
      • 1970-01-01
      • 2020-09-04
      • 1970-01-01
      • 2023-01-04
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多