【发布时间】:2019-02-03 04:43:02
【问题描述】:
我是 spark 流媒体的新手。我正在尝试做一些关于从 kafka 获取数据并加入 hive 表的练习。我不确定如何在 spark 流(不是结构化流)中进行 JOIN。这是我的代码
val ssc = new StreamingContext("local[*]", "KafkaExample", Seconds(1))
val kafkaParams = Map[String, Object](
"bootstrap.servers" -> "dofff2.dl.uk.feefr.com:8002",
"security.protocol" -> "SASL_PLAINTEXT",
"key.deserializer" -> classOf[StringDeserializer],
"value.deserializer" -> classOf[StringDeserializer],
"group.id" -> "1",
"auto.offset.reset" -> "latest",
"enable.auto.commit" -> (false: java.lang.Boolean)
)
val topics = Array("csvstream")
val stream = KafkaUtils.createDirectStream[String, String](
ssc,
PreferConsistent,
Subscribe[String, String](topics, kafkaParams)
)
val strmk = stream.map(record => (record.value,record.timestamp))
现在我想加入 hive 中的一张表。在火花结构化流中,我可以直接调用 spark.table("table nanme") 并进行一些连接,但是在火花流中我该怎么做,因为它的一切都基于 RDD。有人可以帮我吗?
【问题讨论】:
-
走得更远……?
-
没有。我真的不确定如何将它用作数据框..
-
我给出的例子不清楚,参考文献?
-
需要 1 帮助.. 拆分我的值时如何从 kafka 添加时间戳?
-
val rdd1 = strmk.map(line => line.split(',')).map(s => (s(0).toString, s(1).toString,s( 2).toString,s(3).toString,s(4).toString, s(5).toString,s(6).toString,s(7).toString)))
标签: apache-spark spark-streaming