【问题标题】:DataFrame/Dataset join not producing correct results in Spark 2.0/YarnDataFrame/Dataset join 在 Spark 2.0/Yarn 中没有产生正确的结果
【发布时间】:2017-02-15 09:23:03
【问题描述】:

我们有一个在 Hadoop 2.7.2、Centos 7.2 上运行 Apache Spark 2.0 的集群。我们使用 Spark DataFrame/DataSet API 编写了一些新代码,但在将数据写入 Windows Azure 存储 Blob(默认 HDFS 位置)然后读取数据后,我们注意到连接结果不正确。我已经能够通过在集群上运行的以下 sn-p 代码复制该问题。

case class UserDimensions(user: Long, dimension: Long, score: Double)
case class CentroidClusterScore(dimension: Long, cluster: Int, score: Double)

val dims = sc.parallelize(Array(UserDimensions(12345, 0, 1.0))).toDS
val cent = sc.parallelize(Array(CentroidClusterScore(0, 1, 1.0),CentroidClusterScore(1, 0, 1.0),CentroidClusterScore(2, 2, 1.0))).toDS

dims.show
cent.show
dims.join(cent, dims("dimension") === cent("dimension") ).show

输出

+-----+---------+-----+                                                         
| user|dimension|score|
+-----+---------+-----+
|12345|        0|  1.0|
+-----+---------+-----+

+---------+-------+-----+
|dimension|cluster|score|
+---------+-------+-----+
|        0|      1|  1.0|
|        1|      0|  1.0|
|        2|      2|  1.0|
+---------+-------+-----+

+-----+---------+-----+---------+-------+-----+
| user|dimension|score|dimension|cluster|score|
+-----+---------+-----+---------+-------+-----+
|12345|        0|  1.0|        0|      1|  1.0|
+-----+---------+-----+---------+-------+-----+

这是正确的。然而在写入和读取数据之后,我们看到了这个

dims.write.mode("overwrite").save("/tmp/dims2.parquet")
cent.write.mode("overwrite").save("/tmp/cent2.parquet")

val dims2 = spark.read.load("/tmp/dims2.parquet").as[UserDimensions]
val cent2 = spark.read.load("/tmp/cent2.parquet").as[CentroidClusterScore]

dims2.show
cent2.show

dims2.join(cent2, dims2("dimension") === cent2("dimension") ).show

输出

+-----+---------+-----+                                                         
| user|dimension|score|
+-----+---------+-----+
|12345|        0|  1.0|
+-----+---------+-----+

+---------+-------+-----+
|dimension|cluster|score|
+---------+-------+-----+
|        0|      1|  1.0|
|        1|      0|  1.0|
|        2|      2|  1.0|
+---------+-------+-----+

+-----+---------+-----+---------+-------+-----+
| user|dimension|score|dimension|cluster|score|
+-----+---------+-----+---------+-------+-----+
|12345|        0|  1.0|     null|   null| null|
+-----+---------+-----+---------+-------+-----+

但是,使用 RDD API 会产生正确的结果

dims2.rdd.map( row => (row.dimension, row) ).join( cent2.rdd.map( row => (row.dimension, row) ) ).take(5)

res5: Array[(Long, (UserDimensions, CentroidClusterScore))] = Array((0,(UserDimensions(12345,0,1.0),CentroidClusterScore(0,1,1.0))))

我们尝试将输出格式更改为 ORC 而不是 parquet,但我们看到了相同的结果。在本地而非集群上运行 Spark 2.0 不会出现此问题。在 Hadoop 集群的主节点上以本地模式运行 spark 也可以。只有在 YARN 上运行时,我们才会看到这个问题。

这似乎也与这个问题非常相似:https://issues.apache.org/jira/browse/SPARK-10896

【问题讨论】:

    标签: apache-spark apache-spark-sql apache-spark-dataset


    【解决方案1】:

    https://issues.apache.org/jira/browse/SPARK-17806 中提交的拉取请求已修复此问题

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2015-06-23
      • 2017-04-03
      • 2018-04-22
      • 2014-03-14
      • 2011-01-19
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多