【问题标题】:PySpark Streaming + Kafka Word Count not printing any resultsPySpark Streaming + Kafka Word Count 不打印任何结果
【发布时间】:2018-02-26 11:50:14
【问题描述】:

这是我与 Kafka 和 Spark Streaming 的第一次互动,我正在尝试运行下面给出的 WordCount 脚本。正如许多在线博客中给出的那样,该脚本非常标准。但无论出于何种原因,火花流不会打印字数。它没有抛出任何错误,只是不显示计数。 我已经通过控制台消费者测试了该主题,并且消息正确显示。我什至尝试使用 foreachRDD 来查看输入的行,但也没有显示任何内容。

提前致谢!

版本:kafka_2.11-0.8.2.2、Spark2.2.1、spark-streaming-kafka-0-8-assembly_2.11-2.2.1

from __future__ import print_function

import sys
from pyspark import SparkContext
from pyspark import SparkConf
from pyspark.streaming import StreamingContext
from pyspark.streaming.kafka import KafkaUtils
from pyspark.sql.context import SQLContext

sc = SparkContext(appName="PythonStreamingKafkaWordCount")
sc.setCheckpointDir('c:\Playground\spark\logs')
ssc = StreamingContext(sc, 10)
ssc.checkpoint('c:\Playground\spark\logs')

zkQuorum, topic = sys.argv[1:]
print(str(zkQuorum))
print(str(topic))
kvs = KafkaUtils.createStream(ssc, zkQuorum, "spark-streaming-consumer", {topic: 1})
lines = kvs.map(lambda x: x[1])
print(kvs)

counts = lines.flatMap(lambda line: line.split(" ")) \
                  .map(lambda word: (word, 1)) \
                  .reduceByKey(lambda a, b: a+b)
counts.pprint(num=10)
ssc.start()
ssc.awaitTermination()  

生产者代码:

import sys,os
from kafka import KafkaProducer
from kafka.errors import KafkaError
import time

producer = KafkaProducer(bootstrap_servers="localhost:9092")
topic = "KafkaSparkWordCount"

def read_file(fileName):
    with open(fileName) as f:
        print('started reading...')
        contents = f.readlines()
        for content in contents:
            future = producer.send(topic,content.encode('utf-8'))
            try:
                future.get(timeout=10)
            except KafkaError as e:
                print(e)
                break
            print('.',end='',flush=True)
            time.sleep(0.2)

    print('done')       


if __name__== '__main__':
    read_file('C:\\\PlayGround\\spark\\BookText.txt')

【问题讨论】:

  • 您的代码似乎没有问题。来自服务器的数据不是以流方式传输的。检查 Kafka Producer 代码。
  • 嗨 Mayank,感谢您的回复。我在上面也添加了生产者代码,请在这里提出可能有问题的地方。
  • 数据是否足够大,可以继续流式传输一段时间。我认为数据是由 kafka 发送的,并且在 spark 接受它之前就结束了。你能检查一下吗?你能看到这个 Kafka 生产者代码发送的数据吗?意思是, print('done') 不应该打印一段时间..否则没有数据可供 spark 接受。
  • 检查过了,这个过程至少需要 5 分钟才能到达 done 语句。我确保在生产者完成之前启动火花部分。
  • 虽然不确定 createStream 接受什么,但您能否在所有参数都通过后检查一下,因为其他一切似乎都很好。

标签: apache-spark pyspark apache-kafka spark-streaming word-count


【解决方案1】:

你用了多少个内核?

Spark Streaming 至少需要两个内核,一个用于接收器,一个用于处理器。

【讨论】:

  • 我在本地机器上使用默认设置。我的机器有 4 个核心。
猜你喜欢
  • 2015-03-18
  • 1970-01-01
  • 1970-01-01
  • 2020-11-02
  • 2017-10-21
  • 1970-01-01
  • 1970-01-01
  • 2020-08-05
  • 2017-12-11
相关资源
最近更新 更多