【问题标题】:Efficient way to convert Dataframe to RDD in Scala/SPARK?在 Scala/SPARK 中将 Dataframe 转换为 RDD 的有效方法?
【发布时间】:2016-10-30 20:16:14
【问题描述】:

我有一个dataFrame = [CUSTOMER_ID ,itemType, eventTimeStamp, valueType],通过执行以下操作将其转换为RDD[(String, (String, String, Map[String, Int]))]

 val tempFile = result.map( {
     r => {
         val customerId = r.getAs[String]( "CUSTOMER_ID" )
         val itemType = r.getAs[String]( "itemType" )
         val eventTimeStamp = r.getAs[String]( "eventTimeStamp" )
         val valueType = r.getAs[Map[String, Int]]( "valueType" )
         (customerId, (itemType, eventTimeStamp, valueType))
          }
          } )

由于我的投入很大,这需要很长时间。有什么有效的方法可以将 df 转换为 RDD[(String, (String, String, Map[String, Int]))] 吗?

【问题讨论】:

  • 您的输入有多大?
  • DataFrame转成RDD需要多长时间?
  • 您是否尝试在 DataFrame 上设置不同数量的分区?有什么区别吗?
  • 您是否尝试仅使用 result.rdd,而不使用 .map()?它会产生类似的结果吗?它跑得更快吗?
  • Inout 大小为 7TB

标签: scala apache-spark dataframe rdd


【解决方案1】:

您所描述的操作尽可能便宜。做一些getAs 并分配一些元组几乎是免费的。如果速度变慢,由于您的数据量很大(7T),这可能是不可避免的。另请注意,Catalyst 优化无法在 RDD 上执行,因此在 DataFrame 操作下游包含这种.map 通常会阻止其他 Spark 快捷方式。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2018-03-05
    • 2017-06-13
    • 2017-05-13
    • 2016-04-13
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-06-02
    相关资源
    最近更新 更多