【问题标题】:How to join two DataFrames and change column for missing values?如何连接两个 DataFrame 并更改缺失值的列?
【发布时间】:2017-04-21 03:59:53
【问题描述】:
val df1 = sc.parallelize(Seq(
   ("a1",10,"ACTIVE","ds1"),
   ("a1",20,"ACTIVE","ds1"),
   ("a2",50,"ACTIVE","ds1"),
   ("a3",60,"ACTIVE","ds1"))
).toDF("c1","c2","c3","c4")`

val df2 = sc.parallelize(Seq(
   ("a1",10,"ACTIVE","ds2"),
   ("a1",20,"ACTIVE","ds2"),
   ("a1",30,"ACTIVE","ds2"),
   ("a1",40,"ACTIVE","ds2"),
   ("a4",20,"ACTIVE","ds2"))
).toDF("c1","c2","c3","c5")`


df1.show()

// +---+---+------+---+
// | c1| c2|    c3| c4|
// +---+---+------+---+
// | a1| 10|ACTIVE|ds1|
// | a1| 20|ACTIVE|ds1|
// | a2| 50|ACTIVE|ds1|
// | a3| 60|ACTIVE|ds1|
// +---+---+------+---+

df2.show()
// +---+---+------+---+
// | c1| c2|    c3| c5|
// +---+---+------+---+
// | a1| 10|ACTIVE|ds2|
// | a1| 20|ACTIVE|ds2|
// | a1| 30|ACTIVE|ds2|
// | a1| 40|ACTIVE|ds2|
// | a4| 20|ACTIVE|ds2|
// +---+---+------+---+

我的要求是:我需要加入两个数据框。 我的输出数据框应该包含来自 df1 的所有记录以及来自 df2 的记录,这些记录不在 df1 中,仅用于匹配的“c1”。我从 df2 中提取的记录应该在“c3”列更新为非活动状态。

在这个例子中,只有“c1”的匹配值是a1。所以我需要从 df2 中提取 c2=30 和 40 条记录并使其处于非活动状态。

这是输出。

df_output.show()

// +---+---+--------+---+
// | c1| c2|    c3  | c4|
// +---+---+--------+---+
// | a1| 10|ACTIVE  |ds1|
// | a1| 20|ACTIVE  |ds1|
// | a2| 50|ACTIVE  |ds1|
// | a3| 60|ACTIVE  |ds1|
// | a1| 30|INACTIVE|ds1|
// | a1| 40|INACTIVE|ds1|
// +---+---+--------+---+

谁能帮我做这件事。

【问题讨论】:

  • 对于 INACTIVE 记录,c4 值是否从 ds2 更改为 ds1?

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


【解决方案1】:

首先,一件小事。我为df2 中的列使用不同的名称:

val df2 = sc.parallelize(...).toDF("d1","d2","d3","d4")

没什么大不了的,但这让我更容易推理。

现在是有趣的东西。为了清楚起见,我会有点冗长:

val join = df1
.join(df2, df1("c1") === df2("d1"), "inner")
.select($"d1", $"d2", $"d3", lit("ds1").as("d4"))
.dropDuplicates

在这里,我执行以下操作:

  • c1d1 列上的 df1df2 之间的内连接
  • 选择df2 列并在最后一列中简单地“硬编码”ds1 以替换ds2
  • 删除重复项

这基本上只是过滤掉df2 中所有没有df1 中的c1 中具有对应键的所有内容。

接下来我比较:

val diff = join
.except(df1)
.select($"d1", $"d2", lit("INACTIVE").as("d3"), $"d4")

这是一个基本的集合操作,它在join 中找到所有不是df1 中的东西。这些是要停用的项目,因此我选择了所有列,但将第三列替换为硬编码的 INACTIVE 值。

剩下的就是把它们放在一起:

df1.union(diff)

这只是简单地将df1 与我们之前计算的停用值表结合起来以产生最终结果:

+---+---+--------+---+
| c1| c2|      c3| c4|
+---+---+--------+---+
| a1| 10|  ACTIVE|ds1|
| a1| 20|  ACTIVE|ds1|
| a2| 50|  ACTIVE|ds1|
| a3| 60|  ACTIVE|ds1|
| a1| 30|INACTIVE|ds1|
| a1| 40|INACTIVE|ds1|
+---+---+--------+---+

同样,您不需要所有这些中间值。我只是冗长地帮助跟踪整个过程。

【讨论】:

  • val c1Ids = df1.select("c1").as[String].collect()val joinDf = df1.as("t1").join(df2.as("t2"), df1("c1") === df2("c1"), "rightouter").select($"t2.c1", $"t2.c2").distinct().withColumn("c3", lit("INACTIVE")).withColumn("c4",lit("ds1")).filter(not($"c1").isin(c1Ids: _*))
  • val finalDf = df1.unionAll(joinDf) 这应该给出输出。
  • 我确信有多种方法(也许比我的更好)来获得该输出。很高兴为您指明正确的方向。祝你的项目好运!
【解决方案2】:

这是肮脏的解决方案 -

from pyspark.sql import functions as F


# find the rows from df2 that have matching key c1 in df2
df3 = df1.join(df2,df1.c1==df2.c1)\
.select(df2.c1,df2.c2,df2.c3,df2.c5.alias('c4'))\
.dropDuplicates()

df3.show()

+---+---+------+---+
| c1| c2|    c3| c4|
+---+---+------+---+
| a1| 10|ACTIVE|ds2|
| a1| 20|ACTIVE|ds2|
| a1| 30|ACTIVE|ds2|
| a1| 40|ACTIVE|ds2|
+---+---+------+---+

# Union df3 with df1 and change columns c3 and c4 if c4 value is 'ds2'

df1.union(df3).dropDuplicates(['c1','c2'])\
.select('c1','c2',\
        F.when(df1.c4=='ds2','INACTIVE').otherwise('ACTIVE').alias('c3'),
        F.when(df1.c4=='ds2','ds1').otherwise('ds1').alias('c4')
       )\
.orderBy('c1','c2')\
.show()

+---+---+--------+---+
| c1| c2|      c3| c4|
+---+---+--------+---+
| a1| 10|  ACTIVE|ds1|
| a1| 20|  ACTIVE|ds1|
| a1| 30|INACTIVE|ds1|
| a1| 40|INACTIVE|ds1|
| a2| 50|  ACTIVE|ds1|
| a3| 60|  ACTIVE|ds1|
+---+---+--------+---+

【讨论】:

    【解决方案3】:

    享受挑战,这是我的解决方案。

    val c1keys = df1.select("c1").distinct
    val df2_in_df1 = df2.join(c1keys, Seq("c1"), "inner")
    val df2inactive = df2_in_df1.join(df1, Seq("c1", "c2"), "leftanti").withColumn("c3", lit("INACTIVE"))
    scala> df1.union(df2inactive).show
    +---+---+--------+---+
    | c1| c2|      c3| c4|
    +---+---+--------+---+
    | a1| 10|  ACTIVE|ds1|
    | a1| 20|  ACTIVE|ds1|
    | a2| 50|  ACTIVE|ds1|
    | a3| 60|  ACTIVE|ds1|
    | a1| 30|INACTIVE|ds2|
    | a1| 40|INACTIVE|ds2|
    +---+---+--------+---+
    

    【讨论】:

      猜你喜欢
      • 2017-09-14
      • 1970-01-01
      • 1970-01-01
      • 2023-01-13
      • 2022-01-12
      • 1970-01-01
      • 1970-01-01
      • 2022-09-27
      • 1970-01-01
      相关资源
      最近更新 更多