【问题标题】:Converting RDD into Dataframe将RDD转换为Dataframe
【发布时间】:2020-01-24 09:32:22
【问题描述】:

我是 spark/scala 的新手。 我通过从多个路径加载数据在 RDD 下创建了一个。现在我想从中创建数据框以进行进一步的操作。 下面应该是数据框的架构

schema[UserId, EntityId, WebSessionId, ProductId]

rdd.foreach(println)

545456,5615615,DIKFH6545614561456,PR5454564656445454
875643,5485254,JHDSFJD543514KJKJ4
545456,5615615,DIKFH6545614561456,PR5454564656445454
545456,5615615,DIKFH6545614561456,PR5454564656445454
545456,5615615,DIKFH6545614561456,PR54545DSKJD541054
264264,3254564,MNXZCBMNABC5645SAD,PR5142545564542515
732543,8765984,UJHSG4240323545144
564574,6276832,KJDXSGFJFS2545DSAS

谁能帮帮我....!!!

我通过定义模式类和映射到 rdd 来尝试相同但得到错误

"ArrayIndexOutOfBoundsException :3"

【问题讨论】:

  • 您似乎在某些行中有 3 个元素,而在其他行中有 4 个元素。这应该是异常背后的原因。
  • 是的没错..!!!但我正在寻找出路
  • 也许这会有所帮助:stackoverflow.com/questions/50158696/…

标签: scala dataframe apache-spark rdd


【解决方案1】:

如果您将列视为字符串,您可以使用以下内容创建:

import org.apache.spark.sql.Row

val rdd : RDD[Row] = ???

val df = spark.createDataFrame(rdd, StructType(Seq(
  StructField("userId", StringType, false),
  StructField("EntityId", StringType, false),
  StructField("WebSessionId", StringType, false),
  StructField("ProductId", StringType, true))))

请注意,您必须将您的 RDD“映射”到 RDD[Row] 以便编译器允许使用“createDataFrame”方法。对于缺少的字段,您可以在 DataFrame Schema 中将列声明为可为空。

在您的示例中,您使用的是 RDD 方法 spark.sparkContext.textFile()。此方法返回一个 RDD[String],这意味着您的 RDD 的每个元素都是一行。但是,你需要一个 RDD[Row]。所以你需要用逗号分割你的字符串,比如:

val list = 
 List("545456,5615615,DIKFH6545614561456,PR5454564656445454",
   "875643,5485254,JHDSFJD543514KJKJ4", 
   "545456,5615615,DIKFH6545614561456,PR5454564656445454", 
   "545456,5615615,DIKFH6545614561456,PR5454564656445454", 
   "545456,5615615,DIKFH6545614561456,PR54545DSKJD541054", 
   "264264,3254564,MNXZCBMNABC5645SAD,PR5142545564542515", 
"732543,8765984,UJHSG4240323545144","564574,6276832,KJDXSGFJFS2545DSAS")


val FilterReadClicks = spark.sparkContext.parallelize(list)

val rows: RDD[Row] = FilterReadClicks.map(line => line.split(",")).map { arr =>
  val array = Row.fromSeq(arr.foldLeft(List[Any]())((a, b) => b :: a))
  if(array.length == 4) 
    array
  else Row.fromSeq(array.toSeq.:+(""))
}

rows.foreach(el => println(el.toSeq))

val df = spark.createDataFrame(rows, StructType(Seq(
  StructField("userId", StringType, false),
  StructField("EntityId", StringType, false),
  StructField("WebSessionId", StringType, false),
  StructField("ProductId", StringType, true))))

df.show()

+------------------+------------------+------------+---------+
|            userId|          EntityId|WebSessionId|ProductId|
+------------------+------------------+------------+---------+
|PR5454564656445454|DIKFH6545614561456|     5615615|   545456|
|JHDSFJD543514KJKJ4|           5485254|      875643|         |
|PR5454564656445454|DIKFH6545614561456|     5615615|   545456|
|PR5454564656445454|DIKFH6545614561456|     5615615|   545456|
|PR54545DSKJD541054|DIKFH6545614561456|     5615615|   545456|
|PR5142545564542515|MNXZCBMNABC5645SAD|     3254564|   264264|
|UJHSG4240323545144|           8765984|      732543|         |
|KJDXSGFJFS2545DSAS|           6276832|      564574|         |
+------------------+------------------+------------+---------+

使用 rows rdd,您将能够创建数据框。

【讨论】:

  • 嗨,它给出错误作为错误:重载方法值 createDataFrame 和替代:
  • 编辑您的问题并添加您的 RDD 代码以查看发生了什么。
  • val ReadClicks = c.textFile(FlumePath) \\这里的水槽路径包含多个数据源 val FilterReadClicks = ReadClicks.filter(x => ((!x.isEmpty) && (x != null ) && (x.lenght >3))) \\现在我正在尝试将 RDD 覆盖到数据帧中 val df = spark.createDataframe(FilterReadClicks, StructType(Seq(StructField("userId", StringType, false),StructField( "EntityId", StringType, false),StructField("WebSessionId", StringType, false),StructField("ProductId", StringType, true))))
  • 如何创建 RDD[Row] 或者有什么办法可以将 RDD[String] 转换为 RDD[Row]
  • 感谢您的更新...但似乎没有运气..!!我添加了建议的代码
猜你喜欢
  • 1970-01-01
  • 2016-05-29
  • 2017-06-13
  • 2021-06-29
  • 2023-03-13
  • 2018-03-05
  • 2017-10-30
  • 2017-06-02
  • 2018-09-14
相关资源
最近更新 更多