【发布时间】:2017-01-01 03:59:52
【问题描述】:
我想使用 Apache Kafka 和 Spark 流设置一个流应用程序。 Kafka 在单独的 0.9.0.1 版本的 unix 机器上运行,spark v1.6.1 是 hadoop 集群的一部分。
我已经启动了 zookeeper 和 kafka 服务器,并希望使用控制台生产者从日志文件中流式传输消息,并使用直接方法(无接收器)由 spark 流应用程序使用。我已经在 python 中编写了代码并使用以下命令执行:
spark-submit --jars spark-streaming-kafka-assembly_2.10-1.6.1.jar streamingDirectKafka.py
出现以下错误:
/opt/mapr/spark/spark-1.6.1/python/lib/pyspark.zip/pyspark/streaming/kafka.py", line 152, in createDirectStream
py4j.protocol.Py4JJavaError: An error occurred while calling o38.createDirectStreamWithoutMessageHandler.
: java.lang.ClassCastException: kafka.cluster.BrokerEndPoint cannot be cast to kafka.cluster.Broker
你能帮忙吗?
谢谢!!
from pyspark import SparkContext, SparkConf
from pyspark.streaming import StreamingContext
from pyspark.streaming.kafka import KafkaUtils
if __name__ == "__main__":
conf = SparkConf().setAppName("StreamingDirectKafka")
sc = SparkContext(conf = conf)
ssc = StreamingContext(sc, 1)
topic = ['test']
kafkaParams = {"metadata.broker.list": "apsrd7102:9092"}
lines = (KafkaUtils.createDirectStream(ssc, topic, kafkaParams)
.map(lambda x: x[1]))
counts = (lines.flatMap(lambda line: line.split(" "))
.map(lambda word: (word, 1))
.reduceByKey(lambda a, b: a+b))
counts.pprint()
ssc.start()
ssc.awaitTermination()
【问题讨论】:
标签: apache-spark