【问题标题】:Kafka Streaming Using Scala and Spark Error: Exception thrown in awaitResultKafka Streaming 使用 Scala 和 Spark 错误:在 awaitResult 中抛出异常
【发布时间】:2022-01-08 06:01:48
【问题描述】:

我还是 Kafka 的新手,我目前正在尝试生成和读取 .csv 文件,然后使用 IntelliJ IDEA(在 Windows 中运行)和 Kafka 将其流式传输到 Kafka 消费者(动物园管理员、经纪人和消费者)在 WSL 中运行),但我一直没有这样做。

这是我的 build.sbt:

name := "kafka_test"

version := "0.1"

scalaVersion := "2.13.6"

//libraryDependencies += "org.apache.spark" % "spark-streaming_2.11" % "2.2.0"
//libraryDependencies += "org.apache.spark" % "spark-streaming-kafka-0-8_2.11" % "2.1.0"
libraryDependencies += "org.apache.spark" %% "spark-core" % "3.2.0"
libraryDependencies += "org.apache.spark" %% "spark-streaming" % "3.2.0"
//libraryDependencies += "org.apache.spark" %% "spark-streaming-kafka" % "1.6.0"
//libraryDependencies += "org.apache.spark" % "spark-streaming-kafka-0-10_2.11" % "2.2.0"
libraryDependencies += "org.apache.spark" %% "spark-sql" % "3.2.0"
libraryDependencies += "org.apache.spark" %% "spark-hive" % "3.2.0"
libraryDependencies += "org.apache.spark" %% "spark-sql-kafka-0-10" % "3.2.0"

这是我的代码:

import org.apache.log4j.{Level, Logger}
import org.apache.spark.SparkConf
//import org.apache.spark.streaming.kafka.KafkaUtils
import org.apache.spark.streaming.{Seconds, StreamingContext}
import org.apache.spark.streaming._
//import org.apache.spark.streaming.kafka
import org.apache.spark.streaming.StreamingContext._

import java.beans.Statement
import java.sql.{Connection, DriverManager, SQLException}
import scala.io.StdIn
import org.apache.spark.{SparkConf, SparkContext}
import org.apache.spark.sql.{Encoder, Encoders, SparkSession, functions}

import scala.util.control.Breaks._
import org.apache.spark.storage.StorageLevel
import org.apache.spark.sql.functions.{col, count, countDistinct, desc, from_json, when}
import org.apache.spark.sql.types.StructType

import java.io._

import CustomImplicits._


object kafkar {

