【问题标题】:Process single data set with different JSON schema rows using Pyspark使用 Pyspark 处理具有不同 JSON 模式行的单个数据集
【发布时间】: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


【解决方案1】:

如果你链接你的when 语句怎么办? 例如,

df.withColumn("nested_json", F.when(F.col("event_name") =="EventStart",F.from_json(...)).when(F.col("event_name") == "Action1", F. from_json(...)))

【讨论】:

  • 这正是我收到错误when() can only be applied on a Column previously generated by when() function
猜你喜欢
  • 2021-05-17
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-01-23
  • 1970-01-01
  • 2022-12-09
  • 1970-01-01
相关资源
最近更新 更多