【问题标题】:why the data has changed after convert to parquet format testing by union two dataframe?为什么联合两个数据框转换为镶木地板格式测试后数据发生了变化?
【发布时间】:2017-09-14 06:18:40
【问题描述】:

我写了一个函数来操作一个 csv 文件,把它转换成 parquet 格式。 我想知道如何确保数据相同,不会丢失或添加。 所以我为它写了一个测试。但事实证明它们并不相同: 我的逻辑是:

1) 将 csv 制作成数据框 A

2) 并将数据框 A 设置为 parquet 格式,保存到目录。

3) 将 parquet 文件读取为新的数据帧 B

4) 然后 A.union(B).

5)计算 A 和 B 以及 A.union(B)。

如果这三个相同,那么我可以得出它们是相同数据的结论。 但我得到了第三个不同的。

def doJob(sc: SparkContext, data: RDD[String]): DataFrame = {
    logInfo("Extracting omniture data")
    val result = data
      .filter(_.contains("PAGE."))
      .filter(_.contains(".PACKAGE"))
    val sqlsqlContext = new SQLContext(sc)

//just ignore above codes...

    val packagesCsvDF = sqlsqlContext.load("com.databricks.spark.csv", Map("path" -> "file:///D:/test/testsample.csv", "header" -> "true"))
    val sqlContext = new org.apache.spark.sql.hive.HiveContext(sc)
    import sqlContext.implicits._

//
//    // we should have some additional filter here
//    val mydf = packagesDF.groupBy($"page_url").agg(last($"pagename"),last($"prop46"),last($"prop56"),last($"post_evar34"))
//    logInfo("show mydf")
//    mydf.show()

    //TODO
    // save files
    logInfo("Saving omniture packages data to S3")
    if (true) {
      packagesCsvDF
        .repartition(sc.defaultParallelism, col("pagename"))
        .write
        .mode(SaveMode.Append)
        .partitionBy("pagename")
        .parquet("file:///D:/test/parquet")
      logInfo("packagesDF")

    }

    packagesCsvDF//Is this packagesCsvDF have not been changed yet??????
  }

测试:

object ParquetDataTestsSpec {
  def main (args: Array[String] ): Unit = {
    val sc = new SparkContext(new SparkConf().setAppName("parquet data test Logs").setMaster("local"))


    val input = PackagesOmnitureMapReduceJob.formatToJson(sc.textFile("file:///D:/test/option.json", sc.defaultParallelism))
    val df = PackagesOmnitureMapReduceJob.doJob(sc, input)//call the function I want to test in "file:///D:/test/parquet"
    val sqlContext = new SQLContext(sc)

    val SourceCSVDF = sqlContext.load("com.databricks.spark.csv", Map("path" -> "file:///D:/test/testsample.csv", "header" -> "true"))// original 

    val parquetDataFrame = sqlContext.read.parquet("file:///D:/test/parquet") //get the new dataframe

    val dfCount = df.count()
    val SourceCSVDFcount = SourceCSVDF.count()
    val parquetDataCount = parquetDataFrame.count()

    val unionCount = parquetDataFrame.union(SourceCSVDF).count()
    println(dfCount,SourceCSVDFcount,parquetDataCount,unionCount)


  }
}

打印:

(200,200,200,400)

然后我尝试将所有数据帧解析为 json:

parquetDataFrame.write.json("file:///D:/test/parquetDataFrame")
SourceCSVDF.write.json("file:///D:/test/SourceCSVDF")
df.write.json("file:///D:/test/Desktop/df")

当我打开 json 文件时,我发现它们都是一样的。问题是否出在关键字 union 上?

【问题讨论】:

    标签: apache-spark dataframe parquet


    【解决方案1】:
    val unionalldis3 = parquetDataFrame.unionAll(SourceCSVDF).distinct().count()
    

    那就对了……

    但我很困惑。我认为 union() 是不同的 unionAll....

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-07-18
      • 2021-04-25
      • 1970-01-01
      • 1970-01-01
      • 2018-01-04
      • 1970-01-01
      相关资源
      最近更新 更多