【问题标题】:read pretty json format data through spark通过spark读取漂亮的json格式数据
【发布时间】:2021-02-17 12:04:06
【问题描述】:

我们通过scala中的spark读取S3中存在的小时格式的数据。例如,

spark.read.textFile("s3://'Bucket'/'key'/'yyyy'/'MM'/'dd'/'hh'/*").

spark.read.textFile 一次读取一行记录,因此例如 jsonLines 中存在的记录(一行中的完整 json 数据)被读取,并且可以稍后解析以从 json 中检索数据。

现在,我必须读取具有多个 json 但格式漂亮而不是 json 行的数据。使用相同的策略会导致损坏的记录错误。例如通过 spark.read.textFile 读取后获得的 Dataset[String]:

{
"a": 1, 
"b": 2
  }

_corrupt_record|
  +---------------+
  |              {|
  |       "a": 1, |
 |         "b": 2|
   |              }|

输入数据:

{
"key1": "value1",
"key2": "value2"
}
{
"key1": "value1",
 "key2": "value2"
}

预期输出

+------+------+
|key1  |key2  |
+------+------+
|value1|value2|
|value1|value2|
+------+------+

此文件有多个格式漂亮的 json,记录之间的分隔符为换行符。

已经使用的方法

  • spark.read.option("multiline", "true").json("") 。这不起作用,因为多行要求数据以 [{},{}] 的形式存在。

工作方法

val x=sparkSession
.read
.json(sc
  .wholeTextFiles(filePath)
  .values
  .flatMap(x=> {x
  .replace("\n", "")
   .replace("}{", "}}{{")
   .split("\\}\\{")}))

我只是想问是否有更好的方法,因为上述解决方案正在对数据进行一些切片和切块,这可能会导致大数据的性能问题?谢谢

【问题讨论】:

  • 阿曼,我们需要以表格格式查看输入数据和预期输出
  • 同时你可以试试这个答案来检查如何从 json 格式中提取..stackoverflow.com/questions/64640565/…
  • dsk,更新了问题。
  • 抱歉,回答延迟,请您检查解决方案,如果对您有帮助,请帮助接受和投票

标签: json apache-spark pyspark apache-spark-sql


【解决方案1】:

这对您来说可能是一个可行的解决方案,请使用 from_json() 并更正 schema 以正确解析 json

在此处创建数据框

df = spark.createDataFrame([(str([{"key1":"value1","key2":"value2"}, {"key1": "value3", "key2": "value4"}]))],T.StringType())
df.show(truncate=False)
+----------------------------------------------------------------------------+
|value                                                                       |
+----------------------------------------------------------------------------+
|[{'key1': 'value1', 'key2': 'value2'}, {'key1': 'value3', 'key2': 'value4'}]|
+----------------------------------------------------------------------------+

现在,使用 explode() 因为 value/json 列是一个列表,以便正确映射 最后,使用 getItem() 提取列

df = df.withColumn('col', F.from_json("value", T.ArrayType(T.StringType())))
df = df.withColumn("col", F.explode("col"))
df = df.withColumn("col", F.from_json("col", T.MapType(T.StringType(), T.StringType())))
df = df.withColumn("key1", df.col.getItem("key1")).withColumn("key2", df.col.getItem("key2"))
+----------------------------------------------------------------------------+--------------------------------+------+------+
|value                                                                       |col                             |key1  |key2  |
+----------------------------------------------------------------------------+--------------------------------+------+------+
|[{'key1': 'value1', 'key2': 'value2'}, {'key1': 'value3', 'key2': 'value4'}]|[key1 -> value1, key2 -> value2]|value1|value2|
|[{'key1': 'value1', 'key2': 'value2'}, {'key1': 'value3', 'key2': 'value4'}]|[key1 -> value3, key2 -> value4]|value3|value4|
+----------------------------------------------------------------------------+--------------------------------+------+------+

df.show(truncate=False)

【讨论】:

  • 对不起,我不明白答案。我的查询是从 S3 读取以多行 JSONS 形式存在的文件。我不知道文件的列名,因此无法执行任何操作你提到的。
  • 不是问题.. 是否可以共享任何示例文件?上述逻辑用于解压 json 文件/列。此处用于表示的列名仅建议
  • 你确定这是一个示例文件,其中包含来自一个大文件的一些记录。 pastie.org/p/36AaxjeLl7yehj4BUpxdjW.
猜你喜欢
  • 2017-03-28
  • 2023-04-09
  • 1970-01-01
  • 2013-09-23
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多