【问题标题】:Extract data from a Spark RDD, and populate a tuple in scala从 Spark RDD 中提取数据,并在 scala 中填充一个元组
【发布时间】:2021-02-03 11:43:45
【问题描述】:

我在 Hadoop/Spark 框架之上使用 Scala。

其实我的数据是这样的:

RDD[(List[(String, Int)], Long)]

而且,这是该数据湖中前两行的示例:

(List(("COD_LOCALE_PROGETTO",0), ("CUP",1), ("OC_TITOLO_PROGETTO",2), ("OC_SINTESI_PROGETTO",3), ("OC_LINK",4), ("OC_COD_CICLO",5), ("OC_DESCR_CICLO",6), ("OC_COD_TEMA_SINTETICO",7), ("OC_TEMA_SINTETICO",8), ("COD_GRANDE_PROGETTO",9), ("DESCRIZIONE_GRANDE_PROGETTO",10)),0)

(List(("10CAPORTO-POZZUOLI 1",0), ("J86G08000450003",1), ("INTERVENTO C11 2° LOTTO ¿ 1° STRALCIO FUNZIONALE ¿COLLEGAMENTO TRA TANGENZIALE DI NAPOLI (VIA CAMPANA), RETE VIARIA COSTIERA E PORTO DI POZZUOLI""",2), ("INTERVENTO C11 2° LOTTO ¿ 1° STRALCIO FUNZIONALE ¿COLLEGAMENTO TRA TANGENZIALE DI NAPOLI (VIA CAMPANA), RETE VIARIA COSTIERA E PORTO DI POZZUOLI""",3), ("www.opencoesione.gov.it/progetti/10caporto-pozzuoli-1",4), (1,5), ("Ciclo di programmazione 2007-2013",6), ("07",7), ("Trasporti e infrastrutture a rete",8), (" ",9), (" ",10)),1)

在实际情况下,每行持续 194 列,我总共有超过 160 万条记录。

有了这个数据集,我想填充一个新的列表,类型为:

List[(String, Int, Int, Int)]

第一个“Int”是每行的每个字段(COD_LOCALE_PROGETTO,CUP...),第二个字段是每个字段的大小(19、3,...)第三个是位置每个字段,已经编码在变量中,就在字符串之后,最后一个“Int”是整个数据集中每一行的位置。

我试过这个脚本:

     | val Dimensione = item._1.size;
     | for(i <- 0 until Dimensione){
     | ComponentiOpenCoesione :+= (item._1(i)._1.replace("\"","").toString,
     | item._1(i)._1.replace("\"","").toString.size,
     | item._1(i)._2.toInt,
     | item._2.toLong)}
     | })

但它失败了,我称之为“ComponentiOpenCoesione”的元组列表没有填充。

最后,这个变量是这样定义的:

var ComponentiOpenCoesione : List[(String, Int, Int, Long)] = List();

有人可以帮助我吗?如何从 RDD 中提取数据并将其加载到列表中?

非常感谢。

【问题讨论】:

  • 你是如何加载这些数据的,是从表还是从文件中加载的?
  • 数据来自 Hadoop HDFS 分区文件系统
  • 你是怎么加载的,也可以放那个代码吗? & 还有一些解析前的示例文件数据?
  • 你能添加一个小例子,输入和预期输出吗?
  • 很好奇为什么您使用的是 RDD 而不是带有数组和长列的数据框

标签: scala apache-spark


【解决方案1】:

在 Scala 中,返回函数的最后一条语句。在这里,您的函数不会返回任何内容,因为它的最后一条语句是不返回任何内容的 for 循环。

要更正它,您只需将ComponentiOpenCoesione 作为您的最后一条语句。因此,如果您只是打算将您的RDD[(List[(String, Int)], Long)] 映射到RDD[List[(String, Int, Int, Long)]],您的代码应该是:

rdd.map(item => {
  var ComponentiOpenCoesione: List[(String, Int, Int, Long)] = List();
  val Dimensione = item._1.size;
  for (i <- 0 until Dimensione) {
    ComponentiOpenCoesione :+= (item._1(i)._1.replace("\"", "").toString,
      item._1(i)._1.replace("\"", "").toString.size,
      item._1(i)._2.toInt,
      item._2.toLong)
  }
  ComponentiOpenCoesione
})

您可以查看Return in Scala 问题的答案以了解值是如何在 scala 中返回的。

【讨论】:

    猜你喜欢
    • 2017-05-04
    • 2021-10-25
    • 1970-01-01
    • 1970-01-01
    • 2015-10-18
    • 1970-01-01
    • 2012-04-30
    • 2023-04-04
    • 1970-01-01
    相关资源
    最近更新 更多