【发布时间】: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