【问题标题】:pyspark transform json array into multiple columnspyspark 将 json 数组转换为多列
【发布时间】:2021-09-22 20:39:01
【问题描述】:

我正在使用以下代码从 api 读取数据,其中有效负载为 json 格式,使用 pyspark 在 azure databricks 中。所有字段都定义为字符串,但一直运行到 json_tuple 要求所有参数都是字符串错误。

架构:

root
 |-- Payload: array (nullable = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- ActiveDate: string (nullable = true)
 |    |    |-- BusinessId: string (nullable = true)
 |    |    |-- BusinessName: string (nullable = true)

JSON:

 {
    "Payload": 
    [
        {
            "ActiveDate": "2008-11-25",
            "BusinessId": "5678",
            "BusinessName": "ACL"
        },
        {
            "ActiveDate": "2009-03-22",
            "BusinessId": "6789",
            "BusinessName": "BCL"
        }
    ]
}

PySpark:

from pyspark.sql import functions as F
df = df.select(F.col('Payload'), F.json_tuple(F.col('Payload'), 'ActiveDate', 'BusinessId', 'BusinessName') \.alias('ActiveDate', 'BusinessId', 'BusinessName'))
df.write.format("delta").mode("overwrite").saveAsTable("delta_payload")

错误:

AnalysisException: cannot resolve 'json_tuple(`Payload`, 'ActiveDate', 'BusinessId', 'BusinessName')' due to data type mismatch: json_tuple requires that all arguments are strings;

【问题讨论】:

    标签: python-3.x apache-spark pyspark apache-spark-sql


    【解决方案1】:

    从您的架构看来,JSON 已被解析,因此 Payload 属于 ArrayType 而不是包含 JSON 的 StringType,因此出现错误。

    您可能需要explode 而不是json_tuple:

    >>> from pyspark.sql.functions import explode
    >>> df = spark.createDataFrame([{
    ...     "Payload":
    ...     [
    ...         {
    ...             "ActiveDate": "2008-11-25",
    ...             "BusinessId": "5678",
    ...             "BusinessName": "ACL"
    ...         },
    ...         {
    ...             "ActiveDate": "2009-03-22",
    ...             "BusinessId": "6789",
    ...             "BusinessName": "BCL"
    ...         }
    ...     ]
    ... }])
    >>> df.schema
    StructType(List(StructField(Payload,ArrayType(MapType(StringType,StringType,true),true),true)))
    >>> df.select(explode("Payload").alias("x")).select("x.ActiveDate", "x.BusinessName", "x.BusinessId").show()
    +----------+------------+----------+
    |ActiveDate|BusinessName|BusinessId|
    +----------+------------+----------+
    |2008-11-25|         ACL|      5678|
    |2009-03-22|         BCL|      6789|
    +----------+------------+----------+
    

    【讨论】:

    • 嗨@Czaporka,我在尝试使用这行代码时遇到了错误NameError: name 'StructType' is not defined。 StructType(List(StructField(Report_Entry,ArrayType(MapType(StringType,StringType,true),true),true)))
    • 嗨@paone,那一行只是我在交互式解释器会话中输入df.schema 后得到的输出。我包含它只是为了显示我的 DataFrame 的架构。您应该执行的行以>>> 为前缀。最重要的是,您可能需要第一个(import explode)和最后一个带有 select 的。
    • 嗨@Czaporka,谢谢。这样可行。但是我现在遇到的问题是df.show(truncate=False)以表格格式显示数据+----------+------------+--------- -+ |活动日期|企业名称|企业 ID| +----------+------------+----------+ |2008-11-25|访问控制列表| 5678| |2009-03-22| BCL| 6789| +---------+------------+----------+ delta_tbl 结果是数组df.write.format("delta").mode("overwrite").saveAsTable("delta_tbl") spark.sql("SELECT * FROM delta_tbl LIMIT 1")
    • @paone 您是否将select 和explode 的结果分配回df?在我的代码示例中,我再次执行df.select(...).show() 只是为了显示select 将返回什么;但在您的实际代码中,您需要在编写它之前将其结果分配给df,即df = df.select(...); df.write.format(...)...(就像在您的原始代码中一样),或者只是执行df.select(...).write.format(...)...。
    • 嗨@Czaporka,是的,我做到了,现在可以使用。谢谢。
    猜你喜欢
    • 2021-04-06
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-09-29
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多