  def main(args: Array[String]): Unit = {

    Logger.getLogger("org").setLevel(Level.OFF)
    Logger.getLogger("akka").setLevel(Level.OFF)

    println("Program Started")

    val conf = new SparkConf().setMaster("local[4]").setAppName("kafkar")
    val ssc = new StreamingContext(conf, Seconds(2))

    //INITIATE SPARK SESSION//
    System.setProperty("hadoop.home.dir", "C:\\hadoop")
    val spark = SparkSession
      .builder
      .appName("Kafka Streaming")
      .config("spark.master", "local")
      .enableHiveSupport()
      .getOrCreate()
    println("Created Spark Session")
    spark.sparkContext.setLogLevel("ERROR")


    //my kafka topic name is 'mytest'
//    val kafkaStream = KafkaUtils.createStream(ssc, "localhost:2181","spark-streaming-consumer-group", Map("mytest" -> 5) )
//    kafkaStream.print()
//    ssc.start
//    ssc.awaitTermination()


    //GENERATES A .CSV FILE WITH REQUIRED SCHEMA
    val customer_names = List("SpaceX", "Blue Origin", "Orbital Sciences Corporation", "Boeing", "Northrop Grumman Innovation Systems", "Sierra Nevada Corporation", "Scaled Composites", "The Spaceship Company", "NASA", "Lockheed Martin", "ESA", "JAXA", "Rocket Lab", "Virgin Galactic", "Copenhagen Suborbitals", "ROSCOSMOS", "CNSA")
    val customer_countries = List("United States", "United States", "United States", "United States", "United States", "United States", "United States", "United States", "United States", "United States", "France", "Japan", "New Zealand", "England", "Denmark", "Russia", "China")
    val customer_cities = List("Hawthorne CA", "Kent WA", "Dulles VA", "Chicago IL", "Dulles VA", "Sparks NV", "Mojave CA", "Mojave CA", "Washington DC", "Bethesda MD", "Paris", "Tokyo", "Auckland", "London", "Copenhagen", "Moscow", "Beijing")

    val product_names = List("Dragon Capsule", "Falcon 9 Rocket", "Dream Chaser Cargo System", "Biconic Farrier", "Second-stage Fuselage", "Life Support Systems", "Reaction Wheels", "Air Jordans", "Geosynchronous Satellite", "Docking Ports (x3)", "Space Junk")
    val product_categories = List("Rocket", "Rocket", "System", "Part", "Part", "System", "Part", "Misc.", "Satellite", "Part", "Misc.")
    val product_prices = List("$100,000", "$10,000,000", "$1,000,000", "$1000", "$10,000", "$10,000", "$100", "Priceless", "$1,000,000", "$1000", "$0")

    val payment_types = List("Mastercard", "Discover", "Capital One", "Zelle Transfer", "UPI", "Google Wallet", "Apple Pay")

    val failure_reasons = List("Invalid CVV", "Not Enough Balance", "Incorrect Payment Address", "Suspicious Purchase Activity", "They're totally using this to make a bomb...")

    val r = scala.util.Random
    var now = java.time.Instant.now
    //val file = scala.tools.nsc.io.File("transactions.csv")
    val file = new File("input/transactions.csv" )
    val printWriter = new PrintWriter(file)
    file.delete()

    for(i <- 1 to 2000){
      val rand_customer = r.nextInt(customer_names.length)
      val rand_product = r.nextInt(product_names.length)
      val rand_payment = payment_types(r.nextInt(payment_types.length))
      val rand_quantity = (-Math.log(r.nextDouble())*10).toInt + 1
      val rand_txn_id = (r.alphanumeric take 10).mkString
      val rand_success = if (r.nextInt(100) == 0) "N" else "Y"
      val rand_reason = if(rand_success == "Y") " " else failure_reasons(r.nextInt(failure_reasons.length))
      val rand_time_pass = r.nextInt(50000)
      now = now.plusSeconds(rand_time_pass)

      //order_id, customer_id, customer_name, product_id, product_name, product_category, payment_type, qty, price, datetime, country, city, ecommerce_website_name, payment_txn_id, payment_txn_success, failure_reason
      val transaction = List(i, 101 + rand_customer, customer_names(rand_customer), 10001 + rand_product, product_names(rand_product), product_categories(rand_product), rand_payment, rand_quantity, product_prices(rand_product), now, customer_countries(rand_customer), customer_cities(rand_customer), "AllTheSpaceYouNeed.com", rand_txn_id, rand_success, rand_reason).mkString(",")

      //println(transaction)
      //file.appendAll(transaction + "\n")
      printWriter.write(transaction + "\n")
    }


    //    val df = spark.readStream
    //      .format("kafka")
    //      .option("kafka.bootstrap.servers", "localhost:9092")
    //      .option("subscribe", "kafka_test_topic")
    //      .option("startingOffsets", "earliest") // From starting
    //      .load()
    //
    //    df.printSchema()
    //
    //    val personStringDF = df.selectExpr("CAST(value AS STRING)")

    val userSchema = new StructType().add("order_id", "integer").add("customer_id", "integer").add("customer_name", "string")
      .add("product_id", "integer").add("product_name", "string").add("product_category", "string").add("qty", "integer")
      .add("price", "integer").add("datetime", "string").add("country", "string").add("city", "string")
      .add("ecommerce_website_name", "string").add("payment_txn_id", "string")
      .add("payment_txn_success", "string").add("failure_reason", "string")



//    val df = spark.readStream
//      .format("rate")
//      .option("rowsPerSecond", 10)
//      .load()


    //    df.writeStream
    //      .option("checkpointLocation", "/input/")
    //      .toTable("myTable")
    //
    //    // Check the table result
    //    spark.read.table("myTable").show()
    // Write the streaming DataFrame to a table
    /* df.writeStream
       .option("checkpointLocation", "path/to/checkpoint/dir")
       .toTable("myTable")

     spark.read.table("myTable").show()*/
    /*      val df = spark.read.csv("transactions.csv")
          df.show()*/

    //    val csvDF = personStringDF.select(from_json(col("value"), userSchema).as("data"))
    //      .select("data.*")

    //READ THE .CSV FILE
    val csvDF = spark
      .readStream
      .option("sep", ",")
//      .option("kafka.bootstrap.servers", "host1:port1,host2:port2")
//      .option("subscribe", "topic1")
      .schema(userSchema)      // Specify schema of the csv files
      .format("csv")
      .load("input\\")    // Equivalent to format("csv").load("/path/to/directory")


//    csvDF.writeStream
//      .format("console")
//      .outputMode("append")
//      .start()
//      .awaitTermination()

    //WRITE IT TO A KAFKA TOPIC
    csvDF
      .writeStream // use `write` for batch, like DataFrame
      .format("kafka")
      .option("kafka.bootstrap.servers", "localhost:9092")
      .option("topic", "target_topic")
      .option("checkpointLocation", "tmp/vaquarkhan/checkpoint")
      .start()

  }

}

这是我得到的错误:

