【问题标题】:Creating a Spark Structured Streaming Schema for nested Json为嵌套的 Json 创建 Spark 结构化流模式
【发布时间】:2022-10-16 15:39:29
【问题描述】:

我想为我的结构化流作业(在 python 中)定义模式,但我无法以我想要的方式获取数据帧模式。

对于这个 json

{
    "messages": [{
        "IdentityNumber": 1,
        "body": {
            "Alert": "This is the payload"
        },
        "regionNumber": 11000002
    }]
}

我使用下面的代码作为模式

schema1 = StructType([StructField("messages", ArrayType(   
    StructType( 
        [
            StructField("body", StructType( [StructField("Alert", StringType())]) )
        ]
    )
    ,True))])

但我得到我的架构

df-> 消息-> 正文-> 警报

虽然我想要这样的东西

df-> 警报

即一个名为 alert 的单列数据框,它将包含所有作为警报出现的字符串消息。 我应该在我定义的架构中进行哪些更改?

【问题讨论】:

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


    【解决方案1】:

    如果您正在读取有关此架构的数据,则该架构是可以的。

    如果您在读取上述模式中的 json 后需要提取嵌套字段,只需使用点表示法即可。例如:

    df.select(col("messages[0].body.alert"))
    

    如果您需要操作和分解所有数组元素,请查看这篇解释您必须执行的不同选项的文章: https://docs.databricks.com/_static/notebooks/transform-complex-data-types-scala.html

    上面的答案和文章一样在 scala 中,但是大多数 spark sql API 很容易转移到 pySpark。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2019-11-26
      • 2015-01-27
      • 2020-04-29
      • 2022-11-24
      • 1970-01-01
      • 1970-01-01
      • 2018-09-12
      相关资源
      最近更新 更多