You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

PySpark读取Blob容器复杂JSON文件并动态生成表结构求助

PySpark读取嵌套JSON并生成目标表结构实现方案

步骤1:读取Blob存储中的JSON文件

先配置Spark访问Azure Blob存储的权限(如存储账户密钥或SAS令牌),再读取JSON文件:

from pyspark.sql import SparkSession
from pyspark.sql.functions import explode, col
from pyspark.sql.types import StringType

# 初始化SparkSession
spark = SparkSession.builder.appName("NestedJsonToTable").getOrCreate()

# 配置Azure Blob存储(替换为你的存储信息)
spark.conf.set("fs.azure.account.key.<your-storage-account>.blob.core.windows.net", "<your-access-key>")

# 读取单条嵌套JSON,开启multiLine模式
df = spark.read.option("multiLine", "true").json("wasbs://<container-name>@<your-storage-account>.blob.core.windows.net/<file-path>.json")

步骤2:解析嵌套的tables结构

展开tables数组,提取列定义和行数据:

# 展开tables数组,获取目标表对象
table_df = df.select(explode(col("tables")).alias("table"))

# 提取列名列表和类型映射
columns_info = table_df.select("table.columns").first()[0]
column_names = [col["name"] for col in columns_info]
column_type_map = {col["name"]: col["type"] for col in columns_info}

# 提取行数据并映射为列
rows_df = table_df.select(explode(col("table.rows")).alias("row"))
target_df = rows_df.select(*[col("row")[i].alias(column_names[i]) for i in range(len(column_names))])

步骤3:动态转换数据类型

根据JSON中的类型定义,转换对应列的数据类型:

# 定义Spark类型映射
type_mapping = {
    "string": StringType(),
    "datetime": "timestamp"
}

# 遍历列执行类型转换
for col_name, col_type in column_type_map.items():
    target_df = target_df.withColumn(col_name, col(col_name).cast(type_mapping.get(col_type, StringType())))

步骤4:验证结果

查看最终表结构和数据:

target_df.printSchema()
target_df.show(truncate=False)

执行后即可得到与目标结构一致的DataFrame。


内容的提问来源于stack exchange,提问作者coding

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.24 16:02:47