【问题标题】:storing the string from kafkaStream into a variable for processing将 kafkaStream 中的字符串存储到变量中进行处理
【发布时间】:2018-07-20 03:38:29
【问题描述】:

我需要从 Kafka 生产者那里获取消息,并从消息中找到包含 % 的单词并为不同的 % 值生成消息。最后我需要将它发送到 ElasticSearch。

我可以使用kafkaStream.print() 在控制台中查看值,但我需要处理字符串以匹配所需的关键字并生成消息。

我的代码:

package rnd

import org.apache.spark.SparkConf
import kafka.serializer.StringDecoder
import org.apache.spark.sql.SQLContext
import org.apache.spark.streaming.dstream.DStream
import org.apache.spark.streaming.kafka.KafkaUtils
import org.apache.spark.streaming.{Minutes, Seconds, StreamingContext}
import org.apache.spark.{SparkConf, SparkContext}

object WordFind {
  def main(args: Array[String]) {
    val conf = new SparkConf().setMaster("local").setAppName("KafkaReceiver")
    val checkpointDir = "/usr/local/kafka/kafka_2.11-0.11.0.2/checkpoint/"

    import org.apache.spark.streaming.StreamingContext
    import org.apache.spark.streaming.Seconds

    val batchIntervalSeconds = 2
    val ssc = new StreamingContext(conf, Seconds(10))

    import org.apache.spark.streaming.kafka.KafkaUtils

    val kafkaStream = KafkaUtils.createStream(ssc, "localhost:2181", "spark-streaming-consumer-group", Map("wordcounttopic" -> 5))

    val s = kafkaStream.print()
    println(" the words are: " + s)
    ssc.remember(Minutes(1))
    ssc.checkpoint(checkpointDir)
    ssc
    ssc.start()
    ssc.awaitTerminationOrTimeout(batchIntervalSeconds * 5 * 1000)
  }
}

如果我通过 Lafka 生产者传递“使用率为 75%”,我应该在 ElasticSearch 中生成一条消息“将 ram 增加 25%”。

我得到的输出是:

18/02/09 16:38:27 INFO BlockManagerMasterEndpoint: Registering block manager localhost:37879 with 2.4 GB RAM, BlockManagerId(driver, localhost, 37879)
18/02/09 16:38:27 INFO BlockManagerMaster: Registered BlockManager
18/02/09 16:38:27 WARN StreamingContext: spark.master should be set as local[n], n > 1 in local mode if you have receivers to get data, otherwise Spark jobs will not get resources to process the received data.
 ***the words are: ()***

我想要我传递的字符串来代替 's' 中的 ()。

【问题讨论】:

  • kafkaStream.print 不返回 String 进行打印,而是执行实际打印。您要打印的是(),这是Unit 类型的单例值(如果您熟悉的话,您可以在Java 中将其称为void。
  • 那么我如何将值作为字符串获取?我需要对来自 kafka 的消息进行一些模式匹配。如何将值转换为字符串以进行模式匹配?

标签: scala apache-spark elasticsearch apache-kafka spark-streaming


【解决方案1】:

val kafkaStream 是 RecieverInputDStream[(String, String)],其中数据是 (kafkaMetaData, kafkaMessage) 有关详细信息,请参阅 [https://github.com/apache/spark/blob/f830bb9170f6b853565d9dd30ca7418b93a54fe3/external/kafka-0-8/src/main/scala/org/apache/spark/streaming/kafka/KafkaInputDStream.scala#L135]。

我们需要提取元组的第二个并进行模式匹配(即过滤 RecieverInputDStream 找到包含 % 的单词),然后使用 map 生成输出(即不同 % 值的消息)。正如@stefanobaghino 所提到的, print() 函数只是将输出打印到控制台并且不返回任何记录字符串。

例如:

import org.apache.spark.streaming.dstream.ReceiverInputDStream
val kafkaStream: ReceiverInputDStream[(String, String)] = KafkaUtils.createStream(sparkStreamingContext, "localhost:2181",
  "spark-streaming-consumer-group", Map("wordcounttopic" -> 5))

import org.apache.spark.streaming.dstream.DStream
val filteredStream: DStream[(String, String)] = kafkaStream
  .filter(record => record._2.contains("%")) // TODO : pattern matching here

val outputDStream: DStream[String] = filteredStream
  .map(record => record._2.toUpperCase()) // just assuming some operation
outputDStream.print()

使用 outputDStream 写入 ElasticSearch。希望这会有所帮助。

【讨论】:

  • 非常感谢。会试试这个,让你知道。
  • 您好 Pavithran,非常感谢您的帮助。由于我对 kafka 很陌生,您能帮我了解如何为不同的 % 值生成不同的消息吗?例如如果值是
  • 当然@sc,您能否创建一个新问题并对其进行详细描述...因为,在stackoverflow中,每个帖子必须只有一个问题要回答,新帖子将是向所有人开放以回答您的疑问。请提供数据样本,因为百分比计算没有给出清晰的图像。另请阅读:[stackoverflow.com/help/mcve]
  • 嗨 Pavithran 我在下面的链接中发布了一个与我的场景不同的问题:stackoverflow.com/questions/48747604/… 你能帮我写代码吗?如果对场景有任何疑问,请告诉我。谢谢
  • 好的,当然。我也是stackoverflow的新手。尝试如何将其标记为已回答。请帮我解决其他问题
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2014-07-27
  • 2015-07-20
  • 2016-02-18
  • 2020-09-20
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多