【问题标题】:Unable to see messages from Kafka Stream in Spark在 Spark 中无法看到来自 Kafka Stream 的消息
【发布时间】:2018-03-12 03:39:40
【问题描述】:

我刚刚开始使用Pyspark 库测试Kafka StreamSpark

我一直在 Jupyter Notebook 上运行整个设置。 我正在尝试从Twitter Streaming 获取数据。

Twitter 流代码:

import json
import tweepy
from uuid import uuid4
import time
from kafka import KafkaConsumer
from kafka import KafkaProducer

auth = tweepy.OAuthHandler("key", "key")
auth.set_access_token("token", "token")
api = tweepy.API(auth, wait_on_rate_limit=True, retry_count=3, retry_delay=5,
                 retry_errors=set([401, 404, 500, 503]))
class CustomStreamListener(tweepy.StreamListener):
    def __init__(self, api):
        self.api = api
        super(tweepy.StreamListener, self).__init__()

    def on_data(self, tweet):
        print tweet
        # Kafka Producer to send data to twitter topic
        producer.send('twitter', json.dumps(tweet))

    def on_error(self, status_code):
        print status_code
        return True # Don't kill the stream

    def on_timeout(self):
        print 'on_timeout'
        return True # Don't kill the stream
producer = KafkaProducer(bootstrap_servers='localhost:9092')
sapi = tweepy.streaming.Stream(auth, CustomStreamListener(api))
sapi.filter(track=["#party"])

Spark 流式处理代码

from pyspark import SparkContext
from pyspark.streaming import StreamingContext
from pyspark.streaming.kafka import KafkaUtils

sc = SparkContext(appName="PythonSparkStreamingKafka_RM_01").getOrCreate()
sc.setLogLevel("WARN")

streaming_context = StreamingContext(sc, 10)
kafkaStream = KafkaUtils.createStream(streaming_context, 'localhost:2181', 'spark-streaming', {'twitter': 1})  
parsed = kafkaStream.map(lambda v: v)
parsed.count().map(lambda x:'Tweets in this batch: %s' % x).pprint()

streaming_context.start()
streaming_context.awaitTermination()

打印输出:

时间:2017-09-30 11:21:00


时间:2017-09-30 11:21:10


时间:2017-09-30 11:21:20

我做错了什么特定的部分?

【问题讨论】:

  • 你能解决这个错误吗?我面临同样的问题。你能帮帮我吗?

标签: apache-spark pyspark apache-kafka spark-streaming twitter-streaming-api


【解决方案1】:

您还可以使用一些 GUI 工具,例如 Kafdrop。它在调试 kafka 消息时非常方便。您不仅可以查看消息队列,还可以查看它们的偏移量等的分区。

它是一个很好的工具,您应该能够轻松部署它。

这里是链接:https://github.com/HomeAdvisor/Kafdrop

【讨论】:

    【解决方案2】:

    您可以使用以下两个步骤调试应用程序。

    1) 使用 KafkaWordCount 之类的示例消费者来测试是否有数据进来(Kafka 主题是否有消息)

    Kafka 带有一个命令行客户端,该客户端将从文件或标准输入中获取输入,并将其作为消息发送到 Kafka 集群。默认情况下,每行将作为单独的消息发送。

    运行生产者,然后在控制台中输入一些消息以发送到服务器。

         kafka-console-producer.sh \
        --broker-list <brokeer list> \
        --topic <topic name> \
        --property parse.key=true \
        --property key.separator=, \
        --new-producer  
    

    例子:

       > bin/kafka-console-producer.sh --broker-list localhost:9092 --topic test
    

    如果你看到打印消息,那么你在 kafka 中有消息,如果没有,那么你的生产者没有工作

    2) 开启日志记录

      Logger.getLogger("org").setLevel(Level.WARNING);
      Logger.getLogger("akka").setLevel(Level.WARNING);       
      Logger.getLogger("kafka").setLevel(Level.WARNING);
    

    【讨论】:

    • 好的,我会试试这个,但我使用createDirectStream 完成了这项工作,我不知道如何,但完全相同的设置只需使用directStream
    • @NikhilParmar:我也是。我不知道为什么它适用于directStream,但消费者在使用createDirectStream 时没有收到任何消息。我使用 Spark 2.4.5。有人知道吗?
    猜你喜欢
    • 2018-01-09
    • 2021-07-09
    • 1970-01-01
    • 2018-08-29
    • 2018-07-18
    • 1970-01-01
    • 1970-01-01
    • 2018-02-09
    • 2021-04-26
    相关资源
    最近更新 更多