【问题标题】:object FlinkKafkaConsumer010 is not a member of package org.apache.flink.streaming.connectors.kafka对象 FlinkKafkaConsumer010 不是包 org.apache.flink.streaming.connectors.kafka 的成员
【发布时间】:2020-09-19 01:59:08
【问题描述】:

我正在尝试组装一个小程序来使用 apache flink 连接到 kafka 主题。我需要使用 FlinkKafkaConsumer010。

package uimp
import org.apache.flink.streaming.api.scala._
import org.apache.flink.streaming.api.TimeCharacteristic
import org.apache.flink.streaming.util.serialization.SimpleStringSchema
import org.apache.flink.streaming.connectors.kafka.{FlinkKafkaConsumer010}
import java.util.Properties

object Silocompro {
  def main(args: Array[String]): Unit = {
 // set up the execution environment
    val env = StreamExecutionEnvironment.getExecutionEnvironment
    env.setStreamTimeCharacteristic(TimeCharacteristic.ProcessingTime)

    val propertiesTopicDemographic = new Properties()
    propertiesTopicDemographic.setProperty("bootstrap.servers", "bigdata.dataspartan.com:19093")
    propertiesTopicDemographic.setProperty("group.id", "demographic")

    val myConsumerDemographic = new FlinkKafkaConsumer010[String]("topic_demographic", new 
    SimpleStringSchema(), propertiesTopicDemographic)

    val messageStreamDemographic = env
      .addSource(myConsumerDemographic)
      .print()


    env.execute("Flink Scala API Skeleton")

   }
 }

我的问题是当我尝试用这个 build.sbt 组装我的程序时,编译器返回错误“object FlinkKafkaConsumer010 is not a member of package org.apache.flink.streaming.connectors.kafka”:

      ThisBuild / resolvers ++= Seq("Apache Development Snapshot Repository" at 
      "https://repository.apache.org/content/repositories/snapshots/",Resolver.mavenLocal)

      name := "silocompro"

      version := "1.0"

      organization := "uimp"

      ThisBuild / scalaVersion := "2.12.11"

      val flinkVersion = "1.9.0"

      val flinkDependencies = Seq(
         "org.apache.flink" %% "flink-scala" % flinkVersion % "provided",
         "org.apache.flink" %% "flink-streaming-scala" % flinkVersion % "provided",
         "org.apache.flink" %% "flink-core"% flinkVersion % "provided",
         "org.apache.flink" %% "flink-connector-kafka-base" % flinkVersion % "provided",
         "org.apache.flink" %% "flink-clients" % flinkVersion % "provided",
         "org.apache.flink" %% "flink-connector-kafka" % flinkVersion % "provided")

      lazy val root = (project in file(".")).
      settings( libraryDependencies ++= flinkDependencies)


      assembly / mainClass := Some("uimp.Silocompro")

      Compile / run  := Defaults.runTask(Compile / fullClasspath,
                               Compile / run / mainClass,
                               Compile / run / runner
                              ).evaluated

 
      Compile / run / fork := true
      Global / cancelable := true

      assembly / assemblyOption  := (assembly / assemblyOption).value.copy(includeScala = false)

这个依赖错误的原因是什么?

【问题讨论】:

    标签: apache-kafka sbt apache-flink flink-streaming


    【解决方案1】:

    连接器不是 flink-binary 的一部分,这意味着您需要在 compile 范围内有连接器,所以这基本上意味着您需要从这些依赖项中删除 provided。在此设置下,应用程序将在集群上运行。

    但是,如果您想在不启动集群的情况下在本地运行它,那么您应该在 compile 范围内拥有所有 flink 依赖项,即删除所有 provided 范围声明。

    【讨论】:

    • 我已经删除了所有提供的范围声明,但是在我组装我的程序时仍然出现错误,“对象连接器不是包 org.apache.flink.streaming 的成员”
    【解决方案2】:

    最后我遇到了依赖问题。我做了一些动作:

    1. 我添加了一个新的解析器 https://oss.sonatype.org/content/repositories
    2. 我已经卸载了 来自 VS Code 的插件 metal(scala)
    3. 我已经添加了"org.apache.flink"%% "flink-connector-kafka-0.10"%flinkVersion - tomy flinkDependencies

    在此操作之后,我解决了我的库依赖问题。谢谢

    【讨论】:

      猜你喜欢
      • 2022-01-01
      • 2018-02-21
      • 2012-12-20
      • 2018-09-10
      • 2021-12-06
      • 2021-06-03
      • 2016-02-29
      • 2018-10-25
      • 2016-09-25
      相关资源
      最近更新 更多