【问题标题】:Dataframe to RDD piece of code is not working数据框到 RDD 代码段不起作用
【发布时间】:2020-08-16 23:36:24
【问题描述】:

我正在尝试读取每一行数据帧并将行数据转换为自定义 bean 类。但这里的问题是,代码没有被执行。为了检查,我编写了多个打印语句,但df.rdd.map{row=>} 中的打印语句都没有执行,就好像整个代码块被转义了一样。

代码sn-p:

 print("data frame:", df.show()). 

 df.rdd.map(row => {
   // Debugging
   println("Debugging")

  if(row.isNullAt(0)) {
    println("null data")
  } else {
    println(row.get(0).toString)
  }

  val employeeJobData = new EmployeeJobData

  if(row.get(0).toString == null || row.get(0).toString.isEmpty){
    employeeJobData.setEmployeeId("NULL_KEY_VALUE")
  } else {
    employeeJobData.setEmployeeId(row.get(0).toString)
  }
  employeeJobDataList.add(employeeJobData)
  } )

df.show()的输出:

   |employee_id|employee_name|employee_email|paygroup|level|dept_id|
   +-----------+-------------+--------------+--------+-----+-------+
   |13         |         null|          null|    null| null|   null|
   |14         |         null|          null|    null| null|   null|
   |15         |         null|          null|    null| null|   null|
   |16         |         null|          null|    null| null|   null|
   |17         |         null|          null|    null| null|   null|
   +-----------+-------------+--------------+--------+-----+-------+

【问题讨论】:

  • 你能在这里发布完整的代码吗?
  • 在Spark中,在执行任何收集操作之前,它不会执行代码。
  • ...如果/当它被执行,你应该会在执行者的日志中看到输出,因为map 将被分发。

标签: scala apache-spark apache-spark-sql rdd


【解决方案1】:

您可以如下删除不必要的代码并获得java.util.List[EmployeeJobData] 如下

import java.util

object MapToCaseClass {

  def main(args: Array[String]): Unit = {
    val spark = Constant.getSparkSess;

    import spark.implicits._

    val df  = List((12,"name","email@email.com","paygroup","level","dept_id")).toDF()
    val employeeList : util.List[EmployeeJobData] = df
      .map(row => {
        val id = if (null == row.getString(0) || "null".equals(row.getString(0)) || row.getString(0).trim.isEmpty) {
          "NULL_KEY_VALUE"
        } else {
          row.getString(0)
        }
        EmployeeJobData(id, row.getString(1), row.getString(2),
          row.getString(3), row.getString(4), row.getString(5))
      })
      .collectAsList
  }

}

case class EmployeeJobData(employee_id: String, employee_name: String,employee_email: String,paygroup: String,
                           level: String,dept_id: String)

只需将employee_iddept_id 的数据类型(即如果是数字)设置为Long,就可以进一步改善上述情况。对于employee_id,可以避免"null".equals.isEmpty(),并且可以进一步减少代码。

【讨论】:

    猜你喜欢
    • 2017-05-03
    • 1970-01-01
    • 1970-01-01
    • 2019-12-11
    • 1970-01-01
    • 1970-01-01
    • 2020-06-03
    • 2015-11-12
    • 1970-01-01
    相关资源
    最近更新 更多