【问题标题】:How to convert csv into parquet file inside of HDFS如何将 csv 转换为 HDFS 内的镶木地板文件
【发布时间】:2020-08-13 02:36:45
【问题描述】:

我是Big Data 的新手,所以Hadoophdfs 现在对我来说有点消失了,所以我寻求帮助。 现在我有 4 个csv 格式的文件,它们位于HDFS 集群中,我应该使用PythonPARQUET 格式制作它们的4 个副本,但我不知道如何制作它。 我希望你能帮助我解决这个不难的问题。

【问题讨论】:

  • 欢迎来到 StackOverflow!请添加您迄今为止尝试过的内容以及您遇到的具体问题!我们不是来解决您的任务,而是解决一个特定问题!

标签: python csv hadoop hdfs parquet


【解决方案1】:

我把你的例子放在Scala 代码中,但是在Python 中做几乎是一样的。

我也放了一些 cmets 和一些解释

import org.apache.log4j.{Level, Logger}
import org.apache.spark.sql.SparkSession

object ReadCsv {
  val spark = SparkSession
    .builder()
    .appName("ReadCsv")
    .master("local[*]")
    .config("spark.sql.shuffle.partitions","4") //Change to a more reasonable default number of partitions for our data
    .config("spark.app.id","ReadCsv") // To silence Metrics warning
    .getOrCreate()

  val sqlContext = spark.sqlContext

  def main(args: Array[String]): Unit = {

    Logger.getRootLogger.setLevel(Level.ERROR)

    try {

      val df = sqlContext
        .read
        .csv("/path/directory_to_csv_files/") // Here we read the .csv files
        .cache()
      
      df.repartition(4) // we get four files
          .write
          .parquet("/path/directory_to_parquet_files/") // output format file.parquet.snappy by default
      // if we want parquet uncompressed before write we have to do:
      // sqlContext.setConf("spark.sql.parquet.compression.codec", "uncompressed")

      // To have the opportunity to view the web console of Spark: http://localhost:4040/
      println("Type whatever to the console to exit......")
      scala.io.StdIn.readLine()
    } finally {
      spark.stop()
      println("SparkSession stopped")
    }
  }
}

【讨论】:

    猜你喜欢
    • 2014-11-25
    • 2017-01-18
    • 2018-11-09
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-01-04
    • 2019-04-23
    • 1970-01-01
    相关资源
    最近更新 更多