【问题标题】:error mismatching in scalascala中的错误不匹配
【发布时间】:2015-11-16 14:40:35
【问题描述】:

我正在尝试返回 RDD[(String,String,String)],但使用 flatMap 无法做到这一点。我试过(tweetId, tweetBody, gender)(tweetId, tweetBody, gender),但它给了我一个类型不匹配的错误你能指导我知道如何从flatMap返回RDD[(String, String, String)]

override def transform(sqlContext: SQLContext, rdd: RDD[Array[Byte]], config: UserTransformConfig, logger: PhaseLogger): DataFrame = {
    val idColumnName = config.getConfigString("column_name").getOrElse("id")
    val bodyColumnName = config.getConfigString("column_name").getOrElse("body")
    val genderColumnName = config.getConfigString("column_name").getOrElse("gender")

    // convert each input element to a JsonValue
    val jsonRDD = rdd.map(r => byteUtils.bytesToUTF8String(r))

    val hashtagsRDD: RDD[(String,String, String)] = jsonRDD.mapPartitions(r => {
      // register jackson mapper (this needs to be instantiated per partition
      // since it is not serializable)
      val mapper = new ObjectMapper()
      mapper.registerModule(DefaultScalaModule)

      r.flatMap(tweet => tweet match {
        case _ :: tweet =>
        val rootNode = mapper.readTree(tweet)
        val tweetId = rootNode.path("id").asText.split(":")(2)
        val tweetBody = rootNode.path("body").asText
        val tweetVector =  new HashingTF().transform(tweetBody.split(" "))
        val result =genderModel.predict(tweetVector)
        val gender = if(result == 1.0){"Male"}else{"Female"}

        (tweetId, tweetBody, gender)
       // Array(1).map(x => (tweetId, tweetBody, gender))

      })

    })

    val rowRDD: RDD[Row] = hashtagsRDD.map(x => Row(x._1,x._2,x._3))
    val schema = StructType(Array(StructField(idColumnName,StringType, true),StructField(bodyColumnName, StringType, true),StructField(genderColumnName,StringType, true)))
    sqlContext.createDataFrame(rowRDD, schema)
  }
}

【问题讨论】:

  • 请多描述一下您的问题
  • 我正在尝试使用 RDD[String,String,String] 返回,但使用平面地图无法做到这一点。我试过 (tweetId, tweetBody, gender) 和 {tweetId, tweetBody, gender} 但它给了我类型不匹配的错误你能指导我知道如何从平面图中返回 RDD[(String, String, String)] 跨度>
  • 请在您的问题中添加相应信息,提供错误文本并修复格式
  • 已编辑,抱歉

标签: scala apache-spark singlestore


【解决方案1】:

尝试使用map 而不是flatMap。 当参数函数的结果类型为集合或RDD时使用flatMap

即当当前集合的每个元素都映射到零个或多个元素时,使用flatMap。 当当前集合的每个元素都映射到一个元素时,使用map

mapA => B 将符号 Afunctorial types 中的符号 B 交换,即将 RDD[A] 转换为 RDD[B]

flatMap 可以读作 map 然后 flattenmonadic types。例如。你有和RDD[A] 和参数函数是A => RDD[B] 类型的简单map 的结果将是RDD[RDD[B]] 并且这对出现可以通过flatten 简化为RDD[B]

这里是编译成功的例​​子。

import com.fasterxml.jackson.databind.ObjectMapper
import com.fasterxml.jackson.module.scala.DefaultScalaModule
import org.apache.spark.rdd.RDD
import org.apache.spark.sql._
import org.apache.spark.sql.types.{StringType, StructField, StructType}

class UserTransformConfig {
  def getConfigString(name: String): Option[String] = ???
}

class PhaseLogger
object byteUtils {
  def bytesToUTF8String(r: Array[Byte]): String = ???
}

class HashingTF {
  def transform(strs: Array[String]): Array[Double] = ???
}

object genderModel {
  def predict(v: Array[Double]): Double = ???
}

def transform(sqlContext: SQLContext, rdd: RDD[Array[Byte]], config: UserTransformConfig, logger: PhaseLogger): DataFrame = {
  val idColumnName = config.getConfigString("column_name").getOrElse("id")
  val bodyColumnName = config.getConfigString("column_name").getOrElse("body")
  val genderColumnName = config.getConfigString("column_name").getOrElse("gender")

  // convert each input element to a JsonValue
  val jsonRDD = rdd.map(r => byteUtils.bytesToUTF8String(r))

  val hashtagsRDD: RDD[(String, String, String)] = jsonRDD.mapPartitions(r => {
    // register jackson mapper (this needs to be instantiated per partition
    // since it is not serializable)
    val mapper = new ObjectMapper
    mapper.registerModule(DefaultScalaModule)

    r.map { tweet =>
      val rootNode = mapper.readTree(tweet)
      val tweetId = rootNode.path("id").asText.split(":")(2)
      val tweetBody = rootNode.path("body").asText
      val tweetVector = new HashingTF().transform(tweetBody.split(" "))
      val result = genderModel.predict(tweetVector)
      val gender = if (result == 1.0) {"Male"} else {"Female"}

      (tweetId, tweetBody, gender)

    }
  })

  val rowRDD: RDD[Row] = hashtagsRDD.map(x => Row(x._1, x._2, x._3))
  val schema = StructType(Array(StructField(idColumnName, StringType, true), StructField(bodyColumnName, StringType, true), StructField(genderColumnName, StringType, true)))
  sqlContext.createDataFrame(rowRDD, schema)
}

请注意我应该从我的想象中带来多少,因为您没有提供minimum example。一般来说,这样的问题不值得回答

【讨论】:

  • 但是上面代码的结果是RDD[(string,string,strin)]
  • 感谢您为回答我的问题所做的努力。但我需要我们
  • Flatmap( tweet=> (tweetvody,tweetid,gender)) 那怎么办
猜你喜欢
  • 1970-01-01
  • 2012-02-16
  • 1970-01-01
  • 1970-01-01
  • 2022-11-01
  • 2012-06-26
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多