【问题标题】:How to effectively load and process JSON files containing different, evolving schemas如何有效地加载和处理包含不同的、不断发展的模式的 JSON 文件
【发布时间】:2021-10-24 19:14:38
【问题描述】:

我有一个复杂的数据问题,为了论证,不能修改。以下是来自数据库转储的虚拟示例 JSON 文件:

{"payload": {"table": "Table1", "data": {"colA": 1, "colB": 2}}}
{"payload": {"table": "Table2", "data": {"colA": 1, "colC": 2}}}

所以我在目录中的每个文件中都有来自多个表的数据,每个表都有自己的架构。随着在上游添加列,架构会随着时间而改变,因此我们无法定义静态架构,而是依赖架构推断。

这是我当前的工作流程(高级):

  • 将整个 JSON 数据目录加载到一个数据帧中
  • 查找这批更改中存在的唯一表
  • 对于每个表,将数据框过滤到仅该表数据
  • 读取表子集的架构(表 1 为 A+B,表 2 为 A+C)
  • 做一些验证
  • 将记录与目的地合并

示例代码:

import pyspark.sql.functions as F

df = spark.read.text(directory)
with_table_df = (
    df
    .withColumn("table", F.get_json_object('value', '$.payload.table'))
    .withColumn("json_payload_data", F.get_json_object('value', '$.payload.data'))
)
unique_tables = with_table_df.select('table').distinct().rdd.map(lambda r: r[0]).collect()

for table in unique_tables:
    filtered_df = with_table_df.filter(f"table = '{table}'")
    table_schema = spark.read.json(filtered_df.rdd.map(lambda row: row.json_payload_data)).schema

    changes_df = (
        filtered_df
        .withColumn('payload_data', F.from_json('json_payload_data', table_schema))
        .select('payload_data.*')
    )

    # do some validation
    if valid:
        changes_df.write.mode("append").option("mergeSchema", "true").saveAsTable(target_table)

我的问题是我无法使用 spark.read.json() 加载目录,因为它会将超集架构应用于所有记录,并且我无法确定哪些列是 Table1 列,哪些是 @987654325 @ 列。现在我作为文本加载,提取关键 JSON 元素 (payload.table),然后仅在我有相同模式的记录时解析为 JSON。这将有效,但会给驱动程序节点带来大量负载。

但我不认为过滤和迭代 Dataframe 行是一个好方法。我想以某种方式利用foreachPartition 将验证/选择逻辑映射到执行程序节点,但由于使用spark.read.json() 创建JSON 模式的方式(无法序列化到驱动程序节点),我无法做到这一点。

我怎样才能重新设计它以更适合 Spark 架构?

更新:

我希望修改数据创建过程,以便 JSON 文件按表分区,这样我就可以为每个唯一路径简单地 spark.read.json(table_path)

【问题讨论】:

  • 在您当前的示例中,这两行应该是不同的架构?因为对我来说,都是一样的。
  • Table1 有 A 列和 B 列,Table2 有 A 列和 C 列
  • 另一个问题,就像我说的那样,是表格可以随着时间的推移添加列。
  • 你不能用他们的名字定义不同的文件?
  • 很遗憾,没有,如果我可以将记录映射到具有同一张表的文件,那会容易得多。它是按时间分区的所有表的转储。

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


【解决方案1】:

我用你的数据创建了一个虚拟文件。

这是您要避免的简单代码:

df = spark.read.json("test.json")

df.show()
+-----------------+
|          payload|
+-----------------+
|[[1, 2,], Table1]|
|[[1,, 2], Table2]|
+-----------------+

df.printSchema()
root
 |-- payload: struct (nullable = true)
 |    |-- data: struct (nullable = true)
 |    |    |-- colA: long (nullable = true)
 |    |    |-- colB: long (nullable = true)
 |    |    |-- colC: long (nullable = true)
 |    |-- table: string (nullable = true)

这里的问题是,对于每个新的“col*”,您需要将它添加到您的架构和每个 json 行中。它是自动的,但不方便。

为此,您需要稍微欺骗一下您的架构。 Data字段的类型不是struct而是map:

from pyspark.sql import types as T

schm = T.StructType(
    [
        T.StructField(
            "payload",
            T.StructType(
                [
                    T.StructField("data", T.MapType(T.StringType(), T.IntegerType())),
                    T.StructField("table", T.StringType()),
                ]
            ),
        )
    ]
)


df = spark.read.json("test.json", schema=schm)

df.show()
+--------------------+                                                          
|             payload|
+--------------------+
|[[colA -> 1, colB...|
|[[colA -> 1, colC...|
+--------------------+

df.printSchema()
root
 |-- payload: struct (nullable = true)
 |    |-- data: map (nullable = true)
 |    |    |-- key: string
 |    |    |-- value: integer (valueContainsNull = true)
 |    |-- table: string (nullable = true)

【讨论】:

  • 嗯,您和其他答案使用MapType 遵循相同的方法,如果值不是整数怎么办?假设我们添加"colD": "some_value",这是一个字符串?我们可以将这些colAcolB 等提取到它们自己的列名中吗?
  • @TomNash 这是一个真正的用例还是一些理论问题?我们根据您的用例和您提供的示例数据为您提供了最佳答案。但是,如果您正在寻找适合您甚至没有的任何用例的最通用的解决方案,那么答案(如果存在)可能不会是一个有效的解决方案。
  • 是真的,我没有意识到我的文字描述不够具体,说它可以进化和改变,不得不举一个如此详细的例子(数据库表当然可以有列't 整数)。您的回答有助于重新思考问题,但我想我的解决方案必须适用于这个特定问题。
  • @TomNash 使用 map(String, String) 您可以加载任何数据,但会丢失类型。知道 JSON 中的类型只有字符串、布尔值或数字……日期类型不存在。
  • @OneCricketeer 我说的是 JSON - 抱歉不够清楚
【解决方案2】:

您的架构没有什么不同。你有一个带有"table" 字符串和"data" 映射的"payload" 结构

如果您对如何为"data" 定义架构感到困惑,请参阅MapType

data_schema = MapType(StringType(), IntegerType(), False)
payload_schema = StructType(
    [
        StructField("table", StringType(), False),
        StructField("data", data_schema, False),
    ]
)
schema = StructType(
    [
        StructField("payload", payload_schema, False),
    ]
)

【讨论】:

  • data_schema 并不总是字符串和整数的映射,它可以随着时间的推移而演变,可以是整数、日期、字符串等。它是如何工作的?
  • 我想这个问题并不清楚。如果是这种情况,我建议您修复创建 JSON 文件的过程,以便 Map 确实具有一致的值。
  • 对不起,我以为我很清楚“随着在上游添加列,架构会随着时间而改变”。不幸的是上游是我无法控制的,如果我能做到这一点,我不会有任何问题。不过我很欣赏这个答案,它肯定会有所帮助。
  • 很公平。将来会尝试提供任何可能的列的示例。
  • 将尝试更改 JSON 文件的创建过程,以便一个子目录中的所有记录具有相同的架构
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2019-04-06
  • 1970-01-01
  • 2019-04-04
  • 2012-05-24
  • 1970-01-01
  • 2019-05-12
  • 1970-01-01
相关资源
最近更新 更多