【问题标题】:Consume no message from spark-streaming-kafka-0-10 with option kafka.bootstrap.servers使用选项 kafka.bootstrap.servers 不使用来自 spark-streaming-kafka-0-10 的消息
【发布时间】:2019-04-03 16:33:13
【问题描述】:

我正在使用来自 CDH 的 kafka 1.0.1-kafka-3.1.0-SNAPSHOT(hadoop 的 cloudera 发行版)

在我的第 1 批边缘服务器上,我可以生成消息:

kafka-console-producer --broker-list batch-1:9092 --topic MyTopic

感谢 Zookeeper 在我的第一个节点上,我可以使用消息:

kafka-console-consumer --zookeeper data1:2181 --topic MyTopic --from-beginning

但是使用 bootstrap-server 选项 我没有得到 nothing :

kafka-console-consumer --bootstrap-server batch-1:9092 --topic MyTopic --from-beginning

问题是我在 spark 上使用 kafka :

libraryDependencies += "org.apache.spark" %% "spark-streaming-kafka-0-10" % "2.3.0"

val df = spark.readStream
  .format("org.apache.spark.sql.kafka010.KafkaSourceProvider")
  .option("kafka.bootstrap.servers", "batch-1:9092")
  .option("subscribe", "MyTopic")
  .load()

println("Select :")

val df2 = df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)", "CAST(topic AS STRING)")
  .as[(String, String, String)]

println("Show :")

val query = df2.writeStream
  .outputMode("append")
  .format("console")
  .start()

query.awaitTermination()

我在我的边缘做了一个export SPARK_KAFKA_VERSION=0.10。那么

spark2-submit --driver-memory 2G --jars spark-sql-kafka-0-10_2.11-2.3.0.cloudera4.jar --class "spark.streaming.Poc" poc_spark_kafka_2.11-0.0.1.jar

这逼我用kafka.bootstrap.servers,好像已经连接了,但是收不到任何消息。

输出与带有--bootstrap-server 选项的kafka-console-consumer 相同:

18/10/30 16:11:48 INFO utils.AppInfoParser: Kafka version : 0.10.0-kafka-2.1.0
18/10/30 16:11:48 INFO utils.AppInfoParser: Kafka commitId : unknown
18/10/30 16:11:48 INFO streaming.MicroBatchExecution: Starting new streaming query.

然后,什么都没有。 我应该连接到 Zookeeper 吗?如何 ?

是否存在版本冲突,而他们在这里说“结构化流 + Kafka 集成指南(Kafka 代理版本 0.10.0 或更高版本)”:https://spark.apache.org/docs/latest/structured-streaming-kafka-integration.html ?

我错过了什么?

【问题讨论】:

  • 您是否检查过您的代理是否真的在运行?当我使用非常相似的火花流程序进行本地设置时,一切正常。事实上,您的控制台消费者没有显示任何内容,这让我认为生产者存在问题。
  • 从 Kafka 0.9 及更高版本开始,消费者和生产者不应使用 Zookeeper。在您期望 Spark 工作之前,带有 boostrap 服务器的控制台使用者应该工作。此外,除非另有配置,否则 Spark 可能会从您的主题的最新偏移量开始
  • @Elmar Macek,我猜生产者很好,因为我使用 Zookeeper 消费来自控制台消费者的消息。
  • @cricket_007 这正是问题所在,也是我感到失望的原因:它不应该在 1.0 中使用 zookeeper。谢谢你们两位的快速回答。

标签: apache-spark apache-kafka streaming kafka-consumer-api


【解决方案1】:

解决方案

/var/log/kafka/kafka-broker-batch-1.log 说:

2018-10-31 13:40:08,284 ERROR kafka.server.KafkaApis: [KafkaApi-51] Number of alive brokers '1' does not meet the required replication factor '3' for the offsets topic (configured via 'offsets.topic.replication.factor'). This error can be ignored if the cluster is starting up and not all brokers are up yet.

所以我在我的集​​群节点上部署了 3 个代理,并在边缘设置了一个网关,现在它可以与:

kafka-console-producer --broker-list data1:9092,data2:9092,data3:9092 --topic Test

kafka-console-consumer --bootstrap-server data1:9092 --topic Test --from-beginning

Spark 也可以正常工作。

【讨论】:

  • 所以,是的,使用 Zookeeper 会存储偏移量,而不是偏移量主题
猜你喜欢
  • 2019-08-08
  • 2018-01-09
  • 2018-02-09
  • 2017-07-13
  • 2016-11-03
  • 2018-01-13
  • 2020-08-04
  • 2018-06-25
  • 2017-10-30
相关资源
最近更新 更多