【问题标题】:type mismatch; found : Unit required: Array[org.apache.spark.sql.Dataset[org.apache.spark.sql.Row]]类型不匹配;发现:所需单位:数组[org.apache.spark.sql.Dataset[org.apache.spark.sql.Row]]
【发布时间】:2018-10-24 21:43:00
【问题描述】:

为什么下面的代码在return语句出现编译错误,

  def getData(queries: Array[String]): Dataset[Row] = {
    val res = spark.read.format("jdbc").jdbc(jdbcUrl, "", props).registerTempTable("")
    return res
  }

错误,

type mismatch; found : Unit required: Array[org.apache.spark.sql.Dataset[org.apache.spark.sql.Row]]

Scala 版本 2.11.11

Spark 版本 2.0.0

编辑: 实际案例

  def getDataFrames(queries: Array[String]) = {
    val jdbcResult = queries.map(query => {
      val tablename = extractTableName(query)
      if (tablename.contains("1")) {
        spark.sqlContext.read.format("jdbc").jdbc(jdbcUrl1, query, props)
      } else {
        spark.sqlContext.read.format("jdbc").jdbc(jdbcUrl2, query, props)
      }
    })
  }

在这里,我想返回迭代的组合输出,如 Array[Dataset[Row]] 或 Array[DataFrame](但 Dataframe 在 2.0.0 中不可用作为依赖项)。上面的代码有什么魔力吗?或者我该怎么做?

【问题讨论】:

  • registerTempTable 返回Unit 你最好删除registerTempTable 并返回Dataframe,你为什么要返回Array[Dataset[Row]]
  • 我有多个查询,我想创建一个数据框数组。但在问题中有一个错误编辑。

标签: scala apache-spark dataframe


【解决方案1】:

您可以将list 中的dataframes 返回为List[Dataframe]

def getData(queries: Array[String]): List[Dataframe] = {
  val res = spark.read.format("jdbc").jdbc(jdbcUrl, "", props)
  //create multiple dataframe from your queries
  val df1 = ???
  val df2 = ???
  val list = List(df1, df2)
  //You can create a list dynamically with list of quries 
  list
}

registerTempTable 返回Unit 你最好删除registerTempTable 并返回Dataframe,并返回一个list 的数据帧。

更新:

这里是您如何使用查询列表返回数据框列表

def getDataFrames(queries: Array[String]): Array[DataFrame] = {
  val jdbcResult = queries.map(query => {
    val tablename = extractTableName(query)
    val dataframe = if (tablename.contains("1")) {
      spark.sqlContext.read.format("jdbc").jdbc("", query, prop)
    } else {
      spark.sqlContext.read.format("jdbc").jdbc("", query, prop)
    }
    dataframe
  })
  jdbcResult
}

我希望这会有所帮助!

【讨论】:

  • 请查找编辑以了解有关我的要求的更多详细信息。
  • @Krishas 如果可行,您可以检查更新的答案吗?
  • 这消除了编译错误。谢谢你。同样对于 spark 2,我们可以将 Array[DataFrame] 替换为 Array[Dataset[Row]]
  • @Krishas 是的,他们都是一样的
【解决方案2】:

从错误消息中可以清楚地看出您的函数中存在类型不匹配。 registerTempTable() api 创建一个内存表,范围为当前会话,并在 SparkSession 处于活动状态之前保持可访问性。

Check the return type of registerTempTable() api here

将您的代码更改为以下内容以删除错误消息:

def getData(queries: Array[String]): Unit = {
    val res = spark.read.format("jdbc").jdbc(jdbcUrl, "", props).registerTempTable("")

  }

更好的方法是编写如下代码:

val tempName: String = "Name_Of_Temp_View" spark.read.format("jdbc").jdbc(jdbcUrl, "", props).createOrReplaceTempView(tempName)

使用 createOrReplaceTempView() as registerTempTable() 自 Spark 2.0.0 起已弃用

根据您的要求的替代解决方案:

def getData(queries: Array[String], spark: SparkSession): Array[DataFrame] = { spark.read.format("jdbc").jdbc(jdbcUrl, "", props).createOrReplaceTempView("Name_Of_Temp_Table") val result: Array[DataFrame] = queries.map(query => spark.sql(query)) result }

【讨论】:

  • Krishas 您可以提供反馈或接受答案
  • 感谢您的解释,我是 scala 的新手。实际上,我在 getData 方法中有两种类型的 SQL(在 sql 数组中),因此我必须在 if 和 else 中调用 spark.read,那么如何结合结果和返回以及数据帧数组
  • Shakar 的回答是正确的,如下图。下次你能不能分享一下实际的功能,以便回答问题的人更好地了解你的问题
  • 当然,很抱歉造成混乱。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-08-02
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多