Exception in thread "stream execution thread for [id = c9d29df1-5f05-44b3-958e-a99e1a87beb9, runId = e0c3b3e5-05ff-4aee-be56-38ae8016bd40]" org.apache.spark.SparkException: Exception thrown in awaitResult: 
    at org.apache.spark.util.ThreadUtils$.awaitResult(ThreadUtils.scala:301)
    at org.apache.spark.rpc.RpcTimeout.awaitResult(RpcTimeout.scala:75)
    at org.apache.spark.rpc.RpcEndpointRef.askSync(RpcEndpointRef.scala:103)
    at org.apache.spark.rpc.RpcEndpointRef.askSync(RpcEndpointRef.scala:87)
    at org.apache.spark.sql.execution.streaming.state.StateStoreCoordinatorRef.deactivateInstances(StateStoreCoordinator.scala:119)
    at org.apache.spark.sql.streaming.StreamingQueryManager.notifyQueryTermination(StreamingQueryManager.scala:402)
    at org.apache.spark.sql.execution.streaming.StreamExecution.$anonfun$runStream$3(StreamExecution.scala:352)
    at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.scala:18)
    at org.apache.spark.util.UninterruptibleThread.runUninterruptibly(UninterruptibleThread.scala:77)
    at org.apache.spark.sql.execution.streaming.StreamExecution.org$apache$spark$sql$execution$streaming$StreamExecution$$runStream(StreamExecution.scala:333)
    at org.apache.spark.sql.execution.streaming.StreamExecution$$anon$1.run(StreamExecution.scala:209)
Caused by: org.apache.spark.rpc.RpcEnvStoppedException: RpcEnv already stopped.
    at org.apache.spark.rpc.netty.Dispatcher.postMessage(Dispatcher.scala:176)
    at org.apache.spark.rpc.netty.Dispatcher.postLocalMessage(Dispatcher.scala:144)
    at org.apache.spark.rpc.netty.NettyRpcEnv.askAbortable(NettyRpcEnv.scala:242)
    at org.apache.spark.rpc.netty.NettyRpcEndpointRef.askAbortable(NettyRpcEnv.scala:555)
    at org.apache.spark.rpc.netty.NettyRpcEndpointRef.ask(NettyRpcEnv.scala:559)
    at org.apache.spark.rpc.RpcEndpointRef.askSync(RpcEndpointRef.scala:102)
    ... 8 more

Process finished with exit code 0

我已经尝试了多种方法,并且我得到的最远的方法是在 IDEA 内的控制台中编写它。我的不同方法导致了错误,或者代码运行成功但消费者没有收到异常结果。

【问题讨论】:

  • 您是否获得了在 WSL2 外部运行的 Kafka CLI 工具以与在内部运行的代理一起工作?换句话说,您可能会遇到端口转发问题
  • @OneCricketeer 我没有,我只使用过生产者和消费者之间的基本 kafka 命令,但 WSL 中的所有内容都非常简单。如果您有任何可以提供帮助的信息,我将不胜感激。
  • 使用 Kafka 文件夹内的 bin\windows 目录从 WSL 外部运行相同的命令
  • 我在 WSL 中运行了 zookeeper 和 broker,在 windows 中运行了消费者,但它仍然无法正常工作。我也在windows中跑了zookeeper、broker和consumer,还是不行。
  • @OneCricketeer 嘿,我实际上能够解决我的问题。只是想感谢您,因为您的回答解决了一半的问题并为我指明了正确的方向。我只是在运行 IntelliJ IDEA 的 Windows 中运行 Kafka,一旦我的代码运行没有错误,它就可以正常工作。这一直是一个沟通问题,而且老实说我的代码不是很好。

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


【解决方案1】:

我能够使用此视频解决我的代码中的错误:

https://youtu.be/OPTMje7wKmU

所有功劳归于它的创造者。

问题 1:正如 @OneCricketeer 所指出的,我的 IntelliJ IDEA 没有与 WSL 中的 Kafka 正确通信。我通过在 Windows 而不是 WSL 下运行 Kafka 解决了这个问题。

问题 2:我的代码充满了错误。我相信由于我缺乏这门学科的知识,我对我想要达到的目标采取了不正确的方法。多亏了上面的视频,我能够编辑我的代码来完成我最初的打算,即从 IDE 发送消息以供 Kafka 消费者客户端使用。

【讨论】:

  • 我看不出这如何回答您的问题,因为您的原始代码使用的是 Spark。但是,如果您不使用分布式系统,则不需要 Spark
猜你喜欢
  • 2017-09-05
  • 2017-03-19
  • 2018-09-15
  • 2021-04-03
  • 1970-01-01
  • 2017-06-22
  • 1970-01-01
  • 1970-01-01
  • 2014-11-17
相关资源
最近更新 更多