【问题标题】:Pyspark - print messages from KafkaPyspark - 打印来自 Kafka 的消息
【发布时间】:2019-04-21 05:36:25
【问题描述】:

我建立了一个带有生产者和消费者的 kafka 系统,将 json 文件的行作为消息流式传输。

使用 pyspark,我需要分析不同流窗口的数据。为此,我需要查看 pyspark 流式传输的数据......我该怎么做?

为了运行我使用Yannael's Docker 容器的代码。这是我的python代码:

# Add dependencies and load modules
import os
os.environ['PYSPARK_SUBMIT_ARGS'] = '--conf spark.ui.port=4040 --packages org.apache.spark:spark-streaming-kafka-0-8_2.11:2.0.0,com.datastax.spark:spark-cassandra-connector_2.11:2.0.0-M3 pyspark-shell'

from kafka import KafkaConsumer
from random import randint
from time import sleep

# Load modules and start SparkContext  
from pyspark import SparkContext, SparkConf
from pyspark.sql import SQLContext, Row
conf = SparkConf() \
    .setAppName("Streaming test") \
    .setMaster("local[2]") \
    .set("spark.cassandra.connection.host", "127.0.0.1")

try:
    sc.stop()
except:
    pass    

sc = SparkContext(conf=conf) 
sqlContext=SQLContext(sc)
from pyspark.streaming import StreamingContext
from pyspark.streaming.kafka import KafkaUtils

# Create streaming task
ssc = StreamingContext(sc, 0.60)
kafkaStream = KafkaUtils.createStream(ssc, "127.0.0.1:2181", "spark-streaming-consumer", {'test': 1})
ssc.start()

【问题讨论】:

标签: python apache-spark pyspark apache-kafka


【解决方案1】:

您可以致电kafkaStream.pprint(),或了解更多信息about structured streaming,您可以像这样打印

query = kafkaStream \
    .writeStream \
    .outputMode("complete") \
    .format("console") \
    .start()

query.awaitTermination()

我看到你有 endpoints,所以假设你正在写入 Cassandra,你可以使用 Kafka Connect 而不是为此编写 Spark 代码

【讨论】:

  • 谢谢@cricket_007!作为第一次测试,我加入了kafkaStream.pprint(),结果我得到了当前时间……你对如何获得正确的消息有什么建议吗?
  • 不确定我理解你所说的“正确”是什么意思
  • 我不应该将test1 主题中的消息视为输出吗?据我了解,我订阅了kafkaStream
  • 是的,您应该这样做,但前提是它们正被积极地纳入主题。默认情况下,它使用最新的偏移量
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-05-08
  • 1970-01-01
  • 2019-09-01
  • 2019-11-15
  • 1970-01-01
相关资源
最近更新 更多