【问题标题】:Trying to consuming the kafka streams using spark structured streaming尝试使用火花结构化流来使用 kafka 流
【发布时间】:2020-10-25 11:06:20
【问题描述】:

我是 Kafka 流媒体的新手。我使用 python 设置了一个 twitter 监听器,它在 localhost:9092 kafka 服务器中运行。我可以使用 kafka 客户端工具(conduktor)以及使用命令“bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic twitter --from-beginning”来使用监听器产生的流 但是,当我尝试使用 Spark 结构化流式传输使用相同的流时,它没有捕获并抛出错误 - 找不到数据源:kafka。请按照《Structured Streaming + Kafka Integration Guide》的部署部分部署应用; 找到下面的截图

  1. Command output - Consumes Data
  2. Jupyter output for spark consumer - Doesn't consume data

我的生产者或侦听器代码:

auth = tweepy.OAuthHandler("**********", "*************")
auth.set_access_token("*************", "***********************")
# session.set('request_token', auth.request_token)
api = tweepy.API(auth)
class KafkaPushListener(StreamListener):          
    def __init__(self):
        #localhost:9092 = Default Zookeeper Producer Host and Port Adresses
        self.client = pykafka.KafkaClient("0.0.0.0:9092")

    #Get Producer that has topic name is Twitter
        self.producer = self.client.topics[bytes("twitter", "ascii")].get_producer()

    def on_data(self, data):
        #Producer produces data for consumer
        #Data comes from Twitter
        self.producer.produce(bytes(data, "ascii"))
        return True

    def on_error(self, status):
        print(status)
        return True
twitter_stream = Stream(auth, KafkaPushListener())
twitter_stream.filter(track=['#fashion'])

来自 Spark 结构化流的消费者访问

df = spark \
  .readStream \
  .format("kafka") \
  .option("kafka.bootstrap.servers", "localhost:9092") \
  .option("subscribe", "twitter") \
  .load()
df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")

【问题讨论】:

  • 你能发布完整的代码吗? &我在屏幕截图中没有看到任何水槽..您使用的是什么水槽??
  • 已更新代码

标签: apache-spark pyspark apache-kafka spark-structured-streaming spark-streaming-kafka


【解决方案1】:

发现缺少的东西,当我提交 spark-job 时,我必须包含正确的依赖包版本。 我有火花 3.0.0 因此,我包括了 - org.apache.spark:spark-sql-kafka-0-10_2.12:3.0.0 包

【讨论】:

    【解决方案2】:

    添加sink会从kafka开始消费数据。

    检查下面的代码。

    df = spark \
      .readStream \
      .format("kafka") \
      .option("kafka.bootstrap.servers", "localhost:9092") \
      .option("subscribe", "twitter") \
      .load()
    
    query = df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)") \
        .writeStream \
        .outputMode("append") \
        .format("console") \ # here I am using console format .. you may change as per your requirement.
        .start()
    
    query.awaitTermination()
    

    【讨论】:

      猜你喜欢
      • 2020-08-18
      • 2020-03-10
      • 2019-09-20
      • 2021-05-31
      • 2020-02-12
      • 1970-01-01
      • 1970-01-01
      • 2018-12-22
      • 2017-04-27
      相关资源
      最近更新 更多