【问题标题】:Filter is not working as expected when it is applied on a DF(Which is a union of 2 DF's) in Spark过滤器在 Spark 中应用于 DF(这是 2 个 DF 的联合)时未按预期工作
【发布时间】:2017-02-27 22:25:49
【问题描述】:

数据框a:

SN  Hash_id Name
111 11ww11  Airtel
222 null    Idea

数据框 b:

SN  Hash_id Name
333 null    BSNL
444 22ee11  Vodafone

按列名对这些数据帧执行 UnionAll 如下:

def unionByName(a: DataFrame, b: DataFrame): DataFrame = {
    val columns = a.columns.toSet.intersect(b.columns.toSet).map(col).toSeq
    a.select(columns: _*).unionAll(b.select(columns: _*))
} 

结果是:数据框c

SN  Hash_id Name
111 11ww11  Airtel
222 null    Idea
333 null    BSNL
444 22ee11  Vodafone

对数据框 c 执行过滤。

val withHashDF = c.where(c("Hash_id").isNotNull)
val withoutHashDF = c.where(c("Hash_id").isNull)

withHashDF 的结果是:它只给出数据框 a 的结果

111 11ww11  Airtel

数据帧 b 的记录在哈希 id 存在的地方丢失:

444 22ee11  Vodafone

withoutHashDF 的结果是:

222 null    Idea
BSNL 333    null    
null  222    Idea

在此 DF 中,列值与列名不同,计数为 3,应仅为 2。从数据框中“a”行重复。

【问题讨论】:

  • 它应该可以正常工作。似乎是调用 unionByName 方法以及填充数据框 c 的位置的问题。

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


【解决方案1】:

unionByName 获取columns有小变化

改变

val columns = a.columns.toSet.intersect(b.columns.toSet).map(col).toSeq

val columns = a.columns.intersect(b.columns).map(row => new Column(row)).toSeq

那么它应该按预期工作。

看看下面的完整代码sn-p & 结果:

import sparkSession.sqlContext.implicits._
import org.apache.spark.sql.DataFrame
import org.apache.spark.sql.Column

val dataFrameA = Seq(("111", "11ww11", "Airtel"),("222", null, "Idea")).toDF("SN","Hash_id", "Name")
val dataFrameB = Seq(("333", null, "BSNL"),("444", "22ee11", "Vodafone")).toDF("SN","Hash_id", "Name")

def unionByName(a: DataFrame, b: DataFrame): DataFrame = {
  val columns = a.columns.intersect(b.columns).map(row => new Column(row)).toSeq
  a.select(columns: _*).union(b.select(columns: _*))
}

val dataFrameC = unionByName(dataFrameA, dataFrameB)
val withHashDF = dataFrameC.where(dataFrameC("Hash_id").isNotNull)
val withoutHashDF = dataFrameC.where(dataFrameC("Hash_id").isNull)

println("dataFrameC")
dataFrameC.show()

println("withHashDF")
withHashDF.show

println("withoutHashDF")
withoutHashDF.show

输出:

dataFrameC
+---+-------+--------+
| SN|Hash_id|    Name|
+---+-------+--------+
|111| 11ww11|  Airtel|
|222|   null|    Idea|
|333|   null|    BSNL|
|444| 22ee11|Vodafone|
+---+-------+--------+

withHashDF
+---+-------+--------+
| SN|Hash_id|    Name|
+---+-------+--------+
|111| 11ww11|  Airtel|
|444| 22ee11|Vodafone|
+---+-------+--------+

withoutHashDF
+---+-------+----+
| SN|Hash_id|Name|
+---+-------+----+
|222|   null|Idea|
|333|   null|BSNL|
+---+-------+----+

【讨论】:

  • 更改 unionbyname 后,在我的情况下结果也不正确。我怀疑正在发生脏读。但不确定如何处理。
【解决方案2】:

如果 Dataframe(Unionall) 中有重复项,它会为过滤器或 where 子句提供意外结果。使用 distinct 方法消除重复项后,结果符合预期。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2022-01-23
    • 2011-05-22
    • 1970-01-01
    • 2021-09-07
    • 1970-01-01
    • 2014-05-15
    • 1970-01-01
    相关资源
    最近更新 更多