代码之家  ›  专栏  ›  技术社区  ›  Cassie

使用Kafka 2.4.0的MicrobatchExecution Spark结构化流媒体

  •  0
  • Cassie  · 技术社区  · 7 年前

    当我尝试使用Kafka集成运行Spark结构化流式应用程序时,我一直收到以下错误:

    ERROR MicroBatchExecution: Query [id = ff14fce6-71d3-4616-bd2d-40f07a85a74b, runId = 42670f29-21a9-4f7e-abd0-66ead8807282] terminated with error
    java.lang.IllegalStateException: No entry found for connection 2147483647
    

    为什么会这样?可能是一些依赖性问题吗?

    我的build.sbt文件如下:

    name := "SparkAirflowK8s"
    
    version := "0.1"
    
    scalaVersion := "2.12.7"
    
    val sparkVersion = "2.4.0"
    val circeVersion = "0.11.0"
    
    dependencyOverrides += "com.fasterxml.jackson.core" % "jackson-core" % "2.9.8"
    dependencyOverrides += "com.fasterxml.jackson.core" % "jackson-databind" % "2.9.8"
    dependencyOverrides += "com.fasterxml.jackson.module" % "jackson-module-scala_2.12" % "2.9.8"
    
    resolvers += "Spark Packages Repo" at "http://dl.bintray.com/spark-packages/maven"
    resolvers += "confluent" at "http://packages.confluent.io/maven/"
    
    libraryDependencies ++= Seq(
      "org.apache.spark" %% "spark-core" % sparkVersion,
      "org.apache.spark" %% "spark-sql" % sparkVersion,
      "org.apache.spark" %% "spark-streaming" % sparkVersion,
      "org.apache.spark" %% "spark-sql-kafka-0-10" % sparkVersion,
      "org.apache.kafka" %% "kafka" % "2.1.0",
      "org.scalatest" %% "scalatest" % "3.2.0-SNAP10" % "it, test",
      "org.scalacheck" %% "scalacheck" % "1.14.0" % "it, test",
      "io.kubernetes" % "client-java" % "3.0.0" % "it",
      "org.json" % "json" % "20180813",
      "io.circe" %% "circe-core" % circeVersion,
      "io.circe" %% "circe-generic" % circeVersion,
      "io.circe" %% "circe-parser" % circeVersion,
      "org.apache.avro" % "avro" % "1.8.2",
      "io.confluent" % "kafka-avro-serializer" % "5.0.1"
    )
    

    这是代码的一部分:

    val sparkConf = new SparkConf()
          .setMaster(args(0))
          .setAppName("KafkaSparkJob")
    
        val sparkSession = SparkSession
          .builder()
          .config(sparkConf)
          .getOrCreate()
    
    val avroStream = sparkSession.readStream
          .format("kafka")
          .option("kafka.bootstrap.servers", "localhost:9092")
          .option("subscribe", "topic1")
          .load()
    
    val outputExample = avroStream
          .writeStream
          .outputMode("append")
          .format("console")
          .start()
    
        outputExample.awaitTermination()
    
    1 回复  |  直到 7 年前
        1
  •  0
  •   Cassie    7 年前

    我已经将localhost更改为为为kafka部署定义的nodeport服务。现在,这个例外不会出现。

    推荐文章