【问题标题】:Spark DataFrame: How to specify schema when writing as AvroSpark DataFrame:编写为 Avro 时如何指定架构
【发布时间】:2018-02-21 00:35:09
【问题描述】:

我想使用提供的 Avro 架构而不是 Spark 的自动生成架构来编写 Avro 格式的 DataFrame。如何告诉 Spark 在写入时使用我的自定义架构?

【问题讨论】:

    标签: apache-spark spark-dataframe spark-avro


    【解决方案1】:

    https://github.com/databricks/spark-avro/pull/222/ 中应用补丁后,我能够在写入时指定一个架构,如下所示:

    df.write.option("forceSchema", myCustomSchemaString).avro("/path/to/outputDir")
    

    【讨论】:

      【解决方案2】:

      希望下面的方法有所帮助。

      import org.apache.spark.sql.types._
      
      val schema = StructType( StructField("title", StringType, true) ::StructField("averageRating", DoubleType, false) ::StructField("numVotes", IntegerType, false) :: Nil)
      
      titleMappedDF.write.option("avroSchema", schema.toString).avro("/home/cloudera/workspace/movies/avrowithschema")
      
      Example:
      
      
      Download data from below site. https://datasets.imdbws.com/
      Download the movies data title.ratings.tsv.gz
      Copy to below location. /home/cloudera/workspace/movies/title.ratings.tsv.gz
      
      
      Start Spark-shell and type below command.
      
      import org.apache.spark.sql.SQLContext
      val sqlContext = new SQLContext(sc)
      val title = sqlContext.read.text("file:///home/cloudera/Downloads/movies/title.ratings.tsv.gz")
      scala> title.limit(5).show
      +--------------------+
      |               value|
      +--------------------+
      |tconst averageRat...|
      |  tt0000001    5.8 1350|
      |   tt0000002   6.5 157|
      |   tt0000003   6.6 933|
      |    tt0000004  6.4 93|
      +--------------------+
      
      val titlerdd = title.rdd
      
      case class Title(titleId:String, averageRating:Float, numVotes:Int)
      
      val titlefirst = titlerdd.first
      val titleMapped = titlerdd.filter(e=> e!=titlefirst).map(e=> {
         val rowStr = e.getString(0)
         val splitted = rowStr.split("\t")
         val titleId = splitted(0).trim
         val averageRating = scala.util.Try(splitted(1).trim.toFloat) getOrElse(0.0f)
         val numVotes = scala.util.Try(splitted(2).trim.toInt) getOrElse(0)
         Title(titleId, averageRating, numVotes)
      })
      
      val titleMappedDF =  titleMapped.toDF
      
      scala> titleMappedDF.limit(2).show
      +---------+-------------+--------+
      |  titleId|averageRating|numVotes|
      +---------+-------------+--------+
      |tt0000001|          5.8|    1350|
      |tt0000002|          6.5|     157|
      +---------+-------------+--------+
      
      
      import org.apache.spark.sql.types._
      
      val schema = StructType( StructField("title", StringType, true) ::StructField("averageRating", DoubleType, false) ::StructField("numVotes", IntegerType, false) :: Nil)
      
      titleMappedDF.write.option("avroSchema", schema.toString).avro("/home/cloudera/workspace/movies/avrowithschema")
      

      【讨论】:

      猜你喜欢
      • 2016-11-10
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2017-09-09
      • 1970-01-01
      • 1970-01-01
      • 2021-05-02
      • 2018-06-05
      相关资源
      最近更新 更多