【发布时间】:2020-02-18 00:04:06
【问题描述】:
我正在使用 hdp sandbox 2.6.4,并且我在本地机器(主机)上配置了 spark。
我已经使用 shell 登录到 docker 映像并启动了一个简单的控制台使用者。我正在尝试在我的本地机器(而不是 docker 容器)上使用 Spark 来使用它。它没有给出任何错误。但是,它也没有提供任何输出。
def main(args: Array[String]): Unit = {
val spark = SparkSession.builder.appName("demo")
.master("local[*]")
.getOrCreate()
spark.sparkContext.setLogLevel("WARN")
import spark.implicits._
val df = spark.readStream.
format("kafka").
option("kafka.bootstrap.servers", "localhost:6667").
option("subscribe", "test").
load()
val query = df.writeStream
.outputMode("append")
.format("console")
.start()
query.awaitTermination()
}
如果我登录到 docker image 并启动另一个控制台生产者,我就可以使用所有消息。我检查了端口,它从 docker 到主机是开放的。
【问题讨论】:
-
为什么只需要 Spark 和 Kafka 的整个 hdp VM?
-
我没有..我的机器上已经有 hdp 并且想编写一个程序将结果存储回 hive..这个 intelij 代码用于开发
-
为什么不使用 Kafka Connect HDFS 连接器或 Apache Gobblin 呢?
-
@cricket_007 可能是真正的用例,这是用于学习火花流。
-
Spark Streaming 自 Spark 2.4 起已弃用。我假设您的意思是结构化流?无论如何,您的代码看起来都不错。如果您什么也没看到,要么主题是空的,要么您没有看到的 Spark UI 执行程序日志中存在静默网络错误
标签: scala apache-spark apache-kafka spark-structured-streaming