【问题标题】:Error including a column in a join between spark dataframes在火花数据帧之间的连接中包含列时出错
【发布时间】:2020-09-12 11:30:47
【问题描述】:

我在 cleanDFsentiment_df 之间使用 array_contains 连接,效果很好 (from solution 61687997)。我需要在 Result df 中包含一个来自 cleanDF 的新列(“年份”)。

这是连接:

from pyspark.sql import functions

Result = cleanDF.join(sentiment_df, expr("""array_contains(MeaningfulWords,word)"""), how='left')\
                .groupBy("ID")\
                .agg(first("MeaningfulWords").alias("MeaningfulWords")\
                  ,collect_list("score").alias("ScoreList")\
                  ,mean("score").alias("MeanScore"))

这是 Result 结构:

Result.show(5)

#+------------------+--------------------+--------------------+-----------------+
#|                ID|     MeaningfulWords|           ScoreList|        MeanScore|
#+------------------+--------------------+--------------------+-----------------+
#|a0U3Y00000p1IzjUAE|[buen, servicio, ...|        [6.39, 1.82]|            4.105|
#|a0U3Y00000p1KhGUAU|              [mala]|              [2.02]|             2.02|
#|a0U3Y00000p1M1oUAE|[cliente, content...|        [6.39, 8.41]|              7.4|
#|a0U3Y00000p1OnTUAU|[positivo, trato,...|               [8.2]|             8.19|
#|a0U3Y00000p1R5DUAU|[momento, sido, g...|               [6.0]|              6.0|
#+------------------+--------------------+--------------------+-----------------+

我添加了一个.select (36132322) 以包含来自cleanDF的列Year

Result1 = cleanDF.alias('a').join(sentiment_df.alias('b'), expr("""array_contains(a.MeaningfulWords,b.word)"""), how='left')\
                .select(col('a.ID'),col('a.Year'),col('a.MeaningfulWords'),col('b.word'),col('b.score'))\
                .groupBy("ID")\
                .agg(first("a.MeaningfulWords").alias("MeaningfulWords")\
                  ,collect_list("score").alias("ScoreList")\
                  ,mean("score").alias("MeanScore"))

但我进入 Result1**Result** 相同的列:

display(Result1)

#DataFrame[ID: string, MeaningfulWords: array<string>, ScoreList: array<double>, MeanScore: double]

当我尝试在 .agg 函数中包含 Year 时:

Result2 = cleanDF.join(sentiment_df, expr("""array_contains(MeaningfulWords,word)"""), how='left')\
                .groupBy("ID")\
                .agg(first("MeaningfulWords").alias("MeaningfulWords"),first("Year").alias("Year")\
                  ,collect_list("score").alias("ScoreList")\
                  ,mean("score").alias("MeanScore"))

Result2.show()

Py4JJavaError: An error occurred while calling o3205.showString.
: org.apache.spark.SparkException: Exception thrown in awaitResult: 
    at org.apache.spark.util.ThreadUtils$.awaitResult(ThreadUtils.scala:226)
    at org.apache.spark.sql.execution.exchange.BroadcastExchangeExec.doExecuteBroadcast(BroadcastExchangeExec.scala:146)
    at org.apache.spark.sql.execution.SparkPlan$$anonfun$executeBroadcast$1.apply(SparkPlan.scala:144)
    at org.apache.spark.sql.execution.SparkPlan$$anonfun$executeBroadcast$1.apply(SparkPlan.scala:140)
    at org.apache.spark.sql.execution.SparkPlan$$anonfun$executeQuery$1.apply(SparkPlan.scala:155)
    at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151)
    at org.apache.spark.sql.execution.SparkPlan.executeQuery(SparkPlan.scala:152)
    at org.apache.spark.sql.execution.SparkPlan.executeBroadcast(SparkPlan.scala:140)
    at org.apache.spark.sql.execution.joins.BroadcastNestedLoopJoinExec.doExecute(BroadcastNestedLoopJoinExec.scala:343)
    ...
    ...
    ...
        Caused by: org.apache.spark.SparkException: Failed to execute user defined function($anonfun$createTransformFunc$1: (string) => array<string>)
            at org.apache.spark.sql.catalyst.expressions.ScalaUDF.eval(ScalaUDF.scala:1066)
            at org.apache.spark.sql.catalyst.expressions.ScalaUDF$$anonfun$2.apply(ScalaUDF.scala:109)
            at org.apache.spark.sql.catalyst.expressions.ScalaUDF$$anonfun$2.apply(ScalaUDF.scala:107)
            at org.apache.spark.sql.catalyst.expressions.ScalaUDF.eval(ScalaUDF.scala:1063)
    ...
    ...
    Caused by: org.apache.spark.SparkException: Job aborted due to stage failure: Task 2 in stage 411.0 failed 1 times, most recent failure: Lost task 2.0 in stage 411.0 (TID 9719, localhost, executor driver): org.apache.spark.SparkException: Failed to execute user defined function($anonfun$5: (array<string>) => array<string>)
        at org.apache.spark.sql.catalyst.expressions.ScalaUDF.eval(ScalaUDF.scala:1066)
        at org.apache.spark.sql.catalyst.expressions.SimpleHigherOrderFunction$class.eval(higherOrderFunctions.scala:208)
        at org.apache.spark.sql.catalyst.expressions.ArrayFilter.eval(higherOrderFunctions.scala:296)
        at org.apache.spark.sql.catalyst.expressions.GeneratedClass$SpecificUnsafeProjection.apply(Unknown Source)
        at org.apache.spark.sql.catalyst.expressions.GeneratedClass$SpecificUnsafeProjection.apply(Unknown Source)
    ...
    ...
    ... 20 more
    Caused by: java.lang.NullPointerException

我在 spark 2.4.5 上使用 pyspark。

提前感谢您的帮助。

【问题讨论】:

  • 从异常中,我可以看到问题是 java.lang.NullPointerException,如果是,您是否在任何一个数据框中都有该列,您可以过滤空值.. 还显示两个数据框的架构跨度>

标签: apache-spark pyspark apache-spark-sql


【解决方案1】:

Year 列可能有空值,因为它因Caused by: java.lang.NullPointerException 异常而失败。过滤Year 列中的所有空值。

【讨论】:

  • 感谢@Srinivas,我在“Year”列中输入了空值,然后 java.lang.NullPointerException 已解决,现在我在 Result2 中有“Year”列。是否有一些替代方法可以在此联接中包含具有 null 'Year' 的记录?
  • 用零或负数代替null。
  • 你也可以发布 result2 df printschema 吗??
  • 现在我在 Result2 中已经有了“Year”列(我将空值估算为默认值)。这是结构:DataFrame[ID: string, MeaningfulWords: array&lt;string&gt;, Year: int, ScoreList: array&lt;double&gt;, MeanScore: double]。我只是怀疑,如果存在一些替代方法来在此连接中包含具有 null 'Year' 的行,不估算 null 值。提前感谢@Srinivas!
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2020-01-21
  • 2017-05-24
  • 1970-01-01
  • 2016-05-06
  • 2018-11-08
  • 2016-09-23
  • 2019-02-01
相关资源
最近更新 更多