【问题标题】:Spark Streaming is reading from Kafka topic and how to convert nested Json format into dataframeSpark Streaming 正在阅读 Kafka 主题以及如何将嵌套的 Json 格式转换为数据帧
【发布时间】: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"}    

【问题讨论】:

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


【解决方案1】:

IICU,请使用 explode()getItems() 以便从 json 中创建 Dataframe..

在此处创建数据框

a_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"}
df = spark.createDataFrame([(a_json)])
df.show(truncate=False)
+--------------+---------------+-------------------------------------------------------------------------------------+-------------------+-----+
|country       |invoice_no     |items                                                                                |timestamp          |type |
+--------------+---------------+-------------------------------------------------------------------------------------+-------------------+-----+
|United Kingdom|154132541847735|[[quantity -> 2, unit_price -> 2.46, title -> EGG CUP MILKMAID HELGA , SKU -> 23565]]|2020-11-02 20:56:01|ORDER|
+--------------+---------------+-------------------------------------------------------------------------------------+-------------------+-----+

这里的逻辑

df = df.withColumn("items_array", F.explode("items"))
df = df.withColumn("quantity", df.items_array.getItem("quantity")).withColumn("unit_price", df.items_array.getItem("unit_price")).withColumn("title", df.items_array.getItem("title")).withColumn("SKU", df.items_array.getItem("SKU"))
df.select("country", "invoice_no", "quantity","unit_price", "title", "SKU", "timestamp", "timestamp").show(truncate=False)
+--------------+---------------+--------+----------+-----------------------+-----+-------------------+-------------------+
|country       |invoice_no     |quantity|unit_price|title                  |SKU  |timestamp          |timestamp          |
+--------------+---------------+--------+----------+-----------------------+-----+-------------------+-------------------+
|United Kingdom|154132541847735|2       |2.46      |EGG CUP MILKMAID HELGA |23565|2020-11-02 20:56:01|2020-11-02 20:56:01|
+--------------+---------------+--------+----------+-----------------------+-----+-------------------+-------------------+

【讨论】:

  • 这对从 Kafka 构建流式数据帧有何帮助?
猜你喜欢
  • 1970-01-01
  • 2020-08-24
  • 2020-09-05
  • 2018-02-22
  • 1970-01-01
  • 2021-10-12
  • 2017-09-30
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多