【发布时间】: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》的部署部分部署应用; 找到下面的截图
我的生产者或侦听器代码:
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