【发布时间】:2019-01-15 08:11:32
【问题描述】:
我尝试使用 pyspark 将 spark 和 kafka 集成到 Jupyter notebook 中。这是我的工作环境。
Spark 版本:Spark 2.2.1 卡夫卡版本:Kafka_2.11-0.8.2.2 火花流kafka jar:spark-streaming-kafka-0-8-assembly_2.11-2.2.1.jar
我在 spark-defaults.conf 文件中添加了一个 Spark 流式 kafka 程序集 jar 文件。
当我为 pyspark 流式传输启动 streamingContext 时,此错误显示为 无法从 MANIFEST.MF 读取 kafka 版本。
这是我的代码。
from pyspark import SparkContext, SparkConf
from pyspark import SparkContext, SparkConf
from pyspark.streaming import StreamingContext
from pyspark.streaming.kafka import KafkaUtils
import sys
import os
from kafka import KafkaProducer
#Receive data handler
def handler(message):
records = message.collect()
for record in records:
print(record)
#producer.send('receive', str(res))
#producer.flush()
producer = KafkaProducer(bootstrap_servers='slave02:9092')
sc = SparkContext(appName="SparkwithKafka")
ssc = StreamingContext(sc, 1)
#Create Kafka streaming with argv
zkQuorum = 'slave02:2181'
topic = 'send'
kvs = KafkaUtils.createStream(ssc, zkQuorum, "spark-streaming-consumer", {topic:1})
kvs.foreachRDD(handler)
ssc.start()
【问题讨论】:
-
你提交代码的命令是什么?或者你如何将 JAR 加载到 jupyter 中?
-
编辑:我的错,看起来不像 python 的 api 注意:KafkaUtils.createStream 是阅读 kafka 主题的旧方法。你应该使用Kafka 0.10 api
-
@cricket_007 ssc.start() 是起始代码。我在 Jupyter 笔记本中进行了测试。我将 jar 路由附加到 spark-defaults.conf 中,例如 spark.jars spark-streaming-kafka-0-8-assembly_2.11-2.2.1.jar
-
@Bameza 但我知道 spark 与 0.10.0 版本不兼容。 link 这就是我使用 kafka 0.8.0 API 的原因。我可以在 Jupyter notebook 中使用 kafka 0.10.0 API 吗?
-
我发现了这个问题。这是一个警告,我的应用程序不会因此而崩溃。这对我来说很有用。谢谢大家!
标签: apache-spark pyspark apache-kafka spark-streaming spark-streaming-kafka