【问题标题】:Convert csv file to dataframe in Spark 1.5.2 without databricks在没有数据块的 Spark 1.5.2 中将 csv 文件转换为数据帧
【发布时间】:2023-04-10 02:37:02
【问题描述】:

我正在尝试使用 Scala 将 csv 文件转换为 Spark 1.5.2 中的数据帧,而不使用库数据块,因为它是一个社区项目,并且该库不可用。我的方法如下:

var inputPath  = "input.csv"
var text = sc.textFile(inputPath)
var rows = text.map(line => line.split(",").map(_.trim))
var header = rows.first()
var data = rows.filter(_(0) != header(0))
var df = sc.makeRDD(1 to data.count().toInt).map(i => (data.take(i).drop(i-1)(0)(0), data.take(i).drop(i-1)(0)(1), data.take(i).drop(i-1)(0)(2), data.take(i).drop(i-1)(0)(3), data.take(i).drop(i-1)(0)(4))).toDF(header(0), header(1), header(2), header(3), header(4))

这段代码虽然很乱,但运行时不会返回任何错误消息。问题是在尝试显示df中的数据以验证此方法的正确性并稍后尝试在df中进行一些查询时出现问题。执行df.show() 后得到的错误代码是SPARK-5063。我的问题是:

1)为什么不能打印df的内容?

2) 有没有其他更直接的方法可以在不使用库 databricks 的情况下将 csv 转换为 Spark 1.5.2 中的数据帧?

【问题讨论】:

  • “这是一个社区项目”——你是认真的吗?您知道 Databricks 是推动 Spark 开发的公司吗?你知道spark-csv插件已经合并到Spark 2.x核心库中了吗?
  • 问题是我没有机会改变它,因此我正在寻找不使用 Databricks 将 csv 解析为数据帧的替代方法。
  • "改变那个" -- 你是什么意思?您无法下载 JAR(一劳永逸)并将其附加到带有 --jars 的作业以及它的 commons-csv 依赖项?在 CDH 发行版中捆绑 Spark 对我来说效果很好(请注意,对于 Apache 构建,--jars 不适用于 CDH,我必须使用 spark.driver.extraClasspath 道具和明确的 sc.addJar() 作为解决方法)

标签: scala csv apache-spark spark-dataframe


【解决方案1】:

对于 spark 1.5.x 可以使用下面的代码 sn-p 将输入转换为 DF

val sqlContext = new org.apache.spark.sql.SQLContext(sc)
// this is used to implicitly convert an RDD to a DataFrame.
import sqlContext.implicits._

// Define the schema using a case class.
// Note: Case classes in Scala 2.10 can support only up to 22 fields. To work around this limit,
// you can use custom classes that implement the DataClass interface with 5 fields.
case class DataClass(id: Int, name: String, surname: String, bdate: String, address: String)

// Create an RDD of DataClass objects and register it as a table.
val peopleData = sc.textFile("input.csv").map(_.split(",")).map(p => DataClass(p(0).trim.toInt, p(1).trim, p(2).trim, p(3).trim, p(4).trim)).toDF()
peopleData.registerTempTable("dataTable")

val peopleDataFrame = sqlContext.sql("SELECT * from dataTable")

peopleDataFrame.show()

Spark 1.5

【讨论】:

    【解决方案2】:

    你可以这样创建:

    SparkSession spark = SparkSession
                    .builder()
                    .appName("RDDtoDF_Updated")
                    .master("local[2]")
                    .config("spark.some.config.option", "some-value")
                    .getOrCreate();
    
            StructType schema = DataTypes
                    .createStructType(new StructField[] {
                            DataTypes.createStructField("eid", DataTypes.IntegerType, false),
                            DataTypes.createStructField("eName", DataTypes.StringType, false),
                            DataTypes.createStructField("eAge", DataTypes.IntegerType, true),
                            DataTypes.createStructField("eDept", DataTypes.IntegerType, true),
                            DataTypes.createStructField("eSal", DataTypes.IntegerType, true),
                            DataTypes.createStructField("eGen", DataTypes.StringType,true)});
    
    
            String filepath = "F:/Hadoop/Data/EMPData.txt";
            JavaRDD<Row> empRDD = spark.read()
                    .textFile(filepath)
                    .javaRDD()
                    .map(line -> line.split("\\,"))
                    .map(r -> RowFactory.create(Integer.parseInt(r[0]), r[1].trim(),Integer.parseInt(r[2]),
                            Integer.parseInt(r[3]),Integer.parseInt(r[4]),r[5].trim() ));
    
    
            Dataset<Row> empDF = spark.createDataFrame(empRDD, schema);
            empDF.groupBy("edept").max("esal").show();
    

    【讨论】:

      【解决方案3】:

      在 Scala 中使用 Spark。

      import org.apache.spark.sql.Row
      import org.apache.spark.sql.types._
      
      var hiveCtx = new HiveContext(sc)
      var inputPath  = "input.csv"
      var text = sc.textFile(inputPath)
      var rows = text.map(line => line.split(",").map(_.trim)).map(a => Row.fromSeq(a))
      var header = rows.first()
      val schema = StructType(header.map(fieldName => StructField(fieldName.asInstanceOf[String],StringType,true)))
      
      val df = hiveCtx.createDataframe(rows,schema)
      

      这应该可行。

      但是对于创建数据框,建议您使用Spark-CSV

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2016-01-03
        • 2018-05-24
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2021-04-19
        • 2016-09-27
        相关资源
        最近更新 更多