【问题标题】:Filtering inactivated rows in Spark using Scala使用 Scala 过滤 Spark 中的非激活行
【发布时间】:2020-09-07 21:35:28
【问题描述】:

我对 Spark 和 Scala 编程非常陌生,我有一个问题希望一些聪明的人可以帮助我解决。 我有一个名为 users 的表,有 4 列:status、user_id、name、date

行是:

status  user_id name    date
active      1   Peter   2020-01-01
active      2   John    2020-01-01
active      3   Alex    2020-01-01
inactive    1   Peter   2020-02-01
inactive    2   John    2020-01-01

我只需要选择活跃用户。两个用户被停用。只有一个在同一日期被停用。

我的目标是过滤具有非活动状态的行(我可以),并在非活动行与活动行匹配的列时过滤非活动用户。彼得在不同的日期被停用,他没有被过滤。期望的结果是:

1 Peter 2020-01-01
3 Alex 2020-01-01

已过滤非活动状态的行。 John 未激活,因此他的行也被过滤了。

我最接近的是过滤处于非活动状态的用户:

val users = spark.table("db.users")
      .filter(col("status").not Equal("Inactive"))
      .select("user_id", "name", "date")

任何想法或建议如何解决这个问题? 谢谢!

【问题讨论】:

  • 同一日期出现更多次?对于同一用户?
  • 我很确定,notEqual 之间没有空格。 notEqual 是函数。同样在您提供的数据中inactive 不是大小写,在代码中是大小写。请先修正错别字,以便我们了解是否存在实际的潜在问题
  • 是的,可能更多,但为了简化这一点,我们假设行是不同的。
  • 我不允许在提交表单中输入 notEqual。我作为示例给出的代码没有问题,只是不完整。

标签: scala dataframe apache-spark


【解决方案1】:

首先使用 group by 为每个用户和日期检查非活动状态,并将此结果加入原始 df。

val df2 = df.groupBy('user_id, 'date).agg(max('status).as("status"))
  .filter("status = 'inactive'")
  .withColumnRenamed("status", "inactive")

df.join(df2, Seq("user_id", "date"), "left")
  .filter('inactive.isNull)
  .select(df.columns.head, df.columns.tail: _*)
  .show()

+------+-------+-----+----------+
|status|user_id| name|      date|
+------+-------+-----+----------+
|active|      1|Peter|2020-01-01|
|active|      3| Alex|2020-01-01|
+------+-------+-----+----------+

【讨论】:

  • 谢谢!解决方案很棒。一旦我达到 15 声望,就会投票。
猜你喜欢
  • 2017-06-26
  • 2015-06-27
  • 1970-01-01
  • 2016-03-02
  • 2017-08-11
  • 2018-05-16
  • 1970-01-01
  • 2020-03-27
  • 2022-08-20
相关资源
最近更新 更多