【发布时间】:2015-10-03 20:16:29
【问题描述】:
我正在尝试使用 Flink 的 KafkaSource 运行一个简单的测试程序。我正在使用以下内容:
- Flink 0.9
- Scala 2.10.4
- Kafka 0.8.2.1
按照here 和here 的描述,我按照文档测试了KafkaSource(添加了依赖项,将Kafka 连接器flink-connector-kafka 捆绑在插件中)。
下面是我的简单测试程序:
import org.apache.flink.streaming.api.scala._
import org.apache.flink.streaming.connectors.kafka
object TestKafka {
def main(args: Array[String]) {
val env = StreamExecutionEnvironment.getExecutionEnvironment
val stream = env
.addSource(new KafkaSource[String]("localhost:2181", "test", new SimpleStringSchema))
.print
}
}
但是,编译总是报错 KafkaSource not found:
[ERROR] TestKafka.scala:8: error: not found: type KafkaSource
[ERROR] .addSource(new KafkaSource[String]("localhost:2181", "test", new SimpleStringSchema))
我错过了什么?
【问题讨论】:
标签: scala apache-kafka apache-flink