【发布时间】: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