【问题标题】:spark read doesn't work inside Scala UDF functionspark read 在 Scala UDF 函数中不起作用
【发布时间】:2019-04-14 16:51:50
【问题描述】:

我正在尝试使用 spark.read 来获取我的 UDF 中的文件计数,但是当我执行程序时,此时会挂起。

我在数据框的列中调用 UDF。 udf 必须读取一个文件并返回它的计数。但它不起作用。我将一个变量值传递给 UDF 函数。当我删除 spark.read 代码并简单地返回它工作的数字时。但 spark.read 不能通过 UDF 工作

def prepareRowCountfromParquet(jobmaster_pa: String)(implicit spark: SparkSession): Int = {
      print("The variable value is " + jobmaster_pa)
      print("the count is " + spark.read.format("csv").option("header", "true").load(jobmaster_pa).count().toInt)
      spark.read.format("csv").option("header", "true").load(jobmaster_pa).count().toInt
    }
val SRCROWCNT = udf(prepareRowCountfromParquet _)

  df
  .withColumn("SRC_COUNT", SRCROWCNT(lit(keyPrefix))) 

SRC_COUNT 列应该获取文件的行

【问题讨论】:

标签: scala apache-spark


【解决方案1】:

UDF 不能使用 spark 上下文,因为它只存在于驱动程序中并且不可序列化。

一般来说,您需要阅读所有 csv,使用 groupBy 计算计数,然后您可以对 df 进行左连接

【讨论】:

  • 谢谢阿农。我没有调用 udf,而是将 spark 读入列中。有效。感谢您让我知道这个概念
猜你喜欢
  • 2016-12-02
  • 2017-09-22
  • 1970-01-01
  • 2018-02-17
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-11-22
相关资源
最近更新 更多