【问题标题】:Spark: read csv file from s3 using scalaSpark:使用scala从s3读取csv文件
【发布时间】:2015-12-04 21:36:43
【问题描述】:

我正在编写一个 spark 作业,尝试使用 scala 读取文本文件,以下在我的本地计算机上运行良好。

  val myFile = "myLocalPath/myFile.csv"
  for (line <- Source.fromFile(myFile).getLines()) {
    val data = line.split(",")
    myHashMap.put(data(0), data(1).toDouble)
  }

然后我尝试让它在 AWS 上运行,我做了以下操作,但它似乎没有正确读取整个文件。在 s3 上读取此类文本文件的正确方法应该是什么?非常感谢!

val credentials = new BasicAWSCredentials("myKey", "mySecretKey");
val s3Client = new AmazonS3Client(credentials);
val s3Object = s3Client.getObject(new GetObjectRequest("myBucket", "myFile.csv"));

val reader = new BufferedReader(new InputStreamReader(s3Object.getObjectContent()));

var line = ""
while ((line = reader.readLine()) != null) {
      val data = line.split(",")
      myHashMap.put(data(0), data(1).toDouble)
      println(line);
}

【问题讨论】:

    标签: scala amazon-web-services amazon-s3 apache-spark


    【解决方案1】:

    我想我得到了它的工作方式如下:

        val s3Object= s3Client.getObject(new GetObjectRequest("myBucket", "myPath/myFile.csv"));
    
        val myData = Source.fromInputStream(s3Object.getObjectContent()).getLines()
        for (line <- myData) {
            val data = line.split(",")
            myMap.put(data(0), data(1).toDouble)
        }
    
        println(" my map : " + myMap.toString())
    

    【讨论】:

      【解决方案2】:

      使用sc.textFile("s3://myBucket/myFile.csv") 读取 csv 文件。这会给你一个 RDD[String]。把它放到地图中

      val myHashMap = data.collect
                          .map(line => {
                            val substrings = line.split(" ")
                            (substrings(0), substrings(1).toDouble)})
                          .toMap
      

      您可以使用sc.broadcast 广播您的地图,以便在您的所有工作节点上随时可用。

      (请注意,如果您愿意,当然也可以使用 Databricks “spark-csv” 包来读取 csv 文件。)

      【讨论】:

      • 我的实用程序函数需要 myHashMap。所以我的代码是这样的: output = input.map { t => myUtiltyFunction(myHashMap, t)} 是否可以避免每次都将 myHashMap 传递给 myUtiltiyFunction?有没有办法使用广播 myHashMap 并让 myUtitlityFunction 直接知道它?非常感谢!
      • 另外,我不想使用 sc.textFile("s3://myBucket/myFile.csv") 因为我想让代码通用,即使没有火花上下文。谢谢。
      • 你确实意识到,如果你让你的效用函数直接读取地图,并且像你描述的output = input.map { t =&gt; myUtiltyFunction(...)}那样使用效用函数,那么将为你输入的每一行读取和创建地图rdd。我真的不认为你想要那个。另一方面,如果您广播变量(使用sc.broadcast),您只需在驱动程序上读取和创建一次地图,然后您的所有工作人员都可以直接访问它。为什么不想将地图传递给实用程序函数?这对我来说似乎很奇怪。
      • 您确定地图是按输入 rdd 的单行而不是按任务创建的吗?我不想通过 HashMap 的原因是: 1. 代码的清洁度。 2. 我希望在其他一些场景中使用相同的代码,将输入数据读取到实用函数是微不足道的。
      • 如果你使用map而不是mapPartition,它会将效用函数应用到每一行,因此如果你的效用函数负责创建地图,它将为每一行完成。如果您使用mapPartitions,它只会为每个分区创建一次映射,但是(当然,取决于您的数据大小)这仍然很容易最终增加显着的开销(io 从不便宜)。 IMO,您应该专注于编写最适合并行处理 (Spark) 的代码,而不是关心代码的其他琐碎(非并行)用途。
      【解决方案3】:

      即使不使用SparkContext textfile 导入 amazons3 库,也可以实现这一点。使用下面的代码

      import org.apache.hadoop.fs.{FileSystem, Path}
      import org.apache.hadoop.conf.Configuration
      val s3Login = "s3://AccessKey:Securitykey@Externalbucket"
      val filePath = s3Login + "/Myfolder/myscv.csv"
      for (line <- sc.textFile(filePath).collect())
      {
          var data = line.split(",")
          var value1 = data(0)
          var value2 = data(1).toDouble
      }
      

      在上面的代码中,sc.textFile 将从您的文件中读取数据并存储在line RDD 中。然后它将带有, 的每一行拆分为循环内的不同RDD data。然后,您可以使用索引访问此 RDD 中的值。

      【讨论】:

      • 此代码返回错误“java.io.IOException: No FileSystem for scheme: s3”
      • 你能解释一下答案吗,我也收到了java.io.FileNotFoundException
      • Mitkash 和 ibaalf - 请分享您的代码供我调试。可能有一些错字,因为这对我来说非常有用
      猜你喜欢
      • 1970-01-01
      • 2018-04-26
      • 2021-04-27
      • 1970-01-01
      • 1970-01-01
      • 2022-01-13
      • 2019-03-25
      • 2021-12-15
      • 2018-10-18
      相关资源
      最近更新 更多