【问题标题】:Spark extracting values from a RowSpark从行中提取值
【发布时间】:2016-01-05 14:33:12
【问题描述】:

我有以下数据框

val transactions_with_counts = sqlContext.sql(
  """SELECT user_id AS user_id, category_id AS category_id,
  COUNT(category_id) FROM transactions GROUP BY user_id, category_id""")

我正在尝试将行转换为 Rating 对象,但由于 x(0) 返回一个数组,因此失败

val ratings = transactions_with_counts
  .map(x => Rating(x(0).toInt, x(1).toInt, x(2).toInt))

错误:值 toInt 不是 Any 的成员

【问题讨论】:

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


    【解决方案1】:

    让我们从一些虚拟数据开始:

    val transactions = Seq((1, 2), (1, 4), (2, 3)).toDF("user_id", "category_id")
    
    val transactions_with_counts = transactions
      .groupBy($"user_id", $"category_id")
      .count
    
    transactions_with_counts.printSchema
    
    // root
    // |-- user_id: integer (nullable = false)
    // |-- category_id: integer (nullable = false)
    // |-- count: long (nullable = false)
    

    有几种方法可以访问Row 值并保留预期类型:

    1. 模式匹配

      import org.apache.spark.sql.Row
      
      transactions_with_counts.map{
        case Row(user_id: Int, category_id: Int, rating: Long) =>
          Rating(user_id, category_id, rating)
      } 
      
    2. 键入get* 方法,例如getIntgetLong

      transactions_with_counts.map(
        r => Rating(r.getInt(0), r.getInt(1), r.getLong(2))
      )
      
    3. getAs 可以同时使用名称和索引的方法:

      transactions_with_counts.map(r => Rating(
        r.getAs[Int]("user_id"), r.getAs[Int]("category_id"), r.getAs[Long](2)
      ))
      

      它可以用来正确提取用户定义的类型,包括mllib.linalg.Vector。显然,按名称访问需要架构。

    4. 转换为静态类型的Dataset(Spark 1.6+ / 2.0+):

      transactions_with_counts.as[(Int, Int, Long)]
      

    【讨论】:

    • 您提到的上述四种方法中,哪种方法最有效......?
    • @Dilan 模式匹配静态类型选项可能会更慢(后者有一些其他的性能影响)。 getAs[_]get* 应该相似,但使用起来很痛苦。
    • 1. “后者具有其他性能含义”是什么意思......? 2. getAs[_] 和 get* 在性能上是否比模式匹配更好?
    • 我正在使用上面描述的第一种方法 DataFrame 具有 nullable 列,例如 case Row(usrId: Int, usrName: String, null, usrMobile: Int) => ...case Row(usrId: Int, usrName: String, usrAge: Int, null) => ... 这会导致长 case 表达式 (我有几个案例)。有没有更简洁的方法(更简洁、更少样板/重复的东西)来做到这一点?请举例回答。
    • 0323 先生的出色回答。
    【解决方案2】:

    使用数据集,您可以按如下方式定义评级:

    case class Rating(user_id: Int, category_id:Int, count:Long)
    

    这里的 Rating 类有一个列名“count”而不是 zero323 建议的“rating”。因此评级变量分配如下:

    val transactions_with_counts = transactions.groupBy($"user_id", $"category_id").count
    
    val rating = transactions_with_counts.as[Rating]
    

    这样您就不会在 Spark 中遇到运行时错误,因为您的 评级类列名称与 Spark 在运行时生成的“计数”列名称相同。

    【讨论】:

      【解决方案3】:

      要访问一行Dataframe的值,您需要使用Dataframerdd.collect和for循环。

      考虑您的 Dataframe 如下所示。

      val df = Seq(
            (1,"James"),    
            (2,"Albert"),
            (3,"Pete")).toDF("user_id","name")
      

      在您的数据框之上使用rdd.collectrow 变量将包含 rdd 行类型的 Dataframe 的每一行。要从行中获取每个元素,请使用 row.mkString(",") 它将包含逗号分隔值中的每一行的值。使用split 函数(内置函数),您可以使用索引访问rdd 行的每一列值。

      for (row <- df.rdd.collect)
      {   
          var user_id = row.mkString(",").split(",")(0)
          var category_id = row.mkString(",").split(",")(1)       
      }
      

      dataframe.foreach 循环相比,上面的代码看起来要大一些,但使用上面的代码可以更好地控制逻辑。

      【讨论】:

      • 为什么要执行收集来应用转换?如果数据不适合内存,这将导致驱动程序崩溃
      • 投反对票,因为据我所知,为每个值拆分行效率不高。使用 get* 方法将避免对每个值进行拆分
      猜你喜欢
      • 2017-07-06
      • 2015-11-04
      • 1970-01-01
      • 1970-01-01
      • 2018-06-12
      • 1970-01-01
      • 2022-11-11
      • 2019-09-22
      • 2021-01-03
      相关资源
      最近更新 更多