【发布时间】: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