【问题标题】:Get data from nested json in kafka stream pyspark从kafka流pyspark中的嵌套json获取数据
【发布时间】:2021-02-14 19:34:49
【问题描述】:

我有一个kafka生产者以

格式发送大量数据
{
  '1000': 
    {
       '3': 
        {
           'seq': '1', 
           'state': '2', 
           'CMD': 'XOR' 
        }
    },
 '1001': 
    {
       '5': 
        {
           'seq': '2', 
           'state': '2', 
           'CMD': 'OR' 
        }
    },
 '1003': 
    {
       '5': 
        {
           'seq': '3', 
           'state': '4', 
           'CMD': 'XOR' 
        }
    }
}

.... 我想要的数据在最后一个循环中:{'seq': '1', 'state': '2', 'CMD': 'XOR'} 上面循环中的键('1000' and '3') 是可变的。请注意,上述值仅作为示例。原始数据集很大,有很多可变键。只有最终循环中的键{'seq', 'state', 'CMD'} 是常量。

我尝试使用通用格式来读取数据,但由于上面的循环具有可变键,因此我得到了不正确的数据,我不确定如何定义架构来解析这种数据格式。

我试图实现的输出是格式的数据框

seq    state     CMD
----------------------
 1       2       XOR
 2       2        OR
 3       4       XOR

【问题讨论】:

  • 最终结果如何?你愿意把它作为一个数据框吗?对预期输出的更多说明将在这里有所帮助
  • @dsk... 我已经添加了预期的输出。感谢您的评论。
  • 您现在能否检查一下下面的解决方案是否是您正在寻找的东西 - 如果您能接受并投票,将不胜感激.. :)

标签: json dataframe pyspark apache-kafka spark-streaming


【解决方案1】:

这对您来说可能是一个有效的解决方案 - 使用 explode()getItem() 如下 -

在此处将 json 加载到 Dataframe 中

a_json={
  '1000': 
    {
       '3': 
        {
           'seq': '1', 
           'state': '2', 
           'CMD': 'XOR' 
        }
    }
}
df = spark.createDataFrame([(a_json)])
df.show(truncate=False)

+-----------------------------------------+
|1000                                     |
+-----------------------------------------+
|[3 -> [CMD -> XOR, state -> 2, seq -> 1]]|
+-----------------------------------------+

这里的逻辑

df = df.select("*", F.explode("1000").alias("x", "y"))
df = df.withColumn("seq", df.y.getItem("seq")).withColumn("state", df.y.getItem("state")).withColumn("CMD", df.y.getItem("CMD"))
df.show(truncate=False)


 +-----------------------------------------+---+----------------------------------+---+-----+---+
|1000                                     |x  |y                                 |seq|state|CMD|
+-----------------------------------------+---+----------------------------------+---+-----+---+
|[3 -> [CMD -> XOR, state -> 2, seq -> 1]]|3  |[CMD -> XOR, state -> 2, seq -> 1]|1  |2    |XOR|
+-----------------------------------------+---+----------------------------------+---+-----+---+

根据更多输入更新代码

#Assuming that all the json columns are in a single column, hence making it an array column first.
df = df.withColumn("array_col", F.array("1000", "1001", "1003"))
#Then explode and getItem
df = df.withColumn("explod_col", F.explode("array_col"))
df = df.select("*", F.explode("explod_col").alias("x", "y"))
df_final = df.withColumn("seq", df.y.getItem("seq")).withColumn("state", df.y.getItem("state")).withColumn("CMD", df.y.getItem("CMD"))
df_final.select("seq","state","CMD").show()
|seq|state|CMD|
+---+-----+---+
|  1|    2|XOR|
|  2|    2| OR|
|  3|    4|XOR|
+---+-----+---+

【讨论】:

  • 嗨,前两个循环的键是可变的,所以使用 F.explode("1000") 不能解决问题。我已经更新了问题中的示例。
  • 另外,由于我从 Kafka 流中获取数据,我将值存储在数据框列“值”中。 df = df.select("*", F.explode("value").alias("x", "y")) 也给我一个错误:由于数据无法解析 'explode(value)'类型不匹配:函数explode的输入应该是数组或映射类型,而不是字符串;
  • 我已经解决了,请您检查更新的答案,这里的想法是首先从值列创建一个爆炸列 - df = df.withColumn("explod_col", F.explode( "value")) 并再次将键和值再次分隔在两个不同的列中 - df = df.select("*", F.explode("explod_col").alias("x", "y"))
  • 数据集很大,1000、1001、1003 并不是其中唯一的值。有成千上万个这样的键。问题中的示例只是为了说明数据集的格式。
  • 好友使用 df.columns 来排列所有列,最终您可以修改或传递一个列表作为您自己的 .. df = df.withColumn("array_col", F.array(df.columns ))
猜你喜欢
  • 1970-01-01
  • 2020-01-21
  • 1970-01-01
  • 1970-01-01
  • 2021-10-20
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2011-03-18
相关资源
最近更新 更多