【发布时间】: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