【问题标题】:Spark Streaming df.writeStream generate no outputSpark Streaming df.writeStream 不生成输出
【发布时间】: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


【解决方案1】:

我怀疑虚拟机内 Docker 容器内的 kafka 代理的adverted.listeners 没有配置为外部连接

【讨论】:

  • Kakfa 代理端口 6667 已从 Docker 到主机开放。
  • 我认为您不了解广告 listeners 变量的用途。端口转发很重要,但不是唯一重要的配置
猜你喜欢
  • 2017-04-04
  • 1970-01-01
  • 2017-10-25
  • 1970-01-01
  • 2019-01-15
  • 2021-02-28
  • 2019-02-19
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多