【发布时间】:2021-10-30 14:40:10
【问题描述】:
我正在使用 PySpark,我需要处理附加到单个数据框中的日志文件。大多数列看起来正常,但其中一列在 {} 中有 JSON 字符串。基本上,每一行都是一个单独的事件,对于 JSON 字符串,我可以应用单独的模式。但我不知道这里处理数据的最佳方式是什么。
此表稍后将有助于以我需要的方式聚合事件。
我尝试使用函数withColumn 并使用from_json。它成功地用于单个列:
from pyspark.sql.types import *
import pyspark.sql.functions as F
df = (df
.withColumn("nested_json",
F.when(F.col("event_name") == "EventStart",F.from_json("json_string","Name String, Version Int, Id Int")))
当我查询nested_json 时,它为我的第一行做了我想要的。但它是应用于整个列的架构,我想处理每一行取决于event_name
我很天真并尝试这样做:
from pyspark.sql.types import *
import pyspark.sql.functions as F
df = (df
.withColumn("nested_json",
F.when(F.col("event_name") == "EventStart",F.from_json("json_string","Name String, Version Int, Id Int"))
F.when(F.col("event_name") == "Action1",F.from_json("json_string","Name String, Version Int, UserName String, PosX int, PosY int"))
)
这在when() can only be applied on a Column previously generated by when() function 下运行失败
我假设,我的第一个 withColumn 为整个列应用了架构。
我还有哪些其他选项可以应用基于 event_name 值和扁平值的 JSON 架构?
【问题讨论】:
-
对于没有
when函数的所有行,如何使用"Name string, Version int, Id int, UserName string, PosX int . . ."之类的整个架构? -
这个选项确实有效。
标签: python json apache-spark pyspark databricks