【问题标题】:Transform DStream RDD using external data使用外部数据转换 DStream RDD
【发布时间】:2017-08-25 16:29:28
【问题描述】:

我们正在开发一个 Spark 流 ETL 应用程序,该应用程序将从 Kafka 获取数据,应用必要的转换并将数据加载到 MongoDB 中。从 Kafka 接收到的数据是 JSON 格式的。根据从 MongoDB 获取的查找数据,将转换应用于 RDD 的每个元素(JSON 字符串)。由于查找数据发生变化,我需要为每个批处理间隔获取它。使用 SqlContext.read 从 MongoDB 读取查找数据。我无法在 DStream.transform 中使用 SqlContext.read,因为 SqlContext 不可序列化,因此我无法将其广播到工作节点。现在我尝试使用 DStream.foreachRDD 在其中从 MongoDB 获取数据并将查找数据广播给工作人员。 RDD 元素上的所有转换都在 rdd.map 闭包内执行,该闭包利用广播数据并执行转换并返回 RDD。然后将 RDD 转换为数据帧并写入 MongoDB。目前,此应用程序运行速度很慢。

PS:如果我将获取查找数据的部分代码移出 DStream.foreachRDD 并添加 DStream.transform 以应用转换,并让 DStream.foreachRDD 仅将数据插入 MongoDB,性能非常好。但是使用这种方法,查找数据不会针对每个批次间隔进行更新。

我正在寻求帮助以了解这是否是一种好方法,并且我正在寻找一些指导来提高性能。

以下是伪代码

package com.testing


object Pseudo_Code {
  def main(args: Array[String]) {

    val sparkConf = new SparkConf().setAppName("Pseudo_Code")
      .setMaster("local[4]")

    val sc = new SparkContext(sparkConf)
    sc.setLogLevel("ERROR")

    val sqlContext = new SQLContext(sc)
    val ssc = new StreamingContext(sc, Seconds(1))

    val mongoIP = "127.0.0.1:27017"

    val DBConnectionURI = "mongodb://" + mongoIP + "/" + "DBName"

    val bootstrap_server_config = "127.0.0.100:9092"
    val zkQuorum = "127.0.0.101:2181"

    val group = "streaming"

    val TopicMap = Map("sampleTopic" -> 1)


    val KafkaDStream = KafkaUtils.createStream(ssc, zkQuorum,  group,  TopicMap).map(_._2)

     KafkaDStream.foreachRDD{rdd => 
       if (rdd.count() > 0) {

       //This lookup data has information required to perform transformation
       //This information keeps changing, so the data should be fetched for every batch interval

       val lookup1 = sqlContext.read.format("com.mongodb.spark.sql.DefaultSource")
        .option("spark.mongodb.input.uri", DBConnectionURI)
        .option("spark.mongodb.input.collection", "lookupCollection1")
        .load()

      val broadcastLkp1 = sc.broadcast(lookup1)

      val new_rdd = rdd.map{elem => 
      val json: JValue = parse(elem)

      //Foreach element in rdd, there are some values that should be looked up from the broadcasted lookup data
      //"value" extracted from json
      val param1 = broadcastLkp1.value.filter(broadcastLkp1.value("key")==="value").select("param1").distinct()
      val param1ReplaceNull = if(param1.count() == 0){
                                  "constant"
                                }
                                else{
                                  param1.head().getString(0)
                                }
      //create a new JSON with a different structure
      val new_JSON = """"""

      compact(render(new_JSON))
     }

     val targetSchema = new StructType(Array(StructField("key1",StringType,true)
                                                  ,StructField("key2",TimestampType,true)))
     val transformedDf = sqlContext.read.schema(targetSchema).json(new_rdd)


     transformedDf.write
          .option("spark.mongodb.output.uri",DBConnectionURI)
          .option("collection", "tagetCollectionName")
          .mode("append").format("com.mongodb.spark.sql").save()
       }
   }

    // Run the streaming job
    ssc.start()
    ssc.awaitTermination()
  }




}

【问题讨论】:

  • 你有一个有趣的问题。需要考虑的事项:您可以从您的 kafka 和您的 mongoDB 流式传输吗?如果是这种情况,那么您可以同时处理两个 DStream。
  • @MichelLemay 你有关于如何从 mongoDB 流式传输的示例。我可以试一试。现在,我可以按照stackoverflow.com/questions/37638519/… 中提供的一些说明向前迈进一点。我创建了一个 DStream.foreachRDD,我在其中重新加载查找数据,然后是 DStream.transform,其中使用了查找数据并返回一个新的 RDD,然后是另一个 foreachRDD 将数据插入到 mongoDB。这可行,但性能很差。
  • 您是否尝试过 dataframe API 中可用的 from_json 函数来进行 json 转换?您可以尝试结构化流式传输(如果您的驱动程序支持 .writeStream)。 val msgSchema = Encoders.product[Message].schema val ds = df .select(from_json($"value".cast("string"), msgSchema).as[Message])
  • @sgireddy 你能给我一个工作的例子吗?
  • 我之前的评论只是一个想法。不知道有没有这样的。但是,我认为这可以使用自定义接收器 api 来完成:spark.apache.org/docs/latest/streaming-custom-receivers.html

标签: scala apache-spark spark-streaming


【解决方案1】:

经过研究,有效的解决方案是在广播数据帧被工作人员读取后对其进行缓存。以下是我为提高性能而必须做的代码更改。

val new_rdd = rdd.map{elem => 
      val json: JValue = parse(elem)

      //Foreach element in rdd, there are some values that should be looked up from the broadcasted lookup data
      //"value" extracted from json
      val lkp_bd = broadcastLkp1.value
      lkp_bd.cache()
      val param1 = lkp_bd.filter(broadcastLkp1.value("key")==="value").select("param1").distinct()
      val param1ReplaceNull = if(param1.count() == 0){
                                  "constant"
                                }
                                else{
                                  param1.head().getString(0)
                                }
      //create a new JSON with a different structure
      val new_JSON = """"""

      compact(render(new_JSON))
     }

附带说明,这种方法在集群上运行时会出现问题。访问广播的数据帧时抛出空指针异常。我创建了另一个线程。 Spark Streaming - null pointer exception while accessing broadcast variable

【讨论】:

  • 这可能是因为变量评估发生在广播变量可用性之前。您可能希望将变量标记为惰性,以便它们等到第一次使用。另一种选择是通过直接从执行程序访问您需要的内容来完全避免广播变量。对象 GetMetaData { @transient 惰性 val metaData = getData; def getData = ... } 关键是将您的对象标记为瞬态,以便 spark 序列化程序忽略它并将其包装在对象中。
猜你喜欢
  • 2015-09-02
  • 2017-06-29
  • 2017-04-05
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2015-06-08
  • 2016-10-31
  • 2018-09-08
相关资源
最近更新 更多