【问题标题】:unable to read kafka topic data using spark无法使用 spark 读取 kafka 主题数据
【发布时间】:2020-09-18 04:23:13
【问题描述】:

我在我创建的名为"sampleTopic" 的主题之一中有如下数据

sid,Believer  

第一个参数是username,第二个参数是用户经常收听的song name。现在,我已经开始使用上面提到的主题名称zookeeperKafka serverproducer。我已经使用CMD 为该主题输入了上述数据。现在,我想阅读 spark 中的主题执行一些聚合,并将其写回流。以下是我的代码:

package com.sparkKafka
import org.apache.spark.SparkContext
import org.apache.spark.SparkConf
import org.apache.spark.sql.SparkSession
object SparkKafkaTopic {
  def main(args: Array[String]) {
    val spark = SparkSession.builder().appName("SparkKafka").master("local[*]").getOrCreate()
    println("hey")
    val df = spark
      .readStream
      .format("kafka")
      .option("kafka.bootstrap.servers", "localhost:9092")
      .option("subscribe", "sampleTopic1")
      .load()
    val query = df.writeStream
      .outputMode("append")
      .format("console")
      .start().awaitTermination()


  }
}

但是,当我执行上面的代码时,它给出了:

    +----+--------------------+------------+---------+------+--------------------+-------------+
| key|               value|       topic|partition|offset|           timestamp|timestampType|
+----+--------------------+------------+---------+------+--------------------+-------------+
|null|[73 69 64 64 68 6...|sampleTopic1|        0|     4|2020-05-31 12:12:...|            0|
+----+--------------------+------------+---------+------+--------------------+-------------+

下面也有无限循环消息

20/05/31 11:56:12 INFO Fetcher: [Consumer clientId=consumer-1, groupId=spark-kafka-source-0d6807b9-fcc9-4847-abeb-f0b81ab25187--264582860-driver-0] Resetting offset for partition sampleTopic1-0 to offset 4.
20/05/31 11:56:12 INFO Fetcher: [Consumer clientId=consumer-1, groupId=spark-kafka-source-0d6807b9-fcc9-4847-abeb-f0b81ab25187--264582860-driver-0] Resetting offset for partition sampleTopic1-0 to offset 4.
20/05/31 11:56:12 INFO Fetcher: [Consumer clientId=consumer-1, groupId=spark-kafka-source-0d6807b9-fcc9-4847-abeb-f0b81ab25187--264582860-driver-0] Resetting offset for partition sampleTopic1-0 to offset 4.
20/05/31 11:56:12 INFO Fetcher: [Consumer clientId=consumer-1, groupId=spark-kafka-source-0d6807b9-fcc9-4847-abeb-f0b81ab25187--264582860-driver-0] Resetting offset for partition sampleTopic1-0 to offset 4.
20/05/31 11:56:12 INFO Fetcher: [Consumer clientId=consumer-1, groupId=spark-kafka-source-0d6807b9-fcc9-4847-abeb-f0b81ab25187--264582860-driver-0] Resetting offset for partition sampleTopic1-0 to offset 4.
20/05/31 11:56:12 INFO Fetcher: [Consumer clientId=consumer-1, groupId=spark-kafka-source-0d6807b9-fcc9-4847-abeb-f0b81ab25187--264582860-driver-0] Resetting offset for partition sampleTopic1-0 to offset 4.

我需要如下输出:

根据 Srinivas 的建议修改后,我得到以下输出:

不确定这里到底出了什么问题。请指导我完成它。

【问题讨论】:

  • 感谢您的建议。建议的依赖项对我有用,但我在这里面临一个新问题。请完成问题
  • 你能在 kafka 中添加你的数据架构吗??
  • 怎么做?实际上,我对这种集成方式很陌生。架构将是两列的字符串。我修改了我的输出
  • 或者你能发布你的 kafka 消息吗??
  • 哪条消息?您要在主题中发布的数据?

标签: apache-spark hadoop apache-kafka streaming


【解决方案1】:

尝试将spark-sql-kafka 库添加到您的构建文件中。检查下面。

build.sbt

libraryDependencies += "org.apache.spark" %% "spark-sql-kafka-0-10" % "2.3.0"  
// Change to Your spark version 

pom.xml

<dependency>
    <groupId>org.apache.spark</groupId>
    <artifactId>spark-sql-kafka-0-10_2.11</artifactId>
    <version>2.3.0</version>    // Change to Your spark version
</dependency>

如下更改您的代码

    package com.sparkKafka
    import org.apache.spark.SparkContext
    import org.apache.spark.SparkConf
    import org.apache.spark.sql.SparkSession
    import org.apache.spark.sql.types._
    import org.apache.spark.sql.functions._
    case class KafkaMessage(key: String, value: String, topic: String, partition: Int, offset: Long, timestamp: String)

    object SparkKafkaTopic {

      def main(args: Array[String]) {
        //val spark = SparkSession.builder().appName("SparkKafka").master("local[*]").getOrCreate()
        println("hey")
        val spark = SparkSession.builder().appName("SparkKafka").master("local[*]").getOrCreate()
        import spark.implicits._
        val mySchema = StructType(Array(
          StructField("userName", StringType),
          StructField("songName", StringType)))
        val df = spark
          .readStream
          .format("kafka")
          .option("kafka.bootstrap.servers", "localhost:9092")
          .option("subscribe", "sampleTopic1")
          .load()

        val query = df
          .as[KafkaMessage]
          .select(split($"value", ",")(0).as("userName"),split($"value", ",")(1).as("songName"))
          .writeStream
          .outputMode("append")
          .format("console")
          .start()
          .awaitTermination()
      }
    }

     /*
        +------+--------+
        |userid|songname|
        +------+--------+
        |   sid|Believer|
        +------+--------+
       */

      }
    }

【讨论】:

  • 添加新评论并根据面临的新问题修改问题。请帮帮我。
  • 当我合并你的代码时,它给出了一个错误提示 unable to find encoder for type [KafkaMessage]
  • 它在 [kafkamessage] 行上给出以下错误def as[U](implicit evidence$2: Encoder[U]): Dataset[U] :: Experimental :: Returns a new Dataset where each record has been mapped on to the specified type. The method used to map columns depend on the type of U: ◾When U is a class, fields for the class will be mapped to columns of the same name(case sensitivity is determined by spark.sql.caseSensitive).
  • 是的。它有效,但无法为消息添加架构。以及如何为到达的输入消息添加一些窗口?检查我编辑的问题输出
  • 我已经用我得到的输出修改了问题
【解决方案2】:

spark-sql-kafka jar 丢失,它正在实现“kafka”数据源。

您可以使用配置选项添加 jar 或构建包含 spark-sql-kafka jar 的 fat jar。请使用相关版本的jar

val spark = SparkSession.builder()
  .appName("SparkKafka").master("local[*]")
  .config("spark.jars","/path/to/spark-sql-kafka-xxxxxx.jar")
  .getOrCreate()

【讨论】:

  • 添加了一条新评论,并根据面临的新问题修改了问题。请帮帮我。
猜你喜欢
  • 2023-03-19
  • 2021-12-21
  • 2021-09-12
  • 2018-01-09
  • 2021-06-11
  • 1970-01-01
  • 1970-01-01
  • 2019-11-04
  • 2023-03-18
相关资源
最近更新 更多