【发布时间】:2021-02-15 15:15:33
【问题描述】:
我能够从 Kafka 主题读取数据,并能够使用 spark 流在控制台上打印数据。
我希望数据采用数据框格式。
这是我的代码:
spark = SparkSession \
.builder \
.appName("StructuredSocketRead") \
.getOrCreate()
spark.sparkContext.setLogLevel('ERROR')
lines = spark \
.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers","********") \
.option("subscribe","******") \
.option("startingOffsets", "earliest") \
.load()
readable = lines.selectExpr("CAST(value AS STRING)")
query = readable \
.writeStream \
.outputMode("append") \
.format("console") \
.option("truncate", "False") \
.start()
query.awaitTermination()
输出为 JSON 文件格式。如何将其转换为数据框?请在下面找到输出:
{"items": [{"SKU": "23565", "title": "EGG CUP MILKMAID HELGA ", "unit_price": 2.46, "quantity": 2}], "type": "ORDER", "country": "United Kingdom", "invoice_no": 154132541847735, "timestamp": "2020-11-02 20:56:01"}
【问题讨论】:
-
好吧,你已经反序列化为一个字符串,所以现在你需要定义一个模式并将其应用到
readabledatabricks.com/blog/2017/04/26/… -
您要查找的输出格式是什么?同时参考这个,我在这里回答过的类似问题 - stackoverflow.com/questions/64640565/…
标签: apache-spark pyspark apache-kafka apache-spark-sql spark-structured-streaming