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
相关产品推荐
相关产品推荐

