【问题标题】:How to do JOIN on streamed data from kafka on spark streaming如何在火花流上加入来自 kafka 的流数据
【发布时间】: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


【解决方案1】:

你需要变换

需要这样的东西:

val dataset: RDD[String, String] = ... // From Hive
val windowedStream = stream.window(Seconds(20))... // From dStream
val joinedStream = windowedStream.transform { rdd => rdd.join(dataset) }

来自手册:

变换操作(以及它的变体,如 transformWith) 允许在 DStream 上应用任意 RDD-to-RDD 函数。它 可用于应用任何未公开的 RDD 操作 DStream API。例如,加入每个批次的功能 具有另一个数据集的数据流不直接暴露在 DStream API。但是,您可以轻松地使用转换来执行此操作。这个 实现了非常强大的可能性。

可以在此处找到一个示例: How to join a DStream with a non-stream file?

以下指南有帮助:https://spark.apache.org/docs/2.2.0/streaming-programming-guide.html

【讨论】:

    猜你喜欢
    • 2019-10-11
    • 1970-01-01
    • 2016-08-21
    • 2020-03-10
    • 2017-04-27
    • 2019-07-07
    • 2016-12-24
    • 2015-10-13
    • 2021-05-05
    相关资源
    最近更新 